video_monitor/app/analysis/process_worker.py

242 lines
8.5 KiB
Python
Raw Permalink Normal View History

2026-08-30 22:22:11 +08:00
"""Monitor 分析子进程入口
每路摄像头在独立进程中运行 CameraPipeline绕过 GIL与主进程通过 Queue 通信
- event_queue: 上报检测/区域/追踪事件
- cmd_queue: stop / reload_zones
- status_dict: 共享状态Manager.dict openStatus 读取 FPS
"""
import json
import logging
import multiprocessing as mp
import os
import queue
import time
logger = logging.getLogger("analysis.process_worker")
def _algorithm_spec_from_dict(d):
return d
2026-09-04 18:16:14 +08:00
def _build_detectors_in_process(algorithm_specs, infer_req_q=None, infer_resp_q=None,
response_channel=None):
2026-08-30 22:22:11 +08:00
"""在子进程中构造检测器;若提供 infer_req_q/infer_resp_q 则走主进程共享推理池"""
detectors = []
if infer_req_q is not None and infer_resp_q is not None:
from app.analysis.remote_detector import RemoteDetector
for spec in algorithm_specs:
detectors.append({
"algorithm_id": spec.get("id", 0),
"algorithm_name": spec.get("name", ""),
2026-09-04 18:16:14 +08:00
"engine": RemoteDetector(spec, infer_req_q, infer_resp_q,
response_channel=response_channel),
"target_labels": spec.get("target_labels") or [],
"device": spec.get("device") or "cpu",
2026-08-30 22:22:11 +08:00
})
else:
from app.analysis.worker_pool import DetectorWorkerPool
pool = DetectorWorkerPool()
class _AlgoObj(object):
pass
for spec in algorithm_specs:
o = _AlgoObj()
for k, v in spec.items():
setattr(o, k, v)
eng = pool.get_detector(o)
if eng:
detectors.append({
"algorithm_id": spec.get("id", 0),
"algorithm_name": spec.get("name", ""),
"engine": eng,
2026-09-04 18:16:14 +08:00
"target_labels": spec.get("target_labels") or [],
"device": spec.get("device") or "cpu",
2026-08-30 22:22:11 +08:00
})
return detectors
def pipeline_process_main(config, event_queue, cmd_queue, status_dict,
infer_req_q=None, infer_resp_q=None):
"""子进程主函数spawn 入口config 必须为纯 dict"""
from app.utils.Logger import LOG_FORMAT
logging.basicConfig(level=logging.INFO, format=LOG_FORMAT)
log = logging.getLogger("analysis.process_worker")
from app.analysis.pipeline import CameraPipeline
from app.analysis.motion import MotionDetector
stream_id = config["stream_id"]
stream_code = config.get("stream_code", str(stream_id))
rtsp_url = config["rtsp_url"]
target_fps = config.get("target_fps", 5)
analyze_fps = config.get("analyze_fps", target_fps)
zones = config.get("zones") or []
algorithm_specs = config.get("algorithms") or []
use_shared_inference = bool(config.get("use_shared_inference", True))
storage_alarm_dir = config.get("storage_alarm_dir") or ""
static_dir = config.get("static_dir") or ""
def on_event(ev):
try:
event_queue.put(("event", ev), timeout=2.0)
except Exception as e:
log.warning("事件入队失败: %s", e)
def on_track_snapshot(sid, frame_index, active, has_motion):
try:
event_queue.put(("touch", {
"stream_id": sid,
"active": active,
"has_motion": has_motion,
}), timeout=1.0)
except Exception:
pass
2026-09-04 18:16:14 +08:00
def on_preview(payload):
try:
event_queue.put(("preview", payload), timeout=0.2)
except Exception:
pass
2026-08-30 22:22:11 +08:00
detectors = _build_detectors_in_process(
algorithm_specs,
infer_req_q if use_shared_inference else None,
infer_resp_q if use_shared_inference else None,
2026-09-04 18:16:14 +08:00
config.get("response_channel"),
2026-08-30 22:22:11 +08:00
)
if not detectors and not algorithm_specs:
log.info("pipeline[%s] 无算法,仅运动检测", stream_code)
motion = MotionDetector()
pipeline = CameraPipeline(
stream_id=stream_id,
stream_code=stream_code,
rtsp_url=rtsp_url,
detectors=detectors,
motion=motion,
target_fps=target_fps,
analyze_fps=analyze_fps,
on_event=on_event,
on_track_snapshot=on_track_snapshot,
2026-09-04 18:16:14 +08:00
on_preview=on_preview,
alarm_enabled=not bool(config.get("preview_only", False)),
2026-08-30 22:22:11 +08:00
zone_polygons=zones,
storage_alarm_dir=storage_alarm_dir,
static_dir=static_dir,
)
pipeline._algorithm_name = ", ".join(d.get("name", "") for d in algorithm_specs) or "motion-only"
import threading
2026-09-04 18:16:14 +08:00
finished = threading.Event()
def wait_started():
while not finished.is_set():
if pipeline._started_event.wait(0.1):
return True
return False
2026-08-30 22:22:11 +08:00
def status_reporter():
# 等待 pipeline.run() 启动_running 在 run() 里才置 True
# 否则 while pipeline._running 条件不满足会立即退出,导致 status_dict 永远为空)
2026-09-04 18:16:14 +08:00
if not wait_started():
return
while not finished.is_set():
2026-08-30 22:22:11 +08:00
try:
st = pipeline.status()
status_dict[str(stream_id)] = st
except Exception:
2026-09-04 18:16:14 +08:00
log.exception("pipeline[%s] 状态上报失败", stream_code)
finished.wait(1.0)
2026-08-30 22:22:11 +08:00
reporter = threading.Thread(target=status_reporter, name="status-%s" % stream_id, daemon=True)
reporter.start()
def cmd_listener():
2026-09-04 18:16:14 +08:00
if not wait_started():
return
while not finished.is_set():
2026-08-30 22:22:11 +08:00
try:
cmd = cmd_queue.get(timeout=0.5)
except queue.Empty:
continue
if not cmd:
continue
try:
if cmd.get("cmd") == "stop":
pipeline.stop()
break
if cmd.get("cmd") == "reload_zones":
pipeline.set_zone_polygons(cmd.get("zones") or [])
if cmd.get("analyze_fps") is not None:
pipeline.set_analyze_fps(cmd.get("analyze_fps"))
2026-09-04 18:16:14 +08:00
if cmd.get("cmd") == "set_analyze_fps":
pipeline.set_analyze_fps(cmd.get("analyze_fps") or 1)
2026-08-30 22:22:11 +08:00
except Exception as e:
log.exception("pipeline[%s] 命令处理失败: %s", stream_code, e)
cmd_thread = threading.Thread(target=cmd_listener, name="cmd-%s" % stream_id, daemon=True)
cmd_thread.start()
log.info("pipeline 子进程启动 stream=%s url=%s", stream_code, rtsp_url)
try:
pipeline.run()
finally:
2026-09-04 18:16:14 +08:00
finished.set()
reporter.join(timeout=2)
cmd_thread.join(timeout=1)
2026-08-30 22:22:11 +08:00
try:
status_dict.pop(str(stream_id), None)
except Exception:
pass
log.info("pipeline 子进程退出 stream=%s", stream_code)
class PipelineProcessHandle(object):
"""主进程侧对子进程的封装"""
def __init__(self, stream_id, process, event_queue, cmd_queue, status_dict):
self.stream_id = stream_id
self.process = process
self.event_queue = event_queue
self.cmd_queue = cmd_queue
self.status_dict = status_dict
self.running = True
def stop(self, timeout=5):
self.running = False
try:
self.cmd_queue.put({"cmd": "stop"}, timeout=1.0)
except Exception:
pass
try:
self.process.join(timeout=timeout)
if self.process.is_alive():
2026-09-04 18:16:14 +08:00
logger.warning("pipeline[%s] 未在 %ss 内停止,强制结束子进程", self.stream_id, timeout)
2026-08-30 22:22:11 +08:00
self.process.terminate()
self.process.join(timeout=2)
except Exception:
pass
def reload_zones(self, zones, analyze_fps=None):
try:
payload = {"cmd": "reload_zones", "zones": zones}
if analyze_fps is not None:
payload["analyze_fps"] = analyze_fps
self.cmd_queue.put(payload, timeout=1.0)
except Exception as e:
logger.warning("reload_zones 发送失败: %s", e)
2026-09-04 18:16:14 +08:00
def set_analyze_fps(self, analyze_fps):
try:
self.cmd_queue.put({"cmd": "set_analyze_fps", "analyze_fps": float(analyze_fps)}, timeout=1.0)
except Exception as e:
logger.warning("set_analyze_fps 发送失败: %s", e)
2026-08-30 22:22:11 +08:00
def status(self):
try:
return self.status_dict.get(str(self.stream_id))
except Exception:
return None