解决并发添加视频导致草稿异常的问题。

This commit is contained in:
Hommy
2026-03-31 19:36:16 +08:00
parent 728772c6a0
commit 8449cd2d71
7 changed files with 1268 additions and 34 deletions
+9 -6
View File
@@ -16,6 +16,7 @@ from src.schemas.easy_create_material import EasyCreateMaterialResponse
from src.schemas.save_draft import SaveDraftResponse
from src.schemas.create_draft import CreateDraftResponse
from fastapi import APIRouter, Request, Depends
import asyncio
from src.schemas.create_draft import CreateDraftRequest, CreateDraftResponse
from src.schemas.add_videos import AddVideosRequest, AddVideosResponse
from src.schemas.add_audios import AddAudiosRequest, AddAudiosResponse
@@ -86,13 +87,14 @@ def save_draft(sdr: SaveDraftRequest) -> SaveDraftResponse:
return SaveDraftResponse(draft_url=draft_url)
@router.post(path="/add_videos", response_model=AddVideosResponse)
def add_videos(avr: AddVideosRequest) -> AddVideosResponse:
async def add_videos(avr: AddVideosRequest) -> AddVideosResponse:
"""
向剪映草稿添加视频 (v1版本)
向剪映草稿添加视频 (v1 版本,带并发锁保护)
使用异步锁机制防止同一草稿的并发写操作导致文件损坏
"""
# 调用service层处理业务逻辑
draft_url, track_id, video_ids, segment_ids = service.add_videos(
# 调用 service 层处理业务逻辑(异步版本,带锁保护)
draft_url, track_id, video_ids, segment_ids = await service.add_videos_async(
draft_url=avr.draft_url,
video_infos=avr.video_infos,
scene_timelines=[{"start": t.start, "end": t.end} for t in avr.scene_timelines] if avr.scene_timelines else None,
@@ -100,7 +102,8 @@ def add_videos(avr: AddVideosRequest) -> AddVideosResponse:
scale_x=avr.scale_x,
scale_y=avr.scale_y,
transform_x=avr.transform_x,
transform_y=avr.transform_y
transform_y=avr.transform_y,
lock_timeout=30.0 # 30 秒超时
)
return AddVideosResponse(draft_url=draft_url, track_id=track_id, video_ids=video_ids, segment_ids=segment_ids)
+1 -1
View File
@@ -1,5 +1,5 @@
from .create_draft import create_draft
from .add_videos import add_videos
from .add_videos import add_videos, add_videos_async, _add_videos_internal
from .add_audios import add_audios
from .add_images import add_images
from .add_sticker import add_sticker
+150 -27
View File
@@ -1,6 +1,6 @@
from src.pyJianYingDraft.video_segment import VideoSegment
import asyncio
from src.utils.logger import logger
from src.pyJianYingDraft import ScriptFile, trange, IntroType
import src.pyJianYingDraft as draft
@@ -12,6 +12,7 @@ from src.utils.download import download
import config
import json
from typing import List, Dict, Any, Tuple, Optional
from src.utils.draft_lock_manager import get_draft_lock_manager
def add_videos(
@@ -25,38 +26,38 @@ def add_videos(
transform_y: int = 0
) -> Tuple[str, str, List[str], List[str]]:
"""
添加视频到剪映草稿的业务逻辑
添加视频到剪映草稿的业务逻辑(同步版本,兼容旧代码)
Args:
draft_url: "" // [必选] 草稿URL
draft_url: "" // [必选] 草稿 URL
video_infos: [
{
"video_url": "https://example.com/video1.mp4", // [必选] 视频文件的URL地址
"video_url": "https://example.com/video1.mp4", // [必选] 视频文件的 URL 地址
"width": 1920, // [可选] 视频宽度,不传则自动获取视频文件尺寸
"height": 1080, // [可选] 视频高度,不传则自动获取视频文件尺寸
"start": 0.0, // [必选] 视频在时间轴上的开始时间 (微秒)
"end": 12000000.0, // [必选] 视频在时间轴上的结束时间 (微秒)
"duration": 12000000.0, // [可选] 视频总时长(微秒),如果不传则默认为end-start
"mask": "", // 遮罩类型[可选],默认值为None
"transition": "", // 转场效果名称[可选],默认值为None
"transition_duration": 500000.0, // 转场持续时间(微秒)[可选],默认值为500000
"volume": 1.0, // 音量大小[0, 10][可选],默认值为1.0,10为最大音量
"duration": 12000000.0, // [可选] 视频总时长 (微秒),如果不传则默认为 end-start
"mask": "", // 遮罩类型 [可选],默认值为 None
"transition": "", // 转场效果名称 [可选],默认值为 None
"transition_duration": 500000.0, // 转场持续时间 (微秒)[可选],默认值为 500000
"volume": 1.0, // 音量大小 [0, 10][可选],默认值为 1.010 为最大音量
}
] // [必选]
scene_timelines: [ // [可选] 场景时间线数组,用于视频变速,与video_infos一一对应
scene_timelines: [ // [可选] 场景时间线数组,用于视频变速,与 video_infos 一一对应
{
"start": 0, // [必选] 场景开始时间(微秒)
"end": 6000000 // [必选] 场景结束时间(微秒)
"start": 0, // [必选] 场景开始时间 (微秒)
"end": 6000000 // [必选] 场景结束时间 (微秒)
}
]
// 变速原理:speed = (video.end - video.start) / (scene_timeline.end - scene_timeline.start)
// 示例:视频时间轴 0-2000000(2秒),场景时间线 0-1000000(1秒),则视频以2倍速播放
// 如果不提供scene_timelines或对应项为None,视频以正常速度(1.0倍)播放
alpha: 全局透明度[0, 1][可选],默认值为1.0
scale_x: X轴缩放比例[可选],默认值为1.0
scale_y: Y轴缩放比例[可选],默认值为1.0
transform_x: X轴位置偏移(像素)[可选],默认值为0
transform_y: Y轴位置偏移(像素)[可选],默认值为0
// 示例:视频时间轴 0-2000000(2 秒),场景时间线 0-1000000(1 秒),则视频以 2 倍速播放
// 如果不提供 scene_timelines 或对应项为 None,视频以正常速度 (1.0 倍) 播放
alpha: 全局透明度 [0, 1][可选],默认值为 1.0
scale_x: X 轴缩放比例 [可选],默认值为 1.0
scale_y: Y 轴缩放比例 [可选],默认值为 1.0
transform_x: X 轴位置偏移 (像素)[可选],默认值为 0
transform_y: Y 轴位置偏移 (像素)[可选],默认值为 0
Returns:
"draft_url": "https://capcut-mate.jcaigc.cn/openapi/capcut-mate/v1/get_draft?draft_id=...",
@@ -69,9 +70,131 @@ def add_videos(
Raises:
CustomException: 视频批量添加失败
"""
logger.info(f"add_videos, draft_url: {draft_url}, video_infos: {video_infos}, scene_timelines: {scene_timelines}, alpha: {alpha}, scale_x: {scale_x}, scale_y: {scale_y}, transform_x: {transform_x}, transform_y: {transform_y}")
# 调用内部处理函数(不获取锁,由外层控制)
return _add_videos_internal(
draft_url=draft_url,
video_infos=video_infos,
scene_timelines=scene_timelines,
alpha=alpha,
scale_x=scale_x,
scale_y=scale_y,
transform_x=transform_x,
transform_y=transform_y
)
# 1. 提取草稿ID
async def add_videos_async(
draft_url: str,
video_infos: str,
scene_timelines: Optional[List[Dict[str, int]]] = None,
alpha: float = 1.0,
scale_x: float = 1.0,
scale_y: float = 1.0,
transform_x: int = 0,
transform_y: int = 0,
lock_timeout: float = 30.0
) -> Tuple[str, str, List[str], List[str]]:
"""
添加视频到剪映草稿的异步版本(带并发锁保护)
功能:
1. 使用 DraftLockManager 防止同一草稿的并发写操作
2. 支持超时控制,避免无限等待
3. 自动释放锁,即使发生异常
Args:
draft_url: 草稿 URL,格式:".../get_draft?draft_id=xxx"
video_infos: JSON 字符串,包含视频信息列表,详见 add_videos 函数
scene_timelines: 场景时间线列表,用于视频变速,与 video_infos 一一对应
alpha: 全局透明度 [0, 1],默认 1.0
scale_x: X 轴缩放比例,默认 1.0
scale_y: Y 轴缩放比例,默认 1.0
transform_x: X 轴位置偏移(像素),默认 0
transform_y: Y 轴位置偏移(像素),默认 0
lock_timeout: 获取锁的超时时间(秒),默认 30 秒
Returns:
tuple: (draft_url, track_id, video_ids, segment_ids)
Raises:
CustomException: 视频添加失败或获取锁超时
asyncio.TimeoutError: 等待锁超时时抛出
Example:
>>> result = await add_videos_async(
... draft_url="http://.../draft_id=123",
... video_infos='[{"video_url":"...", "start":0, "end":5000000}]'
... )
"""
# 提取草稿 ID
draft_id = helper.get_url_param(draft_url, "draft_id")
if not draft_id:
raise CustomException(CustomError.INVALID_DRAFT_URL)
# 获取锁管理器
lock_manager = get_draft_lock_manager()
# 尝试获取锁
try:
await lock_manager.acquire_lock(draft_id, timeout=lock_timeout)
logger.info(f"Lock acquired for draft_id: {draft_id}")
except asyncio.TimeoutError:
logger.error(f"Timeout waiting for lock on draft_id: {draft_id}")
raise CustomException(
CustomError.VIDEO_ADD_FAILED,
f"Failed to acquire lock for draft {draft_id} after {lock_timeout}s"
)
try:
# 执行实际的添加操作
result = _add_videos_internal(
draft_url=draft_url,
video_infos=video_infos,
scene_timelines=scene_timelines,
alpha=alpha,
scale_x=scale_x,
scale_y=scale_y,
transform_x=transform_x,
transform_y=transform_y
)
return result
finally:
# 确保释放锁
await lock_manager.release_lock(draft_id)
logger.info(f"Lock released for draft_id: {draft_id}")
def _add_videos_internal(
draft_url: str,
video_infos: str,
scene_timelines: Optional[List[Dict[str, int]]] = None,
alpha: float = 1.0,
scale_x: float = 1.0,
scale_y: float = 1.0,
transform_x: int = 0,
transform_y: int = 0
) -> Tuple[str, str, List[str], List[str]]:
"""
添加视频的内部处理函数(无锁,需外层控制并发)
此函数不包含锁机制,必须在已获取锁的情况下调用
Args:
draft_url: 草稿 URL
video_infos: 视频信息 JSON 字符串
scene_timelines: 场景时间线列表
alpha: 全局透明度
scale_x: X 轴缩放比例
scale_y: Y 轴缩放比例
transform_x: X 轴位置偏移
transform_y: Y 轴位置偏移
Returns:
tuple: (draft_url, track_id, video_ids, segment_ids)
"""
logger.info(f"_add_videos_internal, draft_url: {draft_url}")
# 1. 提取草稿 ID
draft_id = helper.get_url_param(draft_url, "draft_id")
if (not draft_id) or (draft_id not in DRAFT_CACHE):
raise CustomException(CustomError.INVALID_DRAFT_URL)
@@ -103,16 +226,16 @@ def add_videos(
# 设置 relative_index=10 确保视频轨道在主视频轨道之上,避免与主轨道冲突
script.add_track(track_type=draft.TrackType.video, track_name=track_name, relative_index=10)
# 6. 遍历视频信息,添加视频到草稿中的指定轨道,收集片段ID
# 6. 遍历视频信息,添加视频到草稿中的指定轨道,收集片段 ID
segment_ids = []
current_track_end = 0 # 跟踪当前轨道上的实际结束位置(用于处理变速后的连续性)
for i, video in enumerate(videos):
# 获取对应的场景时间线(如果有)
scene_timeline = scene_timelines[i] if scene_timelines and i < len(scene_timelines) else None
# 自动调整视频的start时间,确保与前一个视频连续(处理变速后的间隙问题)
# 自动调整视频的 start 时间,确保与前一个视频连续(处理变速后的间隙问题)
if i > 0 and current_track_end > 0:
# 使用原始时长计算新的end
# 使用原始时长计算新的 end
original_duration = video['original_end'] - video['original_start']
video['start'] = current_track_end
video['end'] = video['start'] + original_duration
@@ -131,7 +254,7 @@ def add_videos(
# 7. 保存草稿
script.save()
# 8. 获取当前视频轨道id
# 8. 获取当前视频轨道 id
track_id = ""
for key in script.tracks.keys():
if script.tracks[key].name == track_name:
@@ -139,11 +262,11 @@ def add_videos(
break
logger.info(f"draft_id: {draft_id}, track_id: {track_id}")
# 9. 获取当前所有视频资源ID(全局唯一ID)
# 9. 获取当前所有视频资源 ID(全局唯一 ID
video_ids = [video.material_id for video in script.materials.videos]
logger.info(f"draft_id: {draft_id}, video_ids: {video_ids}")
# TODO: 这里还是有点小问题,为什么得到的video_idssegment_ids的结果一样
# TODO: 这里还是有点小问题,为什么得到的 video_idssegment_ids 的结果一样
return draft_url, track_id, video_ids, segment_ids
def add_video_to_draft(
+254
View File
@@ -0,0 +1,254 @@
"""
草稿并发锁管理器
用于防止同一草稿的并发写操作导致文件损坏
"""
import asyncio
from typing import Dict, Optional
from src.utils.logger import logger
class DraftLockManager:
"""
草稿锁管理器 - 单例模式
功能:
1. 为每个草稿 ID 维护一个独立的锁
2. 支持异步获取和释放锁
3. 自动清理已释放的锁以节省内存
4. 提供锁状态查询功能
使用场景:
- add_videos: 防止并发写入同一草稿文件
- add_audios: 防止并发修改同一草稿配置
- save_draft: 防止并发保存导致数据丢失
"""
_instance = None
_init_lock = asyncio.Lock()
def __new__(cls):
"""确保单例模式"""
if cls._instance is None:
cls._instance = super().__new__(cls)
return cls._instance
def __init__(self):
"""初始化管理器"""
# 如果已经初始化过,则跳过
if hasattr(self, '_initialized') and self._initialized:
return
# 存储每个草稿的锁:{draft_id: asyncio.Lock}
self._locks: Dict[str, asyncio.Lock] = {}
# 存储每个锁的持有者数量(用于引用计数)
self._lock_counts: Dict[str, int] = {}
# 初始化锁(用于保护_locks 字典的修改)
self._manager_lock = asyncio.Lock()
# 标记初始化完成
self._initialized = True
logger.info("DraftLockManager initialized")
async def acquire_lock(self, draft_id: str, timeout: Optional[float] = None) -> bool:
"""
获取指定草稿的锁
Args:
draft_id: 草稿 ID
timeout: 超时时间(秒),None 表示无限等待
Returns:
bool: 是否成功获取锁
Raises:
asyncio.TimeoutError: 等待超时时抛出
Example:
>>> lock_manager = DraftLockManager()
>>> success = await lock_manager.acquire_lock("2025092811473036584258", timeout=5.0)
>>> if success:
... try:
... # 执行草稿写操作
... pass
... finally:
... await lock_manager.release_lock("2025092811473036584258")
"""
async with self._manager_lock:
# 如果草稿 ID 没有锁,则创建新锁
if draft_id not in self._locks:
self._locks[draft_id] = asyncio.Lock()
self._lock_counts[draft_id] = 0
lock = self._locks[draft_id]
# 尝试获取锁(带超时)
try:
if timeout is not None:
# 使用 wait_for 实现超时
await asyncio.wait_for(lock.acquire(), timeout=timeout)
else:
# 无限等待
await lock.acquire()
# 增加引用计数
async with self._manager_lock:
self._lock_counts[draft_id] = self._lock_counts.get(draft_id, 0) + 1
logger.debug(f"Lock acquired for draft_id: {draft_id}, count: {self._lock_counts[draft_id]}")
return True
except asyncio.TimeoutError:
logger.warning(f"Timeout waiting for lock on draft_id: {draft_id}")
raise
async def release_lock(self, draft_id: str) -> None:
"""
释放指定草稿的锁
Args:
draft_id: 草稿 ID
Raises:
RuntimeError: 当尝试释放未持有的锁时抛出
KeyError: 当草稿 ID 不存在时抛出
Example:
>>> lock_manager = DraftLockManager()
>>> await lock_manager.acquire_lock("draft-123")
>>> try:
... # 执行写操作
... pass
... finally:
... await lock_manager.release_lock("draft-123")
"""
async with self._manager_lock:
if draft_id not in self._locks:
raise KeyError(f"No lock found for draft_id: {draft_id}")
lock = self._locks[draft_id]
self._lock_counts[draft_id] = max(0, self._lock_counts.get(draft_id, 0) - 1)
# 释放锁(在 manager_lock 之外,避免死锁)
try:
lock.release()
logger.debug(f"Lock released for draft_id: {draft_id}")
except RuntimeError as e:
logger.error(f"Failed to release lock for draft_id {draft_id}: {str(e)}")
raise
def is_locked(self, draft_id: str) -> bool:
"""
检查指定草稿是否被锁定
Args:
draft_id: 草稿 ID
Returns:
bool: 如果草稿被锁定返回 True,否则返回 False
Example:
>>> lock_manager = DraftLockManager()
>>> await lock_manager.acquire_lock("draft-123")
>>> print(lock_manager.is_locked("draft-123")) # True
>>> await lock_manager.release_lock("draft-123")
>>> print(lock_manager.is_locked("draft-123")) # False
"""
if draft_id not in self._locks:
return False
return self._locks[draft_id].locked()
def get_lock_count(self, draft_id: str) -> int:
"""
获取指定草稿的锁持有计数
Args:
draft_id: 草稿 ID
Returns:
int: 锁持有次数(重入次数)
Example:
>>> lock_manager = DraftLockManager()
>>> await lock_manager.acquire_lock("draft-123")
>>> print(lock_manager.get_lock_count("draft-123")) # 1
"""
return self._lock_counts.get(draft_id, 0)
def get_all_locked_drafts(self) -> list:
"""
获取所有当前被锁定的草稿 ID 列表
Returns:
list: 被锁定的草稿 ID 列表
Example:
>>> lock_manager = DraftLockManager()
>>> await lock_manager.acquire_lock("draft-123")
>>> locked = lock_manager.get_all_locked_drafts()
>>> print(locked) # ["draft-123"]
"""
return [
draft_id for draft_id, lock in self._locks.items()
if lock.locked()
]
async def clear_all_locks(self) -> None:
"""
清除所有锁(仅在紧急情况下使用)
Warning: 此方法会强制释放所有锁,可能导致数据不一致
仅应在系统异常或死锁检测时使用
Example:
>>> lock_manager = DraftLockManager()
>>> # 检测到死锁时
>>> await lock_manager.clear_all_locks()
"""
async with self._manager_lock:
released_count = len(self._locks)
self._locks.clear()
self._lock_counts.clear()
if released_count > 0:
logger.warning(f"Cleared all locks, released {released_count} locks")
def get_stats(self) -> dict:
"""
获取锁管理器统计信息
Returns:
dict: 包含锁统计信息的字典
Example:
>>> lock_manager = DraftLockManager()
>>> stats = lock_manager.get_stats()
>>> print(stats) # {"total_locks": 5, "locked_drafts": 2}
"""
locked_count = sum(1 for lock in self._locks.values() if lock.locked())
return {
"total_locks": len(self._locks),
"locked_drafts": locked_count,
"total_holders": sum(self._lock_counts.values())
}
# 全局单例
_draf_lock_manager: Optional[DraftLockManager] = None
def get_draft_lock_manager() -> DraftLockManager:
"""
获取全局草稿锁管理器实例
Returns:
DraftLockManager: 单例锁管理器实例
Example:
>>> lock_manager = get_draft_lock_manager()
>>> await lock_manager.acquire_lock("draft-123")
"""
global _draf_lock_manager
if _draf_lock_manager is None:
_draf_lock_manager = DraftLockManager()
return _draf_lock_manager