feat: 分块断点续传上传 — 20MB/块 + 分块级重试 + 断点查询

Edge 端新增 3 个端点:
- POST /api/edge/video/chunk: 接收单块,保存到 task_{id}/chunk_{index:04d}
- GET /api/edge/video/chunks: 查询已上传分块(断点续传)
- POST /api/edge/video/assemble: 合并全部分块入队

NAS Dispatcher 重写:
- 大文件(>50MB)自动分块上传(20MB/块)
- 每块最多重试 3 次(分块级重试,非整文件级)
- 上传前查询已上传分块,跳过已有的(断点续传)
- 全部上传后调 /assemble 合并入队
- 小文件(<=50MB)走直接上传路径
- max_retries 3→5(文件级重试次数)
- scheduler 切回生产目录

解决: 360MB 视频跨公网单次上传超时/断连问题
This commit is contained in:
ericwyuan
2026-08-20 12:27:30 +08:00
parent fb4ccb64dc
commit b6c13a9047
3 changed files with 351 additions and 25 deletions

View File

@@ -1,7 +1,9 @@
# FAM-Core 配置文件 (NAS 端) - 实际部署配置
# 注: Tailscale 防火墙待修复,当前 edge_url 使用 Oracle 公网 IP
# 异步队列模式: NAS 上传视频到 Edge /api/edge/video/enqueue 入队 → Edge 消费者异步处理
# → NAS Poller 定期从 /api/edge/results 拉取结果写库
# 异步队列模式 + 分块断点续传:
# 小文件(<=50MB): 直接上传 /enqueue
# 大文件(>50MB): 分块(20MB/块)上传 /chunk → /assemble 合并入队
# → NAS Poller 定期从 /api/edge/results 拉取结果写库
# Ollama 未对外暴露chat_handler 通过 FAM-Edge 代理
server:
@@ -18,9 +20,8 @@ database:
scheduler:
scan_interval: 60
# E2E 测试期间指向独立测试目录(正式目录 /volume1/surveillance 有 285 个历史视频,
# 全量建任务会导致 100GB 跨公网上传,待与用户确认回补策略后再切回)
video_dir: "/volume1/web/sentinel-home-ai/e2e-test"
# 正式目录 /volume1/surveillance/Generic_ONVIF-001
video_dir: "/volume1/surveillance/Generic_ONVIF-001"
video_extensions: [".mp4", ".mkv", ".avi"]
file_stable_seconds: 60
camera_name: "客厅"
@@ -28,9 +29,8 @@ scheduler:
dispatcher:
poll_interval: 30
edge_url: "http://129.146.203.203:5000/api/edge/video/enqueue"
max_retries: 3
upload_timeout: 300 # 仅视频上传时间(不含 AI 处理
stale_timeout: 600 # PROCESSING 超时回收10分钟Edge 处理 + 队列等待)
max_retries: 5 # 文件级重试次数分块级重试另计每块3次
stale_timeout: 600 # PROCESSING 超时回收10分钟
poller:
poll_interval: 30

View File

@@ -1,17 +1,21 @@
"""
Dispatcher - 30s 轮询 PENDING 任务,上传视频至 Edge 异步队列
流程(异步队列模式):
流程(异步队列模式 + 分块断点续传:
1. 读取任务对应的本地视频文件
2. multipart 上传至 Edge /api/edge/video/enqueue附已知成员清单等元数据
3. Edge 保存视频 + 入 SQLite 队列,立即返回 202
2a. 小文件 (<=50MB): 直接 multipart 上传至 /enqueue
2b. 大文件 (>50MB): 分块上传 (20MB/块) 至 /chunk支持断点续传最后调 /assemble 合并入队
3. Edge 保存视频 + 入 SQLite 队列,返回 202
4. Dispatcher 标记任务为 PROCESSING已派发等待 Poller 拉取结果)
5. Poller 线程定期从 Edge /api/edge/results 拉取结果,写库后标记 SUCCESS
退避重试: min(60 * (retry_count + 1) * 2, 600) 秒
分块级重试: 每块最多重试 3 次
"""
import os
import io
import time
import math
import threading
import requests
from datetime import datetime, timedelta
@@ -22,20 +26,27 @@ from .. import db_layer
logger = setup_logger('fam-core.dispatcher')
CHUNK_SIZE = 20 * 1024 * 1024 # 20MB per chunk
CHUNK_THRESHOLD = 50 * 1024 * 1024 # files > 50MB use chunked upload
MAX_CHUNK_RETRIES = 3
class Dispatcher:
"""任务下发器30s 轮询(异步队列模式"""
"""任务下发器30s 轮询(异步队列 + 分块断点续传"""
def __init__(self):
cfg = load_config()
self.poll_interval = cfg.get('dispatcher', {}).get('poll_interval', 30)
self.edge_url = cfg.get('dispatcher', {}).get('edge_url', 'http://localhost:5000/api/edge/video/enqueue')
self.edge_url = cfg.get('dispatcher', {}).get('edge_url',
'http://localhost:5000/api/edge/video/enqueue')
self.max_retries = cfg.get('dispatcher', {}).get('max_retries', 3)
# enqueue 模式只需上传时间,不含 AI 处理时间
self.upload_timeout = cfg.get('dispatcher', {}).get('upload_timeout', 300)
self.camera_name = cfg.get('scheduler', {}).get('camera_name', '默认摄像头')
# PROCESSING 超时回收Edge 处理 + 队列等待可能很长
self.stale_timeout = cfg.get('dispatcher', {}).get('stale_timeout', 3600)
self.stale_timeout = cfg.get('dispatcher', {}).get('stale_timeout', 600)
# 推导 Edge base URL
self.edge_base = self.edge_url.rsplit('/api/edge/video/enqueue', 1)[0]
self.chunk_url = f"{self.edge_base}/api/edge/video/chunk"
self.chunks_query_url = f"{self.edge_base}/api/edge/video/chunks"
self.assemble_url = f"{self.edge_base}/api/edge/video/assemble"
self._running = False
self._thread = None
@@ -72,7 +83,7 @@ class Dispatcher:
return payload
def _dispatch_one(self, task):
"""上传视频至 Edge 异步队列,立即返回"""
"""上传视频至 Edge 异步队列(自动选择直接/分块模式)"""
task_id = task['task_id']
video_path = task['video_path']
@@ -80,24 +91,36 @@ class Dispatcher:
db_layer.update_task_status(
task_id, 'FAILED',
error_message=f"视频文件不存在: {video_path}",
failure_stage='upload'
failure_stage='callback'
)
logger.error(f"[task_id={task_id}] 视频文件不存在,标记 FAILED: {video_path}")
return
payload = self._build_payload(task)
size_mb = os.path.getsize(video_path) / (1024 * 1024)
file_size = os.path.getsize(video_path)
size_mb = file_size / (1024 * 1024)
db_layer.update_task_status(task_id, 'PROCESSING')
log_task(logger, task_id, 'dispatcher',
f'上传视频至 Edge 队列: {self.edge_url} ({size_mb:.1f}MB)')
payload = self._build_payload(task)
if file_size > CHUNK_THRESHOLD:
logger.info(f"[task_id={task_id}] 大文件分块上传: {size_mb:.1f}MB, "
f"{math.ceil(file_size / CHUNK_SIZE)}")
self._dispatch_chunked(task, payload, video_path, file_size)
else:
log_task(logger, task_id, 'dispatcher',
f'直接上传: {self.edge_url} ({size_mb:.1f}MB)')
self._dispatch_direct(task, payload, video_path)
def _dispatch_direct(self, task, payload, video_path):
"""小文件直接上传至 /enqueue"""
task_id = task['task_id']
try:
with open(video_path, 'rb') as fh:
resp = requests.post(
self.edge_url,
data=payload,
files={'video': (os.path.basename(video_path), fh, 'video/mp4')},
timeout=(10, 60)
timeout=(10, 120)
)
except requests.RequestException as e:
logger.error(f"[task_id={task_id}] 上传失败: {e}")
@@ -121,6 +144,110 @@ class Dispatcher:
logger.error(f"[task_id={task_id}] Edge 返回异常状态码: {resp.status_code}")
self._schedule_retry(task)
def _dispatch_chunked(self, task, payload, video_path, file_size):
"""大文件分块上传 + 断点续传
1. 查询 Edge 端已上传分块(断点续传)
2. 上传缺失分块(每块最多重试 3 次)
3. 全部分块上传后调用 /assemble 合并入队
"""
task_id = task['task_id']
total_chunks = math.ceil(file_size / CHUNK_SIZE)
filename = os.path.basename(video_path)
# 1. 查询已上传分块(断点续传)
uploaded_set = set()
try:
resp = requests.get(
self.chunks_query_url,
params={'task_id': task_id},
timeout=(10, 15)
)
if resp.status_code == 200:
data = resp.json()
uploaded_set = set(data.get('uploaded_chunks', []))
if uploaded_set:
logger.info(f"[task_id={task_id}] 断点续传: 已有 {len(uploaded_set)}/{total_chunks}")
except requests.RequestException as e:
logger.warning(f"[task_id={task_id}] 查询已上传分块失败(将从头上传): {e}")
# 2. 上传缺失分块
try:
with open(video_path, 'rb') as fh:
for idx in range(total_chunks):
if idx in uploaded_set:
continue
chunk_data = fh.read(CHUNK_SIZE)
if not chunk_data:
break
success = False
for attempt in range(MAX_CHUNK_RETRIES):
try:
cresp = requests.post(
self.chunk_url,
data={
'task_id': str(task_id),
'chunk_index': str(idx),
'total_chunks': str(total_chunks),
'filename': filename,
},
files={'chunk': (f'chunk_{idx}', io.BytesIO(chunk_data))},
timeout=(10, 120)
)
if cresp.status_code == 200:
success = True
break
logger.warning(f"[task_id={task_id}] 分块 {idx} 返回 {cresp.status_code}(尝试 {attempt+1}/{MAX_CHUNK_RETRIES}")
except requests.RequestException as e:
logger.warning(f"[task_id={task_id}] 分块 {idx} 上传失败(尝试 {attempt+1}/{MAX_CHUNK_RETRIES}: {e}")
if attempt < MAX_CHUNK_RETRIES - 1:
time.sleep(5 * (attempt + 1))
if not success:
logger.error(f"[task_id={task_id}] 分块 {idx} 超过最大重试次数,安排文件级重试")
self._schedule_retry(task)
return
if (idx + 1) % 5 == 0 or idx == total_chunks - 1:
logger.info(f"[task_id={task_id}] 分块进度: {idx + 1}/{total_chunks}")
except IOError as e:
logger.error(f"[task_id={task_id}] 读取视频文件失败: {e}")
self._schedule_retry(task)
return
# 3. 合并 + 入队
try:
aresp = requests.post(
self.assemble_url,
data={
'task_id': str(task_id),
'camera_name': payload.get('camera_name', ''),
'event_start_time': payload.get('event_start_time', ''),
'known_members_context': payload.get('known_members_context', ''),
},
timeout=(10, 60)
)
except requests.RequestException as e:
logger.error(f"[task_id={task_id}] 合并请求失败: {e}")
self._schedule_retry(task)
return
if aresp.status_code == 202:
try:
data = aresp.json()
queue_id = data.get('queue_id', '?')
asm_size = data.get('size_mb', '?')
logger.info(f"[task_id={task_id}] 分块合并入队成功 (queue_id={queue_id}, {asm_size}MB),等待 Poller 拉取结果")
except ValueError:
logger.info(f"[task_id={task_id}] 分块合并入队成功,等待 Poller 拉取结果")
return
logger.error(f"[task_id={task_id}] 合并端点返回 {aresp.status_code}: {aresp.text[:200]}")
self._schedule_retry(task)
def _schedule_retry(self, task):
"""调度重试"""
task_id = task['task_id']
@@ -158,7 +285,7 @@ class Dispatcher:
def _run(self):
"""线程主循环"""
logger.info(f"Dispatcher 启动 (enqueue 模式),轮询间隔 {self.poll_interval}s")
logger.info(f"Dispatcher 启动 (enqueue + 分块模式),轮询间隔 {self.poll_interval}s")
while self._running:
try:
self._poll_once()