648 lines
27 KiB
Python
648 lines
27 KiB
Python
import asyncio
|
||
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.payment_tasks",
|
||
"app.tasks.api_generation_tasks",
|
||
"app.tasks.api_recovery_tasks",
|
||
"app.tasks.api_upscale_tasks",
|
||
"app.tasks.scheduled_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["payment-expire-pending-orders-every-minute"] = {
|
||
"task": CeleryTaskName.PAYMENT_EXPIRE_PENDING.value,
|
||
"schedule": 60,
|
||
"options": {
|
||
"queue": RECOVERY_QUEUE,
|
||
"priority": settings.DOWNLOAD_TASK_PRIORITY_RECOVER,
|
||
},
|
||
}
|
||
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},
|
||
CeleryTaskName.PAYMENT_EXPIRE_PENDING.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.PAYMENT_EXPIRE_PENDING.value: {"queue": RECOVERY_QUEUE},
|
||
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
|
||
|
||
|
||
def _parse_schedule_to_celery(schedule_str: str):
|
||
"""将 schedule 字符串解析为 Celery 可识别的调度值。
|
||
|
||
- 纯数字:视为间隔秒数(返回 int)
|
||
- cron 表达式 (5 字段空格分隔):返回 crontab 对象
|
||
"""
|
||
from celery.schedules import crontab
|
||
|
||
s = (schedule_str or "").strip()
|
||
if not s:
|
||
return None
|
||
# 纯数字 → 间隔秒数
|
||
if s.isdigit():
|
||
return int(s)
|
||
# cron 表达式 (分 时 日 月 周)
|
||
parts = s.split()
|
||
if len(parts) == 5:
|
||
try:
|
||
return crontab(
|
||
minute=parts[0],
|
||
hour=parts[1],
|
||
day_of_month=parts[2],
|
||
month_of_year=parts[3],
|
||
day_of_week=parts[4],
|
||
)
|
||
except Exception:
|
||
logger.exception("解析 cron 表达式失败: %s", s)
|
||
return None
|
||
logger.warning("无法解析 schedule 表达式: %s", s)
|
||
return None
|
||
|
||
|
||
@celery_app.on_after_configure.connect # type: ignore
|
||
def _setup_dynamic_beat_tasks(sender, **kwargs):
|
||
"""从数据库加载活跃定时任务并注册到 Beat 调度。
|
||
|
||
通过 @celery_app.on_after_configure.connect 在 Celery 配置完成后执行,
|
||
适用于 Worker 和 Beat 启动场景。
|
||
|
||
注意:此信号运行在 Beat/Worker 主线程,不能使用 run_async(single_loop 在另一线程),
|
||
这里用独立 event loop 同步运行 async 查询,避免 asyncpg Future 跨 loop 错误。
|
||
"""
|
||
if celery_app is None:
|
||
return
|
||
|
||
async def _load():
|
||
from sqlalchemy import select
|
||
|
||
from app.models.base import async_session
|
||
from app.models.scheduled_task import ScheduledTask
|
||
|
||
async with async_session() as db:
|
||
result = await db.execute(
|
||
select(ScheduledTask).where(ScheduledTask.is_active.is_(True))
|
||
)
|
||
return result.scalars().all()
|
||
|
||
active_tasks = None
|
||
try:
|
||
# 在当前线程创建独立 loop 同步运行,不经过 single_loop
|
||
loop = asyncio.new_event_loop()
|
||
try:
|
||
active_tasks = loop.run_until_complete(_load())
|
||
finally:
|
||
loop.close()
|
||
except Exception:
|
||
logger.exception("加载定时任务失败,跳过动态 Beat 注册")
|
||
return
|
||
|
||
for task in active_tasks or []:
|
||
schedule_val = _parse_schedule_to_celery(task.schedule)
|
||
if schedule_val is None:
|
||
logger.warning("定时任务 %s schedule 无效,跳过注册: %s", task.id, task.schedule)
|
||
continue
|
||
beat_key = f"dynamic-scheduled-task-{task.id}"
|
||
sender.conf.beat_schedule[beat_key] = {
|
||
"task": "execute_scheduled_task",
|
||
"schedule": schedule_val,
|
||
"args": (task.id,),
|
||
"options": {"queue": RECOVERY_QUEUE},
|
||
}
|
||
logger.info(
|
||
"动态注册定时任务到 Beat: %s (%s) schedule=%s",
|
||
task.name, task.id, task.schedule,
|
||
)
|
||
|
||
|
||
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()
|