Files
wehub-resource-sync 41b710f9c7
CI / Frontend checks (push) Failing after 0s
CI / Backend tests (push) Failing after 1s
I18n Documentation Sync / sync-docs (push) Failing after 0s
chore: import upstream snapshot with attribution
2026-07-13 12:28:40 +08:00

395 lines
14 KiB
Python

"""视频处理Celery任务
包含WebSocket实时通知和Pipeline适配器集成
"""
import os
import logging
import asyncio
from typing import Dict, Any, Optional
from celery import current_task
from pathlib import Path
from backend.core.celery_app import celery_app
from backend.services.websocket_notification_service import notification_service
from backend.services.processing_service import ProcessingService
from backend.services.pipeline_adapter import create_pipeline_adapter
from backend.core.database import SessionLocal
from backend.models.project import Project, ProjectStatus
from backend.models.task import Task, TaskStatus, TaskType
from datetime import datetime
logger = logging.getLogger(__name__)
def run_async_notification(coro):
"""运行异步通知的辅助函数 - 修复事件循环冲突"""
try:
# 尝试获取现有的事件循环
loop = asyncio.get_event_loop()
if loop.is_running():
# 如果事件循环正在运行,使用线程池执行
import concurrent.futures
import threading
def run_in_thread():
new_loop = asyncio.new_event_loop()
asyncio.set_event_loop(new_loop)
try:
return new_loop.run_until_complete(coro)
finally:
new_loop.close()
with concurrent.futures.ThreadPoolExecutor() as executor:
future = executor.submit(run_in_thread)
return future.result(timeout=10) # 10秒超时
else:
# 如果事件循环没有运行,直接运行
return loop.run_until_complete(coro)
except RuntimeError:
# 没有事件循环,创建新的
loop = asyncio.new_event_loop()
asyncio.set_event_loop(loop)
try:
return loop.run_until_complete(coro)
finally:
loop.close()
# 同一项目同一时间只允许一条流水线在跑。桌面模式下多个入口(/process、/retry、
# 前端自动启动、下载后自动启动)可能并发派发同一项目,过去任务都进 Redis 没人跑
# 所以看不出来;现在任务真的会本地执行,必须防止并发重复执行撞数据库。
import threading as _threading
_active_pipeline_projects: set = set()
_active_pipeline_lock = _threading.Lock()
@celery_app.task(bind=True, name='backend.tasks.processing.process_video_pipeline')
def process_video_pipeline(
self,
project_id: str,
input_video_path: str,
input_srt_path: Optional[str] = None,
) -> Dict[str, Any]:
"""
处理视频流水线任务 - 使用Pipeline适配器
Args:
project_id: 项目ID
input_video_path: 输入视频路径
input_srt_path: 输入SRT路径
Returns:
处理结果
"""
task_id = self.request.id
logger.info(f"开始处理视频流水线: {project_id}, 任务ID: {task_id}")
# 并发去重:同一项目已有流水线在跑就直接跳过本次重复派发
with _active_pipeline_lock:
if project_id in _active_pipeline_projects:
logger.warning(f"项目 {project_id} 已有流水线在运行,跳过重复任务 {task_id}")
return {
"success": False,
"skipped": True,
"project_id": project_id,
"task_id": task_id,
"message": "已有流水线在运行,已跳过重复任务",
}
_active_pipeline_projects.add(project_id)
try:
# 创建数据库会话
db = SessionLocal()
try:
# 创建任务记录
task = Task(
name=f"视频处理流水线",
description=f"处理项目 {project_id} 的完整视频流水线",
task_type=TaskType.VIDEO_PROCESSING,
project_id=project_id,
celery_task_id=task_id,
status=TaskStatus.RUNNING,
progress=0,
current_step="初始化",
total_steps=6
)
db.add(task)
db.commit()
# 发送开始通知
run_async_notification(
notification_service.send_processing_start(project_id, task_id)
)
# 简化的进度系统不需要复杂的回调函数
# 新的进度系统会在流水线内部自动发送进度事件
# 使用简化的Pipeline适配器
from backend.services.simple_pipeline_adapter import create_simple_pipeline_adapter
pipeline_adapter = create_simple_pipeline_adapter(str(project_id), str(task.id))
# 执行Pipeline处理 - 使用异步包装器
import asyncio
result = asyncio.run(pipeline_adapter.process_project_sync(input_video_path, input_srt_path))
# 检查处理结果
if result.get("status") == "failed":
# 处理失败
error_msg = result.get("message", "处理失败")
task.status = TaskStatus.FAILED
task.error_message = error_msg
task.result_data = result
# 更新项目状态为失败
project = db.query(Project).filter(Project.id == project_id).first()
if project:
project.status = ProjectStatus.FAILED
project.updated_at = datetime.utcnow()
logger.info(f"项目状态已更新为失败: {project_id}")
db.commit()
# 失败状态已由简化进度系统自动处理
# 发送错误通知(兼容旧版本) - 已禁用WebSocket通知
# run_async_notification(
# notification_service.send_processing_error(project_id, task_id, error_msg)
# )
return {
"success": False,
"project_id": project_id,
"task_id": task_id,
"error": error_msg,
"result": result
}
else:
# 处理成功
task.status = TaskStatus.COMPLETED
task.progress = 100
task.current_step = "处理完成"
task.result_data = result
# 更新项目状态为已完成
project = db.query(Project).filter(Project.id == project_id).first()
if project:
project.status = ProjectStatus.COMPLETED
project.completed_at = datetime.utcnow()
project.updated_at = datetime.utcnow()
logger.info(f"项目状态已更新为已完成: {project_id}")
db.commit()
# 完成状态已由简化进度系统自动处理
# 发送完成通知(兼容旧版本) - 已禁用WebSocket通知
# run_async_notification(
# notification_service.send_processing_complete(project_id, task_id, result)
# )
logger.info(f"视频流水线处理完成: {project_id}")
return {
"success": True,
"project_id": project_id,
"task_id": task_id,
"result": result,
"message": "视频处理流水线完成"
}
finally:
db.close()
# 释放并发去重锁(无论成功/失败/提前返回都会走到这里)
with _active_pipeline_lock:
_active_pipeline_projects.discard(project_id)
except Exception as e:
error_msg = f"视频流水线处理失败: {str(e)}"
logger.error(error_msg)
# 兜底释放并发锁(正常路径已在内层 finally 释放,这里防止极早期异常泄漏)
with _active_pipeline_lock:
_active_pipeline_projects.discard(project_id)
# 更新任务状态为失败
try:
db = SessionLocal()
task = db.query(Task).filter(Task.celery_task_id == task_id).first()
if task:
task.status = TaskStatus.FAILED
task.error_message = error_msg
# 更新项目状态为失败
project = db.query(Project).filter(Project.id == project_id).first()
if project:
project.status = ProjectStatus.FAILED
project.updated_at = datetime.utcnow()
logger.info(f"项目状态已更新为失败: {project_id}")
db.commit()
db.close()
except Exception as db_error:
logger.error(f"更新任务状态失败: {str(db_error)}")
# 发送错误通知
run_async_notification(
notification_service.send_processing_error(project_id, task_id, error_msg)
)
raise
@celery_app.task(bind=True, name='backend.tasks.processing.process_single_step')
def process_single_step(self, project_id: str, step: str, config: Dict[str, Any]) -> Dict[str, Any]:
"""
处理单个步骤任务
Args:
project_id: 项目ID
step: 步骤名称
config: 处理配置
Returns:
处理结果
"""
task_id = self.request.id
logger.info(f"开始处理单个步骤: {project_id}, 步骤: {step}, 任务ID: {task_id}")
try:
# 发送开始通知
# 发送处理开始通知(兼容旧版本) - 已禁用WebSocket通知
# run_async_notification(
# notification_service.send_processing_start(project_id, task_id)
# )
# 创建数据库会话
db = SessionLocal()
try:
# 创建处理服务
processing_service = ProcessingService(db)
# 根据步骤类型执行不同的处理
if step == "outline":
run_async_notification(
notification_service.send_processing_progress(project_id, task_id, 50, "生成大纲")
)
result = processing_service.generate_outline(project_id, config)
elif step == "timeline":
run_async_notification(
notification_service.send_processing_progress(project_id, task_id, 50, "提取时间轴")
)
result = processing_service.extract_timeline(project_id, config)
elif step == "titles":
run_async_notification(
notification_service.send_processing_progress(project_id, task_id, 50, "生成标题")
)
result = processing_service.generate_titles(project_id, config)
elif step == "clips":
run_async_notification(
notification_service.send_processing_progress(project_id, task_id, 50, "视频切片")
)
result = processing_service.extract_clips(project_id, config)
elif step == "collections":
run_async_notification(
notification_service.send_processing_progress(project_id, task_id, 50, "生成合集")
)
result = processing_service.generate_collections(project_id, config)
else:
raise Exception(f"未知的步骤类型: {step}")
if not result.get("success"):
raise Exception(f"步骤 {step} 处理失败: {result.get('error')}")
# 发送完成通知
run_async_notification(
notification_service.send_processing_complete(project_id, task_id, result)
)
logger.info(f"单个步骤处理完成: {project_id}, 步骤: {step}")
return result
finally:
db.close()
except Exception as e:
error_msg = f"单个步骤处理失败: {str(e)}"
logger.error(error_msg)
# 发送错误通知
run_async_notification(
notification_service.send_processing_error(project_id, task_id, error_msg)
)
raise
@celery_app.task(bind=True, name='backend.tasks.processing.retry_processing_step')
def retry_processing_step(self, project_id: str, step: str, config: Dict[str, Any],
original_task_id: str) -> Dict[str, Any]:
"""
重试处理步骤任务
Args:
project_id: 项目ID
step: 步骤名称
config: 处理配置
original_task_id: 原始任务ID
Returns:
处理结果
"""
task_id = self.request.id
logger.info(f"开始重试处理步骤: {project_id}, 步骤: {step}, 任务ID: {task_id}")
try:
# 发送开始通知
# 发送处理开始通知(兼容旧版本) - 已禁用WebSocket通知
# run_async_notification(
# notification_service.send_processing_start(project_id, task_id)
# )
# 发送重试通知
run_async_notification(
notification_service.send_system_notification(
"retry_started",
"重试开始",
f"正在重试步骤: {step}",
"warning"
)
)
# 调用单个步骤处理
result = process_single_step.apply_async(
args=[project_id, step, config],
task_id=task_id
).get()
# 发送重试成功通知
run_async_notification(
notification_service.send_system_notification(
"retry_success",
"重试成功",
f"步骤 {step} 重试成功",
"success"
)
)
return result
except Exception as e:
error_msg = f"重试处理步骤失败: {str(e)}"
logger.error(error_msg)
# 发送重试失败通知
run_async_notification(
notification_service.send_error_notification(
"retry_failed",
f"步骤 {step} 重试失败",
{"project_id": project_id, "step": step, "error": str(e)}
)
)
raise