优化添加视频的性能问题,减小加锁的范围。

This commit is contained in:
Hommy
2026-04-02 11:05:36 +08:00
parent 4f0b4ca1c9
commit 42f336bf03
2 changed files with 193 additions and 129 deletions
+64 -30
View File
@@ -79,10 +79,37 @@ def add_videos(
scale_x=scale_x,
scale_y=scale_y,
transform_x=transform_x,
transform_y=transform_y
transform_y=transform_y,
prepared_videos=None,
)
def _prepare_videos_local_files(draft_url: str, video_infos: str) -> List[Dict[str, Any]]:
"""
校验草稿、解析 video_infos、规范化时间字段并下载素材到草稿目录。
不含对 ScriptFile 的修改,可在草稿写锁外调用。
"""
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)
draft_dir = os.path.join(config.DRAFT_DIR, draft_id)
draft_video_dir = os.path.join(draft_dir, "assets", "videos")
os.makedirs(name=draft_video_dir, exist_ok=True)
videos = parse_video_data(json_str=video_infos)
if len(videos) == 0:
logger.info(f"No video info, draft_id: {draft_id}")
raise CustomException(CustomError.INVALID_VIDEO_INFO)
for video in videos:
video["original_start"] = video["start"]
video["original_end"] = video["end"]
video["local_video_path"] = download(url=video["video_url"], save_dir=draft_video_dir)
return videos
async def add_videos_async(
draft_url: str,
video_infos: str,
@@ -101,6 +128,7 @@ async def add_videos_async(
1. 使用 DraftLockManager 防止同一草稿的并发写操作
2. 支持超时控制,避免无限等待
3. 自动释放锁,即使发生异常
4. 视频下载在获取锁之前完成,持锁阶段仅修改草稿与写盘
Args:
draft_url: 草稿 URL,格式:".../get_draft?draft_id=xxx"
@@ -117,8 +145,7 @@ async def add_videos_async(
tuple: (draft_url, track_id, video_ids, segment_ids)
Raises:
CustomException: 视频添加失败获取锁超时
asyncio.TimeoutError: 等待锁超时时抛出
CustomException: 视频添加失败,或 `DRAFT_LOCK_TIMEOUT`获取锁超时
Example:
>>> result = await add_videos_async(
@@ -130,24 +157,27 @@ async def add_videos_async(
draft_id = helper.get_url_param(draft_url, "draft_id")
if not draft_id:
raise CustomException(CustomError.INVALID_DRAFT_URL)
# 获取锁管理器
# 解析、规范化与下载在锁外完成,缩短持锁时间
prepared_videos = _prepare_videos_local_files(
draft_url=draft_url,
video_infos=video_infos,
)
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"
CustomError.DRAFT_LOCK_TIMEOUT,
f"Failed to acquire lock for draft {draft_id} after {lock_timeout}s",
)
try:
# 执行实际的添加操作
result = _add_videos_internal(
return _add_videos_internal(
draft_url=draft_url,
video_infos=video_infos,
scene_timelines=scene_timelines,
@@ -155,11 +185,10 @@ async def add_videos_async(
scale_x=scale_x,
scale_y=scale_y,
transform_x=transform_x,
transform_y=transform_y
transform_y=transform_y,
prepared_videos=prepared_videos,
)
return result
finally:
# 确保释放锁
await lock_manager.release_lock(draft_id)
logger.info(f"Lock released for draft_id: {draft_id}")
@@ -172,7 +201,8 @@ def _add_videos_internal(
scale_x: float = 1.0,
scale_y: float = 1.0,
transform_x: int = 0,
transform_y: int = 0
transform_y: int = 0,
prepared_videos: Optional[List[Dict[str, Any]]] = None,
) -> Tuple[str, str, List[str], List[str]]:
"""
添加视频的内部处理函数(无锁,需外层控制并发)
@@ -181,13 +211,14 @@ def _add_videos_internal(
Args:
draft_url: 草稿 URL
video_infos: 视频信息 JSON 字符串
video_infos: 视频信息 JSON 字符串(当 prepared_videos 为 None 时参与解析)
scene_timelines: 场景时间线列表
alpha: 全局透明度
scale_x: X 轴缩放比例
scale_y: Y 轴缩放比例
transform_x: X 轴位置偏移
transform_y: Y 轴位置偏移
prepared_videos: 若已在外部完成解析与下载(含 local_video_path),则直接使用,跳过下载
Returns:
tuple: (draft_url, track_id, video_ids, segment_ids)
@@ -204,18 +235,17 @@ def _add_videos_internal(
draft_video_dir = os.path.join(draft_dir, "assets", "videos")
os.makedirs(name=draft_video_dir, exist_ok=True)
# 3. 解析视频信息
videos = parse_video_data(json_str=video_infos)
if len(videos) == 0:
logger.info(f"No video info, draft_id: {draft_id}")
raise CustomException(CustomError.INVALID_VIDEO_INFO)
if prepared_videos is not None:
videos = prepared_videos
else:
videos = parse_video_data(json_str=video_infos)
if len(videos) == 0:
logger.info(f"No video info, draft_id: {draft_id}")
raise CustomException(CustomError.INVALID_VIDEO_INFO)
for video in videos:
video["original_start"] = video["start"]
video["original_end"] = video["end"]
# 3.5 保存每个视频的原始时间信息(用于变速后的连续性计算)
for video in videos:
video['original_start'] = video['start']
video['original_end'] = video['end']
# 3.6 处理场景时间线(可选,已是对象数组)
logger.info(f"Parsed {len(videos)} videos, scene_timelines: {scene_timelines}")
# 4. 从缓存中获取草稿
@@ -314,8 +344,12 @@ def add_video_to_draft(
actual_duration: 视频在轨道上的实际播放时长(微秒),考虑变速后的时长
"""
try:
# 0. 下载视频
video_path = download(url=video['video_url'], save_dir=draft_video_dir)
video_path = video.get("local_video_path")
if video_path:
if not os.path.isfile(video_path):
raise CustomException(CustomError.VIDEO_ADD_FAILED, f"Missing local file: {video_path}")
else:
video_path = download(url=video["video_url"], save_dir=draft_video_dir)
# 1. 创建视频素材
video_material = draft.VideoMaterial(video_path)
+129 -99
View File
@@ -53,8 +53,9 @@ class TestAddVideosAsync:
patch('src.service.add_videos.add_video_to_draft') as mock_add, \
patch('src.service.add_videos.download') as mock_download:
# 设置 mock
# 设置 mockprepare 阶段会执行 draft_id in DRAFT_CACHE
mock_get_param.return_value = "test-draft-001"
mock_cache.__contains__.return_value = True
mock_parse.return_value = [{
'video_url': 'https://example.com/video1.mp4',
'width': 1920,
@@ -109,16 +110,28 @@ class TestAddVideosAsync:
# 先获取锁并不释放
await lock_manager.acquire_lock(draft_id)
# 尝试获取同一个草稿的锁(应该超时)
with pytest.raises(CustomException) as exc_info:
await add_videos_async(
draft_url=f"http://localhost/v1/get_draft?draft_id={draft_id}",
video_infos='[]',
lock_timeout=0.1
)
assert "Failed to acquire lock" in str(exc_info.value.args)
prepared = [{
"video_url": "https://example.com/v.mp4",
"start": 0,
"end": 1,
"duration": 1,
"original_start": 0,
"original_end": 1,
"local_video_path": "/tmp/timeout-test.mp4",
}]
# 尝试获取同一个草稿的锁(应该超时);prepare 在锁之前,需绕过网络/缓存
with patch("src.service.add_videos._prepare_videos_local_files", return_value=prepared):
with pytest.raises(CustomException) as exc_info:
await add_videos_async(
draft_url=f"http://localhost/v1/get_draft?draft_id={draft_id}",
video_infos="[]",
lock_timeout=0.1,
)
assert exc_info.value.err == CustomError.DRAFT_LOCK_TIMEOUT
assert "Failed to acquire lock" in exc_info.value.detail
# 清理
await lock_manager.release_lock(draft_id)
@@ -133,53 +146,61 @@ class TestAddVideosAsync:
"start": 0,
"end": 5000000
}])
prep_counter = {"n": 0}
def fake_prepare(draft_url: str, video_infos: str):
prep_counter["n"] += 1
n = prep_counter["n"]
return [{
"video_url": "https://example.com/video.mp4",
"start": 0,
"end": 5000000,
"duration": 5000000,
"original_start": 0,
"original_end": 5000000,
"local_video_path": f"/tmp/video_prep_{n}.mp4",
}]
async def add_video_task(task_id):
with patch('src.service.add_videos.helper.get_url_param') as mock_get_param, \
patch('src.service.add_videos.DRAFT_CACHE') as mock_cache, \
patch('src.service.add_videos.os.makedirs'), \
patch('src.service.add_videos.parse_video_data') as mock_parse, \
patch('src.service.add_videos.add_video_to_draft') as mock_add, \
patch('src.service.add_videos.download') as mock_download:
mock_get_param.return_value = "concurrent-test"
mock_parse.return_value = [{
'video_url': 'https://example.com/video.mp4',
'start': 0,
'end': 5000000,
'duration': 5000000,
'original_start': 0,
'original_end': 5000000
}]
mock_add.return_value = (f"segment-{task_id}", 5000000)
mock_download.return_value = f"/tmp/video_{task_id}.mp4"
mock_script = MagicMock()
mock_script.width = 1920
mock_script.height = 1080
mock_script.tracks = {}
mock_script.materials.videos = []
mock_cache.__getitem__.return_value = mock_script
try:
execution_order.append(f"{task_id}_start")
await add_videos_async(
draft_url=draft_url,
video_infos=video_infos,
lock_timeout=5.0
)
execution_order.append(f"{task_id}_complete")
except Exception as e:
execution_order.append(f"{task_id}_error: {str(e)}")
# 启动 3 个并发任务
tasks = [
add_video_task(1),
add_video_task(2),
add_video_task(3)
]
await asyncio.gather(*tasks)
try:
execution_order.append(f"{task_id}_start")
await add_videos_async(
draft_url=draft_url,
video_infos=video_infos,
lock_timeout=5.0,
)
execution_order.append(f"{task_id}_complete")
except Exception as e:
execution_order.append(f"{task_id}_error: {str(e)}")
with patch("src.service.add_videos.helper.get_url_param", return_value="concurrent-test"), \
patch("src.service.add_videos.DRAFT_CACHE") as mock_cache, \
patch("src.service.add_videos.os.makedirs"), \
patch("src.service.add_videos._prepare_videos_local_files", side_effect=fake_prepare), \
patch("src.service.add_videos.add_video_to_draft") as mock_add:
mock_cache.__contains__.return_value = True
mock_script = MagicMock()
mock_script.width = 1920
mock_script.height = 1080
mock_script.tracks = {}
mock_script.materials.videos = []
mock_cache.__getitem__.return_value = mock_script
call_seq = {"i": 0}
def add_side_effect(*args, **kwargs):
call_seq["i"] += 1
return (f"segment-{call_seq['i']}", 5000000)
mock_add.side_effect = add_side_effect
await asyncio.gather(
add_video_task(1),
add_video_task(2),
add_video_task(3),
)
# 验证任务是串行执行的(每个任务必须等待前一个释放锁)
# 第一个任务必须先完成
@@ -194,48 +215,55 @@ class TestAddVideosAsync:
async def test_concurrent_add_videos_different_drafts(self):
"""测试并发添加视频到不同草稿(可以并行执行)"""
completed_drafts = []
def fake_get_url_param(url: str, key: str):
if key != "draft_id":
return None
return url.split("draft_id=")[-1]
def fake_prepare(draft_url: str, video_infos: str):
did = fake_get_url_param(draft_url, "draft_id")
return [{
"video_url": f"https://example.com/video_{did}.mp4",
"start": 0,
"end": 5000000,
"duration": 5000000,
"original_start": 0,
"original_end": 5000000,
"local_video_path": f"/tmp/video_{did}.mp4",
}]
def make_script():
mock_script = MagicMock()
mock_script.width = 1920
mock_script.height = 1080
mock_script.tracks = {}
mock_script.materials.videos = []
return mock_script
async def add_video_task(draft_id):
with patch('src.service.add_videos.helper.get_url_param') as mock_get_param, \
patch('src.service.add_videos.DRAFT_CACHE') as mock_cache, \
patch('src.service.add_videos.os.makedirs'), \
patch('src.service.add_videos.parse_video_data') as mock_parse, \
patch('src.service.add_videos.add_video_to_draft') as mock_add, \
patch('src.service.add_videos.download') as mock_download:
mock_get_param.return_value = draft_id
mock_parse.return_value = [{
'video_url': f'https://example.com/video_{draft_id}.mp4',
'start': 0,
'end': 5000000,
'duration': 5000000,
'original_start': 0,
'original_end': 5000000
}]
mock_add.return_value = (f"segment-{draft_id}", 5000000)
mock_download.return_value = f"/tmp/video_{draft_id}.mp4"
mock_script = MagicMock()
mock_script.width = 1920
mock_script.height = 1080
mock_script.tracks = {}
mock_script.materials.videos = []
mock_cache.__getitem__.return_value = mock_script
await add_videos_async(
draft_url=f"http://localhost/v1/get_draft?draft_id={draft_id}",
video_infos='[]'
)
completed_drafts.append(draft_id)
# 并发访问 3 个不同的草稿
tasks = [
add_video_task("draft-a"),
add_video_task("draft-b"),
add_video_task("draft-c")
]
await asyncio.gather(*tasks)
await add_videos_async(
draft_url=f"http://localhost/v1/get_draft?draft_id={draft_id}",
video_infos="[]",
)
completed_drafts.append(draft_id)
with patch("src.service.add_videos.helper.get_url_param", side_effect=fake_get_url_param), \
patch("src.service.add_videos._prepare_videos_local_files", side_effect=fake_prepare), \
patch("src.service.add_videos.DRAFT_CACHE") as mock_cache, \
patch("src.service.add_videos.os.makedirs"), \
patch("src.service.add_videos.add_video_to_draft") as mock_add:
mock_cache.__contains__.return_value = True
mock_cache.__getitem__.side_effect = lambda _did: make_script()
mock_add.return_value = ("segment-mock", 5000000)
tasks = [
add_video_task("draft-a"),
add_video_task("draft-b"),
add_video_task("draft-c"),
]
await asyncio.gather(*tasks)
# 所有任务都应该完成
assert len(completed_drafts) == 3
@@ -297,6 +325,7 @@ class TestLockManagerIntegration:
async def test_lock_stats_accuracy(self):
"""测试锁统计信息准确性"""
lock_manager = DraftLockManager()
await lock_manager.clear_all_locks()
draft_ids = ["stats-1", "stats-2", "stats-3"]
# 初始状态
@@ -318,7 +347,8 @@ class TestLockManagerIntegration:
await lock_manager.release_lock(draft_ids[0])
stats = lock_manager.get_stats()
assert stats["total_locks"] == 2 # 锁对象被删除
# release_lock 不删除 _locks 中的条目,total_locks 仍为草稿数
assert stats["total_locks"] == 3
assert stats["locked_drafts"] == 2
assert stats["total_holders"] == 2