"""视频处理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