125 lines
3.9 KiB
Python
125 lines
3.9 KiB
Python
import logging
|
|
|
|
from celery import Celery
|
|
from celery.signals import worker_process_init, worker_process_shutdown, worker_ready
|
|
|
|
from app.config import settings
|
|
from app.models.base import engine
|
|
from app.tasks.async_runner import close_loop, run_async
|
|
|
|
logger = logging.getLogger("video_gen")
|
|
|
|
|
|
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}"
|
|
|
|
|
|
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")
|
|
celery_app.conf.update(
|
|
broker_url=broker_url,
|
|
result_backend=backend_url or broker_url,
|
|
task_serializer="json",
|
|
accept_content=["json"],
|
|
result_serializer="json",
|
|
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,
|
|
worker_prefetch_multiplier=1,
|
|
broker_transport_options={
|
|
"visibility_timeout": 3600,
|
|
"queue_order_strategy": "priority",
|
|
"priority_steps": list(range(10)),
|
|
"sep": ":",
|
|
},
|
|
task_routes={
|
|
"generation.chatapi_create_generation_task": {"queue": "gen_chatapi_create"},
|
|
"generation.poll_generation_task": {"queue": "gen_provider_poll"},
|
|
"generation.download_generation_result_task": {"queue": "gen_result_download"},
|
|
"generation.recover_download_tasks_once": {"queue": "gen_result_download"},
|
|
"generation.recover_generation_tasks_once": {"queue": "gen_result_download"},
|
|
"app.tasks.cleanup.*": {"queue": "default"},
|
|
},
|
|
)
|
|
celery_app.autodiscover_tasks(["app.tasks"])
|
|
else:
|
|
celery_app = None
|
|
|
|
|
|
@worker_ready.connect
|
|
def on_worker_ready(sender=None, **kwargs):
|
|
"""Celery worker 启动时做一次容灾恢复。
|
|
|
|
注意:
|
|
- 不启用 Celery beat。
|
|
- 不要求新增第四条启动命令。
|
|
- 只让 gen_result_download worker 投递恢复任务,避免三个 worker 同时重复扫描。
|
|
"""
|
|
if celery_app is None:
|
|
return
|
|
|
|
hostname = str(getattr(sender, "hostname", "") or "")
|
|
if "gen_result_download" not in hostname:
|
|
return
|
|
|
|
try:
|
|
from app.tasks.generation_recovery_tasks import (
|
|
recover_download_tasks_once,
|
|
recover_generation_tasks_once,
|
|
)
|
|
|
|
countdown = max(0, int(settings.DOWNLOAD_RECOVERY_STARTUP_DELAY_SECONDS or 0))
|
|
|
|
recover_generation_tasks_once.apply_async(
|
|
countdown=countdown,
|
|
queue="gen_result_download",
|
|
priority=settings.DOWNLOAD_TASK_PRIORITY_RECOVER,
|
|
)
|
|
recover_download_tasks_once.apply_async(
|
|
countdown=countdown + 5,
|
|
queue="gen_result_download",
|
|
priority=settings.DOWNLOAD_TASK_PRIORITY_RECOVER,
|
|
)
|
|
except Exception:
|
|
logger.exception("启动容灾恢复任务投递失败")
|
|
|
|
|
|
@worker_process_init.connect
|
|
def on_worker_process_init(**kwargs):
|
|
"""Linux prefork 子进程启动后丢弃 fork 前可能继承的连接池状态。"""
|
|
try:
|
|
run_async(engine.dispose())
|
|
except Exception:
|
|
pass
|
|
|
|
|
|
@worker_process_shutdown.connect
|
|
def on_worker_process_shutdown(**kwargs):
|
|
"""子进程退出前关闭连接池和 event loop。"""
|
|
try:
|
|
run_async(engine.dispose())
|
|
except Exception:
|
|
pass
|
|
|
|
try:
|
|
from app.services.celery_download_recovery_service import close_registry_redis
|
|
|
|
run_async(close_registry_redis())
|
|
except Exception:
|
|
pass
|
|
finally:
|
|
close_loop()
|