解决下载阻塞事件循环的问题。

This commit is contained in:
Hommy
2026-04-09 17:38:21 +08:00
parent 48b83d8e68
commit 5322b955e7
3 changed files with 47 additions and 11 deletions
+16 -3
View File
@@ -10,6 +10,7 @@ from src.utils.download import download
import config
import json
import asyncio
import time
from typing import List, Dict, Any, Tuple, Optional
from src.utils.draft_lock_manager import DraftLockManager
@@ -116,13 +117,25 @@ async def add_audios_async(
if not draft_id:
raise CustomException(CustomError.INVALID_DRAFT_URL)
prepared_audios = _prepare_audios_local_files(draft_url=draft_url, audio_infos=audio_infos)
logger.info(f"[flow:add_audios] prep_start, draft_id: {draft_id}")
# 下载与预处理放到线程池,避免阻塞事件循环导致锁超时漂移
prep_started_at = time.monotonic()
prepared_audios = await asyncio.to_thread(
_prepare_audios_local_files,
draft_url,
audio_infos,
)
logger.info(
f"[flow:add_audios] prep_done, draft_id: {draft_id}, "
f"count: {len(prepared_audios)}, elapsed: {time.monotonic() - prep_started_at:.3f}s"
)
lock_manager = DraftLockManager()
logger.info(f"[flow:add_audios] lock_wait_start, draft_id: {draft_id}, timeout: {lock_timeout}s")
try:
await lock_manager.acquire_lock(draft_id, timeout=lock_timeout)
logger.info(f"Lock acquired for draft_id: {draft_id}")
logger.info(f"[flow:add_audios] lock_acquired, draft_id: {draft_id}")
except asyncio.TimeoutError:
logger.error(f"Timeout waiting for lock on draft_id: {draft_id}")
raise CustomException(
@@ -138,7 +151,7 @@ async def add_audios_async(
)
finally:
await lock_manager.release_lock(draft_id)
logger.info(f"Lock released for draft_id: {draft_id}")
logger.info(f"[flow:add_audios] lock_released, draft_id: {draft_id}")
def validate_and_get_draft_id(draft_url: str) -> str:
+16 -3
View File
@@ -10,6 +10,7 @@ from src.utils.download import download
import config
import json
import asyncio
import time
from typing import List, Dict, Any, Tuple, Optional
from src.utils.draft_lock_manager import DraftLockManager
@@ -225,13 +226,25 @@ async def add_images_async(
if not draft_id:
raise CustomException(CustomError.INVALID_DRAFT_URL)
prepared_images = _prepare_images_local_files(draft_url=draft_url, image_infos=image_infos)
logger.info(f"[flow:add_images] prep_start, draft_id: {draft_id}")
# 下载与预处理放到线程池,避免阻塞事件循环导致锁超时漂移
prep_started_at = time.monotonic()
prepared_images = await asyncio.to_thread(
_prepare_images_local_files,
draft_url,
image_infos,
)
logger.info(
f"[flow:add_images] prep_done, draft_id: {draft_id}, "
f"count: {len(prepared_images)}, elapsed: {time.monotonic() - prep_started_at:.3f}s"
)
lock_manager = DraftLockManager()
logger.info(f"[flow:add_images] lock_wait_start, draft_id: {draft_id}, timeout: {lock_timeout}s")
try:
await lock_manager.acquire_lock(draft_id, timeout=lock_timeout)
logger.info(f"Lock acquired for draft_id: {draft_id}")
logger.info(f"[flow:add_images] lock_acquired, draft_id: {draft_id}")
except asyncio.TimeoutError:
logger.error(f"Timeout waiting for lock on draft_id: {draft_id}")
raise CustomException(
@@ -252,7 +265,7 @@ async def add_images_async(
)
finally:
await lock_manager.release_lock(draft_id)
logger.info(f"Lock released for draft_id: {draft_id}")
logger.info(f"[flow:add_images] lock_released, draft_id: {draft_id}")
def add_image_to_draft(
+15 -5
View File
@@ -11,6 +11,7 @@ from src.utils import helper
from src.utils.download import download
import config
import json
import time
from typing import List, Dict, Any, Tuple, Optional
from src.utils.draft_lock_manager import get_draft_lock_manager
@@ -158,17 +159,26 @@ async def add_videos_async(
if not draft_id:
raise CustomException(CustomError.INVALID_DRAFT_URL)
logger.info(f"[flow:add_videos] prep_start, draft_id: {draft_id}")
# 解析、规范化与下载在锁外完成,缩短持锁时间
prepared_videos = _prepare_videos_local_files(
draft_url=draft_url,
video_infos=video_infos,
# 下载与预处理放到线程池,避免阻塞事件循环导致锁超时漂移
prep_started_at = time.monotonic()
prepared_videos = await asyncio.to_thread(
_prepare_videos_local_files,
draft_url,
video_infos,
)
logger.info(
f"[flow:add_videos] prep_done, draft_id: {draft_id}, "
f"count: {len(prepared_videos)}, elapsed: {time.monotonic() - prep_started_at:.3f}s"
)
lock_manager = get_draft_lock_manager()
logger.info(f"[flow:add_videos] lock_wait_start, draft_id: {draft_id}, timeout: {lock_timeout}s")
try:
await lock_manager.acquire_lock(draft_id, timeout=lock_timeout)
logger.info(f"Lock acquired for draft_id: {draft_id}")
logger.info(f"[flow:add_videos] lock_acquired, draft_id: {draft_id}")
except asyncio.TimeoutError:
logger.error(f"Timeout waiting for lock on draft_id: {draft_id}")
raise CustomException(
@@ -190,7 +200,7 @@ async def add_videos_async(
)
finally:
await lock_manager.release_lock(draft_id)
logger.info(f"Lock released for draft_id: {draft_id}")
logger.info(f"[flow:add_videos] lock_released, draft_id: {draft_id}")
def _add_videos_internal(