Files
xiaoxia-saas/apps/worker/worker_app/tasks/compose_video.py
T
xiaoxia 8427bb6852
CI/CD Pipeline / Validate Code Quality And Tests (push) Failing after 27s
CI/CD Pipeline / Build Production Runtime Images (push) Failing after 1579h39m58s
CI/CD Pipeline / Deploy Production (push) Failing after 1579h39m54s
CI/CD Pipeline / Build & Push Staging (Watchtower auto-deploy) (push) Failing after 1579h39m58s
CI/CD Pipeline / Integration Tests (push) Has been skipped
CI/CD Pipeline / Frontend Lint (push) Has been skipped
CI/CD Pipeline / Production Browser E2E (push) Failing after 1580h11m26s
CI/CD Pipeline / Staging API Integration Tests (push) Failing after 1580h11m30s
CI/CD Pipeline / Staging E2E Tests (push) Failing after 1580h11m30s
feat(unified-render): Phase 2 - Feature Flag + RenderAdapter 适配层 (#231)
## 统一渲染引擎 Phase 2 - Feature Flag 开关 + 适配层

### 变更内容

1. **Feature Flag 开关**
   - API 端 `Settings.RENDER_ENGINE`(默认 `legacy`)
   - Worker 端 `WorkerSettings.render_engine`(默认 `legacy`)
   - 支持环境变量 `RENDER_ENGINE=unified` 一键切换

2. **RenderAdapter 适配层** (`apps/worker/video_processing/render_adapter.py`)
   - EditPlan + EditPlanClips → UnifiedRenderService 输入的完整适配
   - 素材自动下载(OSS → 本地路径映射)
   - 进度回调对接(ProgressCallback)
   - 结果自动上传 OSS
   - `validate_plan()` 兼容旧接口,便于灰度切换

3. **compose_video 任务改造**
   - 根据 `RENDER_ENGINE` 配置分流到 legacy / unified 路径
   - 新引擎结果回写字段对齐(engine/width/height/file_size)
   - 错误处理与重试逻辑保持一致

### 测试
- 新增 16 个单元测试:validate_plan(7) + render_plan(6) + download_assets(3)
- 原有 52 个 unified_render_service 测试全部通过
- 合计 **68 个测试全绿**

### 灰度策略
- 默认 `legacy`,不影响现有功能
- 灰度时设置环境变量 `RENDER_ENGINE=unified` 即可切换
- 后续可支持按 plan_id / user_id 灰度(Phase 3)

---------

Co-authored-by: 灵应 <lingying@coze.email>
Reviewed-on: #231
2026-07-12 19:43:08 +08:00

223 lines
7.8 KiB
Python
Executable File
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""视频合成 Celery 任务 — Phase 8 任务 2.10.
使用 JobService 管理任务生命周期,集成 VideoComposeService 执行合成。
"""
from __future__ import annotations
import logging
import os
import shutil
import subprocess
import tempfile
from datetime import datetime, timezone
from pathlib import Path
from celery.utils.log import get_task_logger
from worker_app.celery_app import celery_app
from worker_app.db import SessionLocal
logger = get_task_logger(__name__)
def _get_job_service():
"""延迟导入 JobService,避免循环依赖。"""
from apps.api.app.services.job_service import JobService
from packages.adapters.sqlalchemy_impl.job_repository import SQLAlchemyJobRepository
db = SessionLocal()
repo = SQLAlchemyJobRepository(db)
return JobService(repo), db
@celery_app.task(
name="worker.compose_video",
bind=True,
max_retries=3,
default_retry_delay=60,
)
def compose_video(self, job_id: str, **kwargs):
"""视频合成任务。
根据 RENDER_ENGINE 配置选择渲染引擎:
- legacy: 旧 VideoComposeService(filter_complex 模式)
- unified: 新 UnifiedRenderService(图层架构)
Args:
job_id: JobService 中的任务 ID
**kwargs: 来自 Job.payload 的额外参数(plan_id, output_path 等)
"""
job_service, db = _get_job_service()
try:
job = job_service.get_job(job_id)
if job is None:
logger.error("Job not found: %s", job_id)
return {"status": "error", "message": f"Job {job_id} not found"}
plan_id = job.payload.get("plan_id", "")
if not plan_id:
job_service.fail_job(job_id, "Missing plan_id in job payload")
return {"status": "error", "message": "Missing plan_id"}
# 判断使用哪个渲染引擎
from worker_app.core.config import get_settings as get_worker_settings
worker_settings = get_worker_settings()
engine = (worker_settings.render_engine or "legacy").lower()
if engine == "unified":
return _compose_with_unified_engine(self, job_service, job, plan_id, db)
else:
return _compose_with_legacy_engine(self, job_service, job, plan_id, db)
except self.retry_exc as exc:
logger.warning("视频合成重试中: job_id=%s, exc=%s", job_id, exc)
raise
except Exception as exc:
logger.exception("视频合成异常: job_id=%s", job_id)
try:
job_service.fail_job(job_id, str(exc)[:500])
except Exception:
logger.exception("更新 Job 失败状态时出错")
raise self.retry(exc=exc, countdown=60)
finally:
db.close()
def _compose_with_legacy_engine(task, job_service, job, plan_id: str, db) -> dict:
"""旧引擎渲染路径(VideoComposeService)。"""
job_id = job.id
# 标记为 running
job_service.update_progress(job_id, progress=10.0, current_stage="初始化合成环境")
# 延迟导入 VideoComposeService
from apps.api.app.services.video_compose_service import VideoComposeService
compose_svc = VideoComposeService(db)
# 校验合成条件
job_service.update_progress(job_id, progress=20.0, current_stage="校验合成条件")
validation = compose_svc.validate_compose(plan_id)
if not validation.valid:
error_msg = "; ".join(validation.errors)
job_service.fail_job(job_id, f"合成校验失败: {error_msg}")
return {"status": "error", "message": error_msg}
# 构建合成命令
job_service.update_progress(job_id, progress=30.0, current_stage="构建 FFmpeg 命令")
_output_dir = os.environ.get("VIDEO_OUTPUT_DIR", os.path.join(tempfile.gettempdir(), "video_output"))
output_path = os.path.join(_output_dir, f"{job_id}.mp4")
compose_cmd = compose_svc.build_compose_command(plan_id, output_path)
# 执行 FFmpeg
job_service.update_progress(job_id, progress=50.0, current_stage="正在执行视频合成")
logger.info("Executing FFmpeg for job %s, plan %s", job_id, plan_id)
try:
subprocess.run(
compose_cmd.command,
check=True,
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
text=True,
timeout=3600,
)
except subprocess.CalledProcessError as e:
job_service.fail_job(job_id, f"FFmpeg 执行失败: {e.stderr[:500]}")
raise
# 上传结果
job_service.update_progress(job_id, progress=80.0, current_stage="上传合成结果")
storage_key = f"rendered/{plan_id}/{job_id}.mp4"
from worker_app.tasks.edit_plan_generation import _upload_to_oss
output_url = _upload_to_oss(Path(output_path), storage_key)
# 更新 Job 状态为完成
result_data = {
"plan_id": plan_id,
"output_path": output_path,
"storage_key": storage_key,
"output_url": output_url or "",
"estimated_duration": compose_cmd.estimated_duration,
"clip_count": len(compose_cmd.clip_chains),
"engine": "legacy",
}
job_service.complete_job(job_id, result=result_data)
logger.info("视频合成完成(legacy): job_id=%s, plan_id=%s", job_id, plan_id)
return {"status": "completed", "job_id": job_id, "result": result_data}
def _compose_with_unified_engine(task, job_service, job, plan_id: str, db) -> dict:
"""新引擎渲染路径(UnifiedRenderService + RenderAdapter)。"""
job_id = job.id
# 标记为 running
job_service.update_progress(job_id, progress=10.0, current_stage="初始化统一渲染引擎")
from video_processing.render_adapter import RenderAdapter
adapter = RenderAdapter(db)
# 校验合成条件
job_service.update_progress(job_id, progress=15.0, current_stage="校验合成条件")
valid, errors, warnings, ready_count, total_count = adapter.validate_plan(plan_id)
if not valid:
error_msg = "; ".join(errors)
job_service.fail_job(job_id, f"合成校验失败: {error_msg}")
return {"status": "error", "message": error_msg}
# 进度回调
def progress_cb(progress: float, stage: str) -> None:
try:
job_service.update_progress(job_id, progress=progress, current_stage=stage)
except Exception:
logger.exception("更新进度失败")
# 执行渲染
job_service.update_progress(job_id, progress=20.0, current_stage="开始渲染")
logger.info("统一渲染引擎开始: job_id=%s plan_id=%s", job_id, plan_id)
result = adapter.render_plan(
plan_id=plan_id,
job_id=job_id,
progress_cb=progress_cb,
)
if not result.success:
job_service.fail_job(job_id, f"渲染失败: {result.error_message}")
raise RuntimeError(result.error_message)
# 更新 Job 状态为完成
result_data = {
"plan_id": plan_id,
"output_path": str(result.output_path) if result.output_path else "",
"storage_key": f"rendered/{plan_id}/{job_id}.mp4",
"output_url": result.output_url,
"estimated_duration": result.duration,
"clip_count": result.clip_count,
"engine": "unified",
"width": result.width,
"height": result.height,
"file_size": result.file_size,
}
job_service.complete_job(job_id, result=result_data)
logger.info("视频合成完成(unified): job_id=%s plan_id=%s duration=%.2fs", job_id, plan_id, result.duration)
return {"status": "completed", "job_id": job_id, "result": result_data}
def _cleanup_output(job_id: str) -> None:
"""清理临时输出文件。"""
try:
_output_dir = os.environ.get("VIDEO_OUTPUT_DIR", os.path.join(tempfile.gettempdir(), "video_output"))
output_path = os.path.join(_output_dir, f"{job_id}.mp4")
if Path(output_path).exists():
Path(output_path).unlink()
except Exception as e:
logger.warning(f"清理输出文件失败: {e}", exc_info=True)