修复celery异步输出json序列化异常BUG

This commit is contained in:
2026-06-26 13:59:58 +08:00
parent 479e1a91fb
commit 9bf181e14c
3 changed files with 41 additions and 23 deletions
+15 -1
View File
@@ -58,6 +58,20 @@ if broker_url:
task_acks_late=True,
task_reject_on_worker_lost=True,
task_track_started=True,
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},
"shot_replicate.analyze_custom_segment_video": {"ignore_result": True},
"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},
},
worker_prefetch_multiplier=1,
broker_transport_options={
"visibility_timeout": 3600,
@@ -168,4 +182,4 @@ def on_worker_process_shutdown(**kwargs):
except Exception:
pass
finally:
close_loop()
close_loop()
@@ -21,7 +21,7 @@ from app.tasks.celery_app import celery_app
MODULE = ModuleCodeEnum.HOT_OPENING_REPLICATE.value
async def _run_image_prompt(project_id: str, step_id: str | None = None):
async def _run_image_prompt(project_id: str, step_id: str | None = None) -> None:
lock_token: str | None = None
if step_id:
lock_token = await acquire_object_lock(object_type=OBJECT_MODULE_STEP, object_id=step_id)
@@ -38,17 +38,17 @@ async def _run_image_prompt(project_id: str, step_id: str | None = None):
await mark_active_started(object_type=OBJECT_MODULE_STEP, object_id=step_id)
try:
async with async_session() as db:
result = await run_image_prompt_optimize(db, project_id=project_id, step_id=step_id)
await run_image_prompt_optimize(db, project_id=project_id, step_id=step_id)
await db.commit()
if step_id:
await cleanup_active_if_terminal(db, object_type=OBJECT_MODULE_STEP, object_id=step_id)
return result
return None
finally:
if step_id:
await release_object_lock(object_type=OBJECT_MODULE_STEP, object_id=step_id, token=lock_token)
async def _run_video_prompt(project_id: str, step_id: str | None = None):
async def _run_video_prompt(project_id: str, step_id: str | None = None) -> None:
lock_token: str | None = None
if step_id:
lock_token = await acquire_object_lock(object_type=OBJECT_MODULE_STEP, object_id=step_id)
@@ -64,18 +64,18 @@ async def _run_video_prompt(project_id: str, step_id: str | None = None):
await mark_active_started(object_type=OBJECT_MODULE_STEP, object_id=step_id)
try:
async with async_session() as db:
result = await run_video_prompt_optimize(db, project_id=project_id, step_id=step_id)
await run_video_prompt_optimize(db, project_id=project_id, step_id=step_id)
await db.commit()
if step_id:
await cleanup_active_if_terminal(db, object_type=OBJECT_MODULE_STEP, object_id=step_id)
return result
return None
finally:
if step_id:
await release_object_lock(object_type=OBJECT_MODULE_STEP, object_id=step_id, token=lock_token)
if celery_app:
@celery_app.task(name="hot_opening.start_image_prompt_optimize", bind=True, max_retries=3, default_retry_delay=30)
@celery_app.task(name="hot_opening.start_image_prompt_optimize", bind=True, max_retries=3, default_retry_delay=30, ignore_result=True)
def start_image_prompt_optimize(self, project_id: str, step_id: str | None = None):
"""手动触发后的图片 AI 提词任务。
@@ -84,11 +84,12 @@ if celery_app:
service 内部已经落库为业务失败的情况不会抛出异常,也不会重复 retry。
"""
try:
return run_async(_run_image_prompt(project_id, step_id))
run_async(_run_image_prompt(project_id, step_id))
return None
except Exception as exc:
raise self.retry(exc=exc) from exc
@celery_app.task(name="hot_opening.start_video_prompt_optimize", bind=True, max_retries=3, default_retry_delay=30)
@celery_app.task(name="hot_opening.start_video_prompt_optimize", bind=True, max_retries=3, default_retry_delay=30, ignore_result=True)
def start_video_prompt_optimize(self, project_id: str, step_id: str | None = None):
"""手动触发后的视频 AI 提词任务。
@@ -97,7 +98,8 @@ if celery_app:
service 内部已经落库为业务失败的情况不会抛出异常,也不会重复 retry。
"""
try:
return run_async(_run_video_prompt(project_id, step_id))
run_async(_run_video_prompt(project_id, step_id))
return None
except Exception as exc:
raise self.retry(exc=exc) from exc
else:
@@ -109,4 +111,4 @@ else:
raise RuntimeError("Celery is disabled")
start_image_prompt_optimize = _DisabledTask()
start_video_prompt_optimize = _DisabledTask()
start_video_prompt_optimize = _DisabledTask()
@@ -21,7 +21,7 @@ from app.tasks.celery_app import celery_app
MODULE = ModuleCodeEnum.SHOT_REPLICATE.value
async def _run_image_prompt(project_id: str, step_id: str | None = None):
async def _run_image_prompt(project_id: str, step_id: str | None = None) -> None:
lock_token: str | None = None
if step_id:
lock_token = await acquire_object_lock(object_type=OBJECT_MODULE_STEP, object_id=step_id)
@@ -37,17 +37,17 @@ async def _run_image_prompt(project_id: str, step_id: str | None = None):
await mark_active_started(object_type=OBJECT_MODULE_STEP, object_id=step_id)
try:
async with async_session() as db:
result = await run_image_prompt_optimize(db, project_id=project_id, step_id=step_id)
await run_image_prompt_optimize(db, project_id=project_id, step_id=step_id)
await db.commit()
if step_id:
await cleanup_active_if_terminal(db, object_type=OBJECT_MODULE_STEP, object_id=step_id)
return result
return None
finally:
if step_id:
await release_object_lock(object_type=OBJECT_MODULE_STEP, object_id=step_id, token=lock_token)
async def _run_video_prompt(project_id: str, step_id: str | None = None):
async def _run_video_prompt(project_id: str, step_id: str | None = None) -> None:
lock_token: str | None = None
if step_id:
lock_token = await acquire_object_lock(object_type=OBJECT_MODULE_STEP, object_id=step_id)
@@ -63,11 +63,11 @@ async def _run_video_prompt(project_id: str, step_id: str | None = None):
await mark_active_started(object_type=OBJECT_MODULE_STEP, object_id=step_id)
try:
async with async_session() as db:
result = await run_video_prompt_optimize(db, project_id=project_id, step_id=step_id)
await run_video_prompt_optimize(db, project_id=project_id, step_id=step_id)
await db.commit()
if step_id:
await cleanup_active_if_terminal(db, object_type=OBJECT_MODULE_STEP, object_id=step_id)
return result
return None
finally:
if step_id:
await release_object_lock(object_type=OBJECT_MODULE_STEP, object_id=step_id, token=lock_token)
@@ -75,7 +75,7 @@ async def _run_video_prompt(project_id: str, step_id: str | None = None):
if celery_app:
@celery_app.task(name="shot_replicate.start_image_prompt_optimize", bind=True, max_retries=3, default_retry_delay=30)
@celery_app.task(name="shot_replicate.start_image_prompt_optimize", bind=True, max_retries=3, default_retry_delay=30, ignore_result=True)
def start_image_prompt_optimize(self, project_id: str, step_id: str | None = None):
"""手动触发后的图片 AI 提词任务。
@@ -84,12 +84,13 @@ if celery_app:
service 内部已经落库为业务失败的情况不会抛出异常,也不会重复 retry。
"""
try:
return run_async(_run_image_prompt(project_id, step_id))
run_async(_run_image_prompt(project_id, step_id))
return None
except Exception as exc:
raise self.retry(exc=exc) from exc
@celery_app.task(name="shot_replicate.start_video_prompt_optimize", bind=True, max_retries=3, default_retry_delay=30)
@celery_app.task(name="shot_replicate.start_video_prompt_optimize", bind=True, max_retries=3, default_retry_delay=30, ignore_result=True)
def start_video_prompt_optimize(self, project_id: str, step_id: str | None = None):
"""手动触发后的视频 AI 提词任务。
@@ -98,7 +99,8 @@ if celery_app:
service 内部已经落库为业务失败的情况不会抛出异常,也不会重复 retry。
"""
try:
return run_async(_run_video_prompt(project_id, step_id))
run_async(_run_video_prompt(project_id, step_id))
return None
except Exception as exc:
raise self.retry(exc=exc) from exc
@@ -112,4 +114,4 @@ else:
raise RuntimeError("Celery is disabled")
start_image_prompt_optimize = _DisabledTask()
start_video_prompt_optimize = _DisabledTask()
start_video_prompt_optimize = _DisabledTask()