548 lines
23 KiB
Python
548 lines
23 KiB
Python
import logging
|
|
|
|
from celery import Celery
|
|
from celery.signals import (
|
|
celeryd_init,
|
|
heartbeat_sent,
|
|
worker_init,
|
|
worker_process_init,
|
|
worker_process_shutdown,
|
|
worker_ready,
|
|
worker_shutdown,
|
|
)
|
|
|
|
from app.config import settings
|
|
from app.enums.celery_queue import CeleryQueue, CeleryTaskName
|
|
from app.models.base import engine
|
|
from app.tasks.async_runner import close_loop, run_async
|
|
|
|
logger = logging.getLogger("video_gen")
|
|
|
|
|
|
# 显式注册所有 Celery 任务模块,避免新增任务文件后 worker 启动时未注册任务。
|
|
# 不再依赖 app.tasks.__init__ 内部 import,也不再依赖 autodiscover_tasks。
|
|
CELERY_TASK_IMPORTS = (
|
|
"app.tasks.generation_create_tasks",
|
|
"app.tasks.generation_poll_tasks",
|
|
"app.tasks.generation_download_tasks",
|
|
"app.tasks.generation_recovery_tasks",
|
|
"app.tasks.video_upscale_tasks",
|
|
"app.tasks.hot_opening_replicate_tasks",
|
|
"app.tasks.shot_replicate_tasks",
|
|
"app.tasks.shot_replicate_flow_tasks",
|
|
"app.tasks.module_async_recovery_tasks",
|
|
"app.tasks.module_generation_v2_tasks",
|
|
"app.tasks.private_portrait_asset_tasks",
|
|
"app.tasks.vp_v3_asset_tasks",
|
|
"app.tasks.celery_runtime_tasks",
|
|
"app.tasks.credit_tasks",
|
|
"app.tasks.api_generation_tasks",
|
|
"app.tasks.api_recovery_tasks",
|
|
"app.tasks.api_upscale_tasks",
|
|
)
|
|
|
|
|
|
RECOVERY_QUEUE = settings.CELERY_RECOVERY_QUEUE or CeleryQueue.GEN_RECOVERY.value
|
|
|
|
|
|
def _derive_redis_db(url: str, db_no: int) -> str:
|
|
if not url:
|
|
return url
|
|
import re
|
|
|
|
if re.search(r"/\d+$", url):
|
|
return re.sub(r"/\d+$", f"/{db_no}", url)
|
|
return url.rstrip("/") + f"/{db_no}"
|
|
|
|
|
|
def _beat_schedule() -> dict:
|
|
schedule: dict = {}
|
|
if bool(getattr(settings, "POLL_DUE_DISPATCH_ENABLED", True)):
|
|
schedule["dispatch-due-poll-tasks-every-minute"] = {
|
|
"task": CeleryTaskName.DISPATCH_DUE_POLL.value,
|
|
"schedule": max(1, int(settings.POLL_DUE_DISPATCH_INTERVAL_SECONDS or 60)),
|
|
"options": {
|
|
"queue": RECOVERY_QUEUE,
|
|
"priority": settings.DOWNLOAD_TASK_PRIORITY_RECOVER,
|
|
},
|
|
}
|
|
schedule["generation-create-recovery"] = {
|
|
"task": CeleryTaskName.RECOVER_CREATE.value,
|
|
"schedule": max(1, int(settings.GENERATION_CREATE_RECOVERY_INTERVAL_SECONDS or 60)),
|
|
"options": {
|
|
"queue": RECOVERY_QUEUE,
|
|
"priority": settings.DOWNLOAD_TASK_PRIORITY_RECOVER,
|
|
},
|
|
}
|
|
schedule["module-async-recovery"] = {
|
|
"task": CeleryTaskName.MODULE_ASYNC_RECOVERY.value,
|
|
"schedule": max(1, int(settings.MODULE_ASYNC_RECOVERY_INTERVAL_SECONDS or 60)),
|
|
"options": {
|
|
"queue": RECOVERY_QUEUE,
|
|
"priority": settings.DOWNLOAD_TASK_PRIORITY_RECOVER,
|
|
},
|
|
}
|
|
schedule["video-upscale-recovery-every-minute"] = {
|
|
"task": CeleryTaskName.VIDEO_UPSCALE_RECOVER.value,
|
|
"schedule": 60,
|
|
"options": {
|
|
"queue": RECOVERY_QUEUE,
|
|
"priority": settings.DOWNLOAD_TASK_PRIORITY_RECOVER,
|
|
},
|
|
}
|
|
schedule["api-generation-recovery-every-minute"] = {
|
|
"task": "api_generation.recover_tasks_once",
|
|
"schedule": 60,
|
|
"options": {
|
|
"queue": RECOVERY_QUEUE,
|
|
"priority": settings.DOWNLOAD_TASK_PRIORITY_RECOVER,
|
|
},
|
|
}
|
|
schedule["generation-download-recovery"] = {
|
|
"task": CeleryTaskName.RECOVER_DOWNLOAD.value,
|
|
"schedule": max(1, int(settings.DOWNLOAD_RECOVERY_INTERVAL_SECONDS or 60)),
|
|
"options": {"queue": RECOVERY_QUEUE, "priority": settings.DOWNLOAD_TASK_PRIORITY_RECOVER},
|
|
}
|
|
schedule["shot-split-recovery"] = {
|
|
"task": CeleryTaskName.SHOT_SPLIT_RECOVERY.value,
|
|
"schedule": 60,
|
|
"options": {"queue": RECOVERY_QUEUE, "priority": settings.DOWNLOAD_TASK_PRIORITY_RECOVER},
|
|
}
|
|
schedule["shot-analysis-recovery"] = {
|
|
"task": CeleryTaskName.SHOT_ANALYSIS_RECOVERY.value,
|
|
"schedule": 60,
|
|
"options": {"queue": RECOVERY_QUEUE, "priority": settings.DOWNLOAD_TASK_PRIORITY_RECOVER},
|
|
}
|
|
schedule["celery-runtime-reconcile"] = {
|
|
"task": CeleryTaskName.CELERY_RUNTIME_RECONCILE.value,
|
|
"schedule": max(60, int(settings.CELERY_RUNTIME_RECONCILE_INTERVAL_SECONDS or 300)),
|
|
"options": {"queue": RECOVERY_QUEUE, "priority": settings.DOWNLOAD_TASK_PRIORITY_RECOVER},
|
|
}
|
|
schedule["celery-runtime-registry-gc"] = {
|
|
"task": CeleryTaskName.CELERY_RUNTIME_GC.value,
|
|
"schedule": max(60, int(settings.CELERY_RUNTIME_GC_INTERVAL_SECONDS or 600)),
|
|
"options": {"queue": RECOVERY_QUEUE, "priority": settings.DOWNLOAD_TASK_PRIORITY_RECOVER},
|
|
}
|
|
schedule["private-portrait-sync-due-assets-every-minute"] = {
|
|
"task": CeleryTaskName.PRIVATE_PORTRAIT_SYNC_DUE_ASSETS.value,
|
|
"schedule": 60,
|
|
"options": {"queue": CeleryQueue.GEN_PRIVATE_PORTRAIT.value},
|
|
}
|
|
schedule["private-portrait-recover-remote-deletes-every-5-minutes"] = {
|
|
"task": CeleryTaskName.PRIVATE_PORTRAIT_RECOVER_REMOTE_DELETES.value,
|
|
"schedule": 300,
|
|
"options": {"queue": CeleryQueue.GEN_PRIVATE_PORTRAIT.value},
|
|
}
|
|
schedule["credit-maintenance-every-minute"] = {
|
|
"task": CeleryTaskName.CREDIT_MAINTENANCE.value,
|
|
"schedule": 60,
|
|
"options": {"queue": CeleryQueue.GEN_CREDIT_MAINTENANCE.value},
|
|
}
|
|
schedule["vp-v3-sync-due-assets-every-minute"] = {
|
|
"task": CeleryTaskName.VP_V3_SYNC_DUE_ASSETS.value,
|
|
"schedule": 60,
|
|
"options": {"queue": CeleryQueue.GEN_PRIVATE_PORTRAIT.value},
|
|
}
|
|
schedule["vp-v3-recover-remote-deletes-every-5-minutes"] = {
|
|
"task": CeleryTaskName.VP_V3_RECOVER_REMOTE_DELETES.value,
|
|
"schedule": 300,
|
|
"options": {"queue": CeleryQueue.GEN_PRIVATE_PORTRAIT.value},
|
|
}
|
|
return schedule
|
|
|
|
|
|
broker_url = settings.CELERY_BROKER_URL or (_derive_redis_db(settings.REDIS_URL, 1) if settings.REDIS_URL else "")
|
|
backend_url = settings.CELERY_RESULT_BACKEND or (_derive_redis_db(settings.REDIS_URL, 2) if settings.REDIS_URL else "")
|
|
|
|
if broker_url:
|
|
celery_app = Celery("videogen", include=CELERY_TASK_IMPORTS)
|
|
celery_app.conf.update(
|
|
broker_url=broker_url,
|
|
result_backend=backend_url or broker_url,
|
|
imports=CELERY_TASK_IMPORTS,
|
|
task_serializer="json",
|
|
accept_content=["json"],
|
|
result_serializer="json",
|
|
result_expires=max(60, int(settings.CELERY_RESULT_EXPIRES_SECONDS or 7200)),
|
|
task_store_errors_even_if_ignored=True,
|
|
timezone="Asia/Shanghai",
|
|
enable_utc=True,
|
|
task_soft_time_limit=600,
|
|
task_time_limit=900,
|
|
task_acks_late=True,
|
|
task_reject_on_worker_lost=True,
|
|
task_track_started=True,
|
|
beat_schedule=_beat_schedule(),
|
|
task_annotations={
|
|
# 生成链路任务以数据库状态为准,不依赖 Celery result backend。
|
|
# 这里忽略结果可避免任务误返回 ORM / 非 JSON 对象时触发结果序列化失败。
|
|
# "generation.chatapi_create_generation_task": {"ignore_result": True},
|
|
# "generation.poll_generation_task": {"ignore_result": True},
|
|
# "generation.download_generation_result_task": {"ignore_result": True},
|
|
"hot_opening.start_image_prompt_optimize": {"ignore_result": True},
|
|
"hot_opening.start_video_prompt_optimize": {"ignore_result": True},
|
|
"shot_replicate.analyze_original_video": {
|
|
"ignore_result": True,
|
|
"soft_time_limit": int(settings.SHOT_ANALYSIS_SOFT_TIME_LIMIT_SECONDS or 3720),
|
|
"time_limit": int(settings.SHOT_ANALYSIS_TIME_LIMIT_SECONDS or 3900),
|
|
},
|
|
"shot_replicate.analyze_custom_segment_video": {
|
|
"ignore_result": True,
|
|
"soft_time_limit": int(settings.SHOT_ANALYSIS_SOFT_TIME_LIMIT_SECONDS or 3720),
|
|
"time_limit": int(settings.SHOT_ANALYSIS_TIME_LIMIT_SECONDS or 3900),
|
|
},
|
|
"shot_replicate.split_one_segment": {"ignore_result": True},
|
|
"shot_replicate.start_image_prompt_optimize": {"ignore_result": True},
|
|
"shot_replicate.start_video_prompt_optimize": {"ignore_result": True},
|
|
"module_generation_v2.start_video_prompt_optimize": {"ignore_result": True},
|
|
CeleryTaskName.VIDEO_UPSCALE_EXECUTE_LOCAL.value: {
|
|
"ignore_result": True,
|
|
"soft_time_limit": max(60, int(settings.VIDEO_UPSCALE_LOCAL_TIMEOUT_SECONDS or 3600)) + 60,
|
|
"time_limit": max(60, int(settings.VIDEO_UPSCALE_LOCAL_TIMEOUT_SECONDS or 3600)) + 300,
|
|
},
|
|
CeleryTaskName.VIDEO_UPSCALE_SUBMIT_REMOTE.value: {"ignore_result": True},
|
|
CeleryTaskName.VIDEO_UPSCALE_POLL_REMOTE.value: {"ignore_result": True},
|
|
CeleryTaskName.VIDEO_UPSCALE_DOWNLOAD_REMOTE_RESULT.value: {"ignore_result": True},
|
|
CeleryTaskName.VIDEO_UPSCALE_FINALIZE.value: {"ignore_result": True},
|
|
CeleryTaskName.VIDEO_UPSCALE_RECOVER.value: {"ignore_result": True},
|
|
CeleryTaskName.RECOVER_CREATE.value: {"ignore_result": True},
|
|
CeleryTaskName.MODULE_ASYNC_RECOVERY.value: {"ignore_result": True},
|
|
CeleryTaskName.SHOT_ANALYSIS_RECOVERY.value: {"ignore_result": True},
|
|
CeleryTaskName.SHOT_SPLIT_RECOVERY.value: {"ignore_result": True},
|
|
CeleryTaskName.CREDIT_MAINTENANCE.value: {"ignore_result": True},
|
|
},
|
|
worker_prefetch_multiplier=1,
|
|
worker_cancel_long_running_tasks_on_connection_loss=True,
|
|
broker_transport_options={
|
|
"visibility_timeout": max(
|
|
3600,
|
|
int(settings.VIDEO_UPSCALE_LOCAL_TIMEOUT_SECONDS or 3600) + 600,
|
|
int(settings.VIDEO_UPSCALE_REMOTE_RESULT_DOWNLOAD_TIMEOUT_SECONDS or 600) + 600,
|
|
int(settings.SHOT_ANALYSIS_TIME_LIMIT_SECONDS or 3900) + 600,
|
|
),
|
|
"queue_order_strategy": "priority",
|
|
"priority_steps": list(range(10)),
|
|
"sep": ":",
|
|
},
|
|
task_routes={
|
|
CeleryTaskName.CHATAPI_CREATE.value: {"queue": CeleryQueue.GEN_CHATAPI_CREATE.value},
|
|
CeleryTaskName.POLL_GENERATION.value: {"queue": CeleryQueue.GEN_PROVIDER_POLL.value},
|
|
CeleryTaskName.DOWNLOAD_GENERATION_RESULT.value: {"queue": CeleryQueue.GEN_RESULT_DOWNLOAD.value},
|
|
CeleryTaskName.VIDEO_UPSCALE_EXECUTE_LOCAL.value: {
|
|
"queue": settings.VIDEO_UPSCALE_LOCAL_QUEUE or CeleryQueue.GEN_VIDEO_UPSCALE_LOCAL.value
|
|
},
|
|
CeleryTaskName.VIDEO_UPSCALE_SUBMIT_REMOTE.value: {
|
|
"queue": settings.VIDEO_UPSCALE_REMOTE_QUEUE or CeleryQueue.GEN_VIDEO_UPSCALE_REMOTE.value
|
|
},
|
|
CeleryTaskName.VIDEO_UPSCALE_POLL_REMOTE.value: {
|
|
"queue": settings.VIDEO_UPSCALE_REMOTE_QUEUE or CeleryQueue.GEN_VIDEO_UPSCALE_REMOTE.value
|
|
},
|
|
CeleryTaskName.VIDEO_UPSCALE_DOWNLOAD_REMOTE_RESULT.value: {
|
|
"queue": settings.VIDEO_UPSCALE_REMOTE_QUEUE or CeleryQueue.GEN_VIDEO_UPSCALE_REMOTE.value
|
|
},
|
|
CeleryTaskName.VIDEO_UPSCALE_FINALIZE.value: {
|
|
"queue": settings.VIDEO_UPSCALE_LOCAL_QUEUE or CeleryQueue.GEN_VIDEO_UPSCALE_LOCAL.value
|
|
},
|
|
CeleryTaskName.VIDEO_UPSCALE_RECOVER.value: {"queue": RECOVERY_QUEUE},
|
|
CeleryTaskName.DISPATCH_DUE_POLL.value: {"queue": RECOVERY_QUEUE},
|
|
"hot_opening.start_image_prompt_optimize": {"queue": CeleryQueue.GEN_CHATAPI_CREATE.value},
|
|
"hot_opening.start_video_prompt_optimize": {"queue": CeleryQueue.GEN_CHATAPI_CREATE.value},
|
|
CeleryTaskName.SHOT_ANALYZE_ORIGINAL.value: {"queue": CeleryQueue.GEN_SHOT_ANALYSIS.value},
|
|
CeleryTaskName.SHOT_ANALYZE_CUSTOM_SEGMENT.value: {"queue": CeleryQueue.GEN_SHOT_ANALYSIS.value},
|
|
CeleryTaskName.SHOT_SPLIT_ONE.value: {"queue": CeleryQueue.GEN_SHOT_SPLIT.value},
|
|
"shot_replicate.start_image_prompt_optimize": {"queue": CeleryQueue.GEN_CHATAPI_CREATE.value},
|
|
"shot_replicate.start_video_prompt_optimize": {"queue": CeleryQueue.GEN_CHATAPI_CREATE.value},
|
|
"module_generation_v2.start_video_prompt_optimize": {"queue": CeleryQueue.GEN_CHATAPI_CREATE.value},
|
|
# 恢复扫描统一走独立队列,避免占用下载/轮询/创建业务 worker。
|
|
CeleryTaskName.STARTUP_RECOVERY.value: {"queue": RECOVERY_QUEUE},
|
|
CeleryTaskName.SHOT_SPLIT_RECOVERY.value: {"queue": RECOVERY_QUEUE},
|
|
CeleryTaskName.SHOT_ANALYSIS_RECOVERY.value: {"queue": RECOVERY_QUEUE},
|
|
CeleryTaskName.RECOVER_DOWNLOAD.value: {"queue": RECOVERY_QUEUE},
|
|
CeleryTaskName.RECOVER_GENERATION.value: {"queue": RECOVERY_QUEUE},
|
|
CeleryTaskName.RECOVER_CREATE.value: {"queue": RECOVERY_QUEUE},
|
|
CeleryTaskName.MODULE_ASYNC_RECOVERY.value: {"queue": RECOVERY_QUEUE},
|
|
CeleryTaskName.CELERY_RUNTIME_RECONCILE.value: {"queue": RECOVERY_QUEUE},
|
|
CeleryTaskName.CELERY_RUNTIME_GC.value: {"queue": RECOVERY_QUEUE},
|
|
CeleryTaskName.PRIVATE_PORTRAIT_POLL_ASSET.value: {"queue": CeleryQueue.GEN_PRIVATE_PORTRAIT.value},
|
|
CeleryTaskName.PRIVATE_PORTRAIT_SYNC_DUE_ASSETS.value: {"queue": CeleryQueue.GEN_PRIVATE_PORTRAIT.value},
|
|
CeleryTaskName.PRIVATE_PORTRAIT_DELETE_ASSET.value: {"queue": CeleryQueue.GEN_PRIVATE_PORTRAIT.value},
|
|
CeleryTaskName.PRIVATE_PORTRAIT_DELETE_GROUP.value: {"queue": CeleryQueue.GEN_PRIVATE_PORTRAIT.value},
|
|
CeleryTaskName.PRIVATE_PORTRAIT_DELETE_PROJECT.value: {"queue": CeleryQueue.GEN_PRIVATE_PORTRAIT.value},
|
|
CeleryTaskName.PRIVATE_PORTRAIT_RECOVER_REMOTE_DELETES.value: {"queue": CeleryQueue.GEN_PRIVATE_PORTRAIT.value},
|
|
CeleryTaskName.CREDIT_MAINTENANCE.value: {"queue": CeleryQueue.GEN_CREDIT_MAINTENANCE.value},
|
|
CeleryTaskName.VP_V3_POLL_ASSET.value: {"queue": CeleryQueue.GEN_PRIVATE_PORTRAIT.value},
|
|
CeleryTaskName.VP_V3_SYNC_DUE_ASSETS.value: {"queue": CeleryQueue.GEN_PRIVATE_PORTRAIT.value},
|
|
CeleryTaskName.VP_V3_DELETE_ASSET.value: {"queue": CeleryQueue.GEN_PRIVATE_PORTRAIT.value},
|
|
CeleryTaskName.VP_V3_DELETE_PROJECT.value: {"queue": CeleryQueue.GEN_PRIVATE_PORTRAIT.value},
|
|
CeleryTaskName.VP_V3_RECOVER_REMOTE_DELETES.value: {"queue": CeleryQueue.GEN_PRIVATE_PORTRAIT.value},
|
|
},
|
|
)
|
|
else:
|
|
celery_app = None
|
|
|
|
|
|
async def _try_acquire_startup_recovery_lock() -> bool:
|
|
"""任意 worker 启动时都可尝试抢恢复投递锁,避免依赖 hostname 命名。"""
|
|
from app.services.redis_registry_service import redis_acquire_lock
|
|
|
|
token = await redis_acquire_lock(
|
|
lock_key=settings.CELERY_STARTUP_RECOVERY_LOCK_KEY,
|
|
ttl_seconds=int(settings.CELERY_STARTUP_RECOVERY_LOCK_TTL_SECONDS or 120),
|
|
log_context="celery_startup_recovery",
|
|
)
|
|
return bool(token)
|
|
|
|
|
|
def _worker_name_from_sender(sender=None, **kwargs) -> str:
|
|
candidates = (
|
|
getattr(sender, "hostname", None),
|
|
getattr(sender, "name", None),
|
|
kwargs.get("hostname"),
|
|
kwargs.get("nodename"),
|
|
)
|
|
instance = kwargs.get("instance")
|
|
if instance is not None:
|
|
candidates += (
|
|
getattr(instance, "hostname", None),
|
|
getattr(instance, "name", None),
|
|
)
|
|
for value in candidates:
|
|
normalized = str(value or "").strip()
|
|
if normalized:
|
|
return normalized
|
|
return ""
|
|
|
|
|
|
def _worker_runtime_metadata(sender=None) -> dict:
|
|
queues: list[str] = []
|
|
try:
|
|
consumer = getattr(sender, "consumer", None)
|
|
task_consumer = getattr(consumer, "task_consumer", None)
|
|
for queue in list(getattr(task_consumer, "queues", None) or []):
|
|
name = str(getattr(queue, "name", queue) or "").strip()
|
|
if name and name not in queues:
|
|
queues.append(name)
|
|
except Exception:
|
|
pass
|
|
|
|
pool = getattr(sender, "pool", None)
|
|
pool_type = type(pool).__name__ if pool is not None else None
|
|
configured_concurrency = None
|
|
for value in (
|
|
getattr(pool, "limit", None),
|
|
getattr(sender, "concurrency", None),
|
|
):
|
|
try:
|
|
parsed = int(value)
|
|
except (TypeError, ValueError):
|
|
continue
|
|
if parsed > 0:
|
|
configured_concurrency = parsed
|
|
break
|
|
|
|
return {
|
|
"queues": queues,
|
|
"pool_type": pool_type,
|
|
"configured_concurrency": configured_concurrency,
|
|
}
|
|
|
|
|
|
@celeryd_init.connect
|
|
def on_celeryd_init(sender=None, instance=None, **kwargs):
|
|
"""尽早生成 Worker 主实例 token,确保 prefork 子进程继承。"""
|
|
try:
|
|
from app.services.celery_runtime.worker_service import initialize_worker_main_identity
|
|
|
|
initialize_worker_main_identity(
|
|
_worker_name_from_sender(sender, instance=instance, **kwargs) or None,
|
|
before_pool=True,
|
|
)
|
|
except Exception:
|
|
logger.exception("Celery Worker 主实例身份初始化失败。signal=celeryd_init")
|
|
|
|
|
|
@worker_init.connect
|
|
def on_worker_init(sender=None, **kwargs):
|
|
"""worker_init 幂等兜底,仍处于进程池创建之前。"""
|
|
try:
|
|
from app.services.celery_runtime.worker_service import initialize_worker_main_identity
|
|
|
|
initialize_worker_main_identity(
|
|
_worker_name_from_sender(sender, **kwargs) or None,
|
|
before_pool=True,
|
|
)
|
|
except Exception:
|
|
logger.exception("Celery Worker 主实例身份初始化失败。signal=worker_init")
|
|
|
|
|
|
@worker_ready.connect
|
|
def on_worker_ready(sender=None, **kwargs):
|
|
"""注册当前 Worker 主实例,并协调实例级与全局启动恢复。"""
|
|
if celery_app is None:
|
|
return
|
|
|
|
worker_name = _worker_name_from_sender(sender, **kwargs)
|
|
metadata = _worker_runtime_metadata(sender)
|
|
current_identity = None
|
|
|
|
try:
|
|
from app.services.celery_runtime.worker_service import (
|
|
initialize_worker_main_identity,
|
|
register_worker_instance,
|
|
)
|
|
|
|
# 正常情况下 token 已在 worker_init 前创建;这里仅做 late fallback。
|
|
initialize_worker_main_identity(worker_name or None, before_pool=False)
|
|
current_identity = run_async(
|
|
register_worker_instance(
|
|
worker_name=worker_name,
|
|
queues=metadata["queues"],
|
|
pool_type=metadata["pool_type"],
|
|
configured_concurrency=metadata["configured_concurrency"],
|
|
)
|
|
)
|
|
except Exception:
|
|
# Redis 或身份注册失败不能阻塞 Worker 启动,任务级执行锁仍会 fail-closed。
|
|
logger.exception("Celery Worker 主实例注册失败。worker_name=%s", worker_name)
|
|
|
|
if current_identity is not None:
|
|
try:
|
|
from app.services.celery_runtime.recovery_service import (
|
|
mark_stale_worker_instance_candidates,
|
|
)
|
|
|
|
run_async(
|
|
mark_stale_worker_instance_candidates(
|
|
worker_name=current_identity.worker_name,
|
|
current_worker_instance_id=current_identity.worker_instance_id,
|
|
supports_targeted_recovery=current_identity.supports_targeted_recovery,
|
|
)
|
|
)
|
|
except Exception:
|
|
logger.exception(
|
|
"Worker 旧主实例任务候选标记失败。worker_name=%s worker_instance_id=%s",
|
|
current_identity.worker_name,
|
|
current_identity.worker_instance_id,
|
|
)
|
|
|
|
if not bool(getattr(settings, "CELERY_STARTUP_RECOVERY_ENABLED", True)):
|
|
logger.info("启动容灾恢复已关闭。CELERY_STARTUP_RECOVERY_ENABLED=false")
|
|
return
|
|
|
|
try:
|
|
if not run_async(_try_acquire_startup_recovery_lock()):
|
|
return
|
|
except Exception:
|
|
logger.exception("启动容灾恢复锁获取失败,已跳过本次自动恢复投递")
|
|
return
|
|
|
|
try:
|
|
from app.services.celery_runtime.recovery_service import set_startup_barrier
|
|
from app.tasks.generation_recovery_tasks import startup_recovery_once
|
|
from app.tasks.api_recovery_tasks import api_generation_recover_tasks_once
|
|
|
|
run_async(set_startup_barrier())
|
|
countdown = max(0, int(settings.CELERY_STARTUP_RECOVERY_DELAY_SECONDS or 30))
|
|
startup_recovery_once.apply_async(
|
|
countdown=countdown,
|
|
queue=RECOVERY_QUEUE,
|
|
priority=settings.DOWNLOAD_TASK_PRIORITY_RECOVER,
|
|
)
|
|
# API v3 任务恢复(延迟 35 秒执行,避免与其他恢复任务冲突)
|
|
api_generation_recover_tasks_once.apply_async(
|
|
countdown=countdown + 5,
|
|
queue=RECOVERY_QUEUE,
|
|
priority=settings.DOWNLOAD_TASK_PRIORITY_RECOVER,
|
|
)
|
|
logger.info(
|
|
"启动容灾恢复协调任务已投递(含 API v3)。queue=%s countdown=%s",
|
|
RECOVERY_QUEUE,
|
|
countdown,
|
|
)
|
|
except Exception:
|
|
logger.exception("启动容灾恢复协调任务投递失败")
|
|
|
|
|
|
@heartbeat_sent.connect
|
|
def on_worker_heartbeat_sent(sender=None, **kwargs):
|
|
"""刷新主实例 heartbeat,并低频扫描已到期的旧主实例。"""
|
|
try:
|
|
from app.services.celery_runtime.worker_service import (
|
|
claim_worker_heartbeat_slot,
|
|
claim_worker_stale_scan_slot,
|
|
heartbeat_current_worker_instance,
|
|
registered_worker_identity,
|
|
)
|
|
|
|
heartbeat_ok = True
|
|
if claim_worker_heartbeat_slot():
|
|
heartbeat_ok = bool(run_async(heartbeat_current_worker_instance()))
|
|
if not heartbeat_ok or not claim_worker_stale_scan_slot():
|
|
return
|
|
|
|
identity = registered_worker_identity()
|
|
if identity is None:
|
|
return
|
|
|
|
from app.services.celery_runtime.recovery_service import (
|
|
mark_stale_worker_instance_candidates,
|
|
)
|
|
|
|
run_async(
|
|
mark_stale_worker_instance_candidates(
|
|
worker_name=identity.worker_name,
|
|
current_worker_instance_id=identity.worker_instance_id,
|
|
supports_targeted_recovery=identity.supports_targeted_recovery,
|
|
emit_duplicate_log=False,
|
|
)
|
|
)
|
|
except Exception:
|
|
logger.warning("Celery Worker 主实例 heartbeat/旧实例扫描失败", exc_info=True)
|
|
|
|
|
|
@worker_shutdown.connect
|
|
def on_worker_shutdown(sender=None, **kwargs):
|
|
"""优雅退出时撤销活跃 Worker key;旧任务集合保留给恢复流程。"""
|
|
try:
|
|
from app.services.celery_runtime.worker_service import unregister_current_worker_instance
|
|
|
|
run_async(unregister_current_worker_instance())
|
|
except Exception:
|
|
logger.debug("Celery Worker 主实例注销失败", exc_info=True)
|
|
finally:
|
|
close_loop()
|
|
|
|
|
|
@worker_process_init.connect
|
|
def on_worker_process_init(**kwargs):
|
|
"""prefork 子进程保留主 token,同时重建执行进程缓存和异步连接。"""
|
|
try:
|
|
from app.services.celery_runtime.worker_service import reset_process_identity_cache
|
|
|
|
reset_process_identity_cache()
|
|
except Exception:
|
|
pass
|
|
|
|
try:
|
|
run_async(engine.dispose())
|
|
except Exception:
|
|
pass
|
|
|
|
|
|
@worker_process_shutdown.connect
|
|
def on_worker_process_shutdown(**kwargs):
|
|
"""子进程退出前关闭连接池、Redis 注册表连接和 event loop。"""
|
|
try:
|
|
run_async(engine.dispose())
|
|
except Exception:
|
|
pass
|
|
|
|
try:
|
|
from app.services.redis_registry_service import close_registry_redis
|
|
|
|
run_async(close_registry_redis())
|
|
except Exception:
|
|
pass
|
|
finally:
|
|
close_loop()
|