From b6c13a904711df9f0d8e6934bd8df857e2505bc5 Mon Sep 17 00:00:00 2001 From: ericwyuan Date: Thu, 20 Aug 2026 12:27:30 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E5=88=86=E5=9D=97=E6=96=AD=E7=82=B9?= =?UTF-8?q?=E7=BB=AD=E4=BC=A0=E4=B8=8A=E4=BC=A0=20=E2=80=94=2020MB/?= =?UTF-8?q?=E5=9D=97=20+=20=E5=88=86=E5=9D=97=E7=BA=A7=E9=87=8D=E8=AF=95?= =?UTF-8?q?=20+=20=E6=96=AD=E7=82=B9=E6=9F=A5=E8=AF=A2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 视频跨公网单次上传超时/断连问题 --- fam-core/config/config.yaml | 16 +- .../src/fam_core/dispatcher/dispatcher.py | 161 ++++++++++++-- .../src/fam_edge/api_gateway/api_gateway.py | 199 ++++++++++++++++++ 3 files changed, 351 insertions(+), 25 deletions(-) diff --git a/fam-core/config/config.yaml b/fam-core/config/config.yaml index 141ca53..7189a28 100644 --- a/fam-core/config/config.yaml +++ b/fam-core/config/config.yaml @@ -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 diff --git a/fam-core/src/fam_core/dispatcher/dispatcher.py b/fam-core/src/fam_core/dispatcher/dispatcher.py index dae8a2d..9270e25 100644 --- a/fam-core/src/fam_core/dispatcher/dispatcher.py +++ b/fam-core/src/fam_core/dispatcher/dispatcher.py @@ -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() diff --git a/fam-edge/src/fam_edge/api_gateway/api_gateway.py b/fam-edge/src/fam_edge/api_gateway/api_gateway.py index 779e959..ca20b42 100644 --- a/fam-edge/src/fam_edge/api_gateway/api_gateway.py +++ b/fam-edge/src/fam_edge/api_gateway/api_gateway.py @@ -129,6 +129,205 @@ def enqueue_task(): return jsonify({"error": str(e)}), 500 +# ========== 分块上传(断点续传)========== + +CHUNK_SIZE = 20 * 1024 * 1024 # 20MB per chunk + + +def _chunk_dir(task_id: int) -> str: + upload_dir = os.environ.get('FAM_UPLOAD_DIR', '/tmp/fam_uploads') + d = os.path.join(upload_dir, f"task_{task_id}") + os.makedirs(d, exist_ok=True) + return d + + +@api_bp.route('/api/edge/video/chunk', methods=['POST']) +def upload_chunk(): + """接收单个分块,保存到 task_{id}/chunk_{index:04d} + + 断点续传:同一 task_id + chunk_index 重复上传会覆盖, + NAS 端可通过 /chunks 查询已上传分块,跳过已有的。 + """ + task_id_raw = request.form.get('task_id') + chunk_index_raw = request.form.get('chunk_index') + total_chunks_raw = request.form.get('total_chunks') + filename = request.form.get('filename', 'video.mp4') + chunk_file = request.files.get('chunk') + + if not task_id_raw or not chunk_index_raw or not chunk_file: + return jsonify({"error": "缺少必填字段: task_id, chunk_index, chunk"}), 400 + + try: + task_id = int(task_id_raw) + chunk_index = int(chunk_index_raw) + total_chunks = int(total_chunks_raw) if total_chunks_raw else 0 + except ValueError: + return jsonify({"error": "task_id/chunk_index 必须是整数"}), 400 + + d = _chunk_dir(task_id) + chunk_path = os.path.join(d, f"chunk_{chunk_index:04d}") + + try: + chunk_file.save(chunk_path) + size_kb = os.path.getsize(chunk_path) / 1024 + + # 写元数据(首次上传时) + meta_path = os.path.join(d, "meta.json") + if not os.path.exists(meta_path): + import json + meta = {"filename": filename, "total_chunks": total_chunks} + with open(meta_path, 'w') as f: + json.dump(meta, f) + + # 统计已上传分块 + uploaded = sorted([ + int(fn.split('_')[1]) for fn in os.listdir(d) + if fn.startswith('chunk_') and len(fn.split('_')) == 2 + ]) + + logger.info(f"[task_id={task_id}] 分块 {chunk_index}/{total_chunks} 上传成功 " + f"({size_kb:.0f}KB, 已上传 {len(uploaded)}/{total_chunks})") + + return jsonify({ + "status": "ok", + "task_id": task_id, + "chunk_index": chunk_index, + "uploaded_count": len(uploaded), + "total_chunks": total_chunks, + }), 200 + + except Exception as e: + logger.error(f"[task_id={task_id}] 分块上传失败: {e}", exc_info=True) + return jsonify({"error": str(e)}), 500 + + +@api_bp.route('/api/edge/video/chunks', methods=['GET']) +def query_chunks(): + """查询已上传分块列表(断点续传:NAS 重启后查询跳过已有分块)""" + task_id_raw = request.args.get('task_id') + if not task_id_raw: + return jsonify({"error": "缺少 task_id"}), 400 + + try: + task_id = int(task_id_raw) + except ValueError: + return jsonify({"error": "task_id 必须是整数"}), 400 + + d = _chunk_dir(task_id) + uploaded = sorted([ + int(fn.split('_')[1]) for fn in os.listdir(d) + if fn.startswith('chunk_') and len(fn.split('_')) == 2 + ]) if os.path.isdir(d) else [] + + total = 0 + meta_path = os.path.join(d, "meta.json") + if os.path.exists(meta_path): + import json + try: + with open(meta_path) as f: + total = json.load(f).get('total_chunks', 0) + except (json.JSONDecodeError, IOError): + pass + + return jsonify({ + "task_id": task_id, + "uploaded_chunks": uploaded, + "uploaded_count": len(uploaded), + "total_chunks": total, + }), 200 + + +@api_bp.route('/api/edge/video/assemble', methods=['POST']) +def assemble_chunks(): + """合并所有分块为完整视频文件,入 SQLite 队列 + + NAS 上传完全部分块后调用此端点触发合并 + 入队。 + """ + task_id_raw = request.form.get('task_id') + if not task_id_raw: + return jsonify({"error": "缺少 task_id"}), 400 + + try: + task_id = int(task_id_raw) + except ValueError: + return jsonify({"error": "task_id 必须是整数"}), 400 + + camera_name = request.form.get('camera_name', '') + event_start_time = request.form.get('event_start_time', '') + known_members_context = request.form.get('known_members_context', '') + + d = _chunk_dir(task_id) + + # 读取元数据 + import json + meta_path = os.path.join(d, "meta.json") + if not os.path.exists(meta_path): + return jsonify({"error": "元数据不存在,请先上传分块"}), 400 + + try: + with open(meta_path) as f: + meta = json.load(f) + except json.JSONDecodeError: + return jsonify({"error": "元数据损坏"}), 500 + + filename = meta.get('filename', 'video.mp4') + total_chunks = meta.get('total_chunks', 0) + + # 检查分块完整性 + chunk_files = sorted([ + fn for fn in os.listdir(d) + if fn.startswith('chunk_') and len(fn.split('_')) == 2 + ]) + + if total_chunks and len(chunk_files) < total_chunks: + missing = total_chunks - len(chunk_files) + return jsonify({ + "error": f"分块不完整: {len(chunk_files)}/{total_chunks},缺 {missing} 块", + "uploaded_count": len(chunk_files), + "total_chunks": total_chunks, + }), 400 + + # 合并分块 + upload_dir = os.environ.get('FAM_UPLOAD_DIR', '/tmp/fam_uploads') + video_filename = f"task_{task_id}_{filename}" + video_path = os.path.join(upload_dir, video_filename) + + try: + with open(video_path, 'wb') as out: + for cf in chunk_files: + chunk_path = os.path.join(d, cf) + with open(chunk_path, 'rb') as chunk_f: + out.write(chunk_f.read()) + + size_mb = os.path.getsize(video_path) / 1024 / 1024 + logger.info(f"[task_id={task_id}] 分块合并完成: {filename} ({size_mb:.1f}MB, {len(chunk_files)} 块)") + + # 清理分块目录 + import shutil + shutil.rmtree(d, ignore_errors=True) + + # 入队 + queue_id = queue_manager.enqueue( + nas_task_id=task_id, + video_filename=filename, + video_path=video_path, + camera_name=camera_name, + event_start_time=event_start_time, + known_members_context=known_members_context, + ) + + return jsonify({ + "status": "queued", + "task_id": task_id, + "queue_id": queue_id, + "size_mb": round(size_mb, 1), + }), 202 + + except Exception as e: + logger.error(f"[task_id={task_id}] 合并失败: {e}", exc_info=True) + return jsonify({"error": str(e)}), 500 + + @api_bp.route('/api/edge/results', methods=['GET']) def get_results(): """返回已完成但未拉取的结果,标记为已交付"""