首页
/ WiFi-DensePose v1 分布式架构全解析:从 WiFi CSI 到实时人体姿态估计的系统设计

WiFi-DensePose v1 分布式架构全解析:从 WiFi CSI 到实时人体姿态估计的系统设计

2026-09-07 17:45:35作者:丁柯新Fawn

导读:本文系统讲解 RuView 仓库中归档的 WiFi-DensePose v1(archive/v1)架构设计文档,该文档描述了一套把 WiFi 信道状态信息(CSI)转化为实时人体姿态估计结果的分布式微服务系统。读完本文,你将掌握这套系统的分层架构、核心组件职责、实时数据处理管线、REST/WebSocket API 设计、时间序列存储与缓存选型、安全与隐私机制,以及容器与 Kubernetes 部署方式,并能依据仓库源码逐一对齐每个架构组件的真实实现位置与关键配置参数。

文档定位与适用边界

本篇文章的主体源自仓库中的 architecture-overview.md,属于 archive/v1 目录下的开发者文档。需要先说明两点背景,避免读者产生误导:

  • archive/v1 是 WiFi-DensePose 早期的纯 Python 实现,已被仓库明确标记为“保留用于研究归档、不再维护、被 v2 Rust 工作区取代”,依据见 archive/v1/DEPRECATED.mdADR-187
  • 因此,下文中的架构描述以该文档的设计蓝图为骨架,同时会在每个组件处指出 archive/v1/src 中的真实实现路径;两者不一致的地方会明确标注,绝不混淆“文档设计”与“实际代码”。

系统架构总览

该文档将 WiFi-DensePose 定义为基于微服务的分布式系统,自上而下分为五层:

┌─────────────────────────────────────────────────────────────────┐
│                        WiFi-DensePose System                    │
├─────────────────────────────────────────────────────────────────┤
│  ┌─────────────────┐  ┌─────────────────┐  ┌─────────────────┐  │
│  │   Client Apps   │  │   Web Dashboard │  │  Mobile Apps    │  │
│  └─────────────────┘  └─────────────────┘  └─────────────────┘  │
├─────────────────────────────────────────────────────────────────┤
│                        API Gateway                              │
├─────────────────────────────────────────────────────────────────┤
│  ┌─────────────────┐  ┌─────────────────┐  ┌─────────────────┐  │
│  │   REST API      │  │  WebSocket API  │  │   MQTT Broker   │  │
│  └─────────────────┘  └─────────────────┘  └─────────────────┘  │
├─────────────────────────────────────────────────────────────────┤
│                      Processing Layer                           │
│  ┌─────────────────┐  ┌─────────────────┐  ┌─────────────────┐  │
│  │ Pose Estimation │  │    Tracking     │  │   Analytics     │  │
│  │    Service      │  │    Service      │  │    Service      │  │
│  └─────────────────┘  └─────────────────┘  └─────────────────┘  │
├─────────────────────────────────────────────────────────────────┤
│                       Data Layer                                │
│  ┌─────────────────┐  ┌─────────────────┐  ┌─────────────────┐  │
│  │ CSI Processor   │  │  Data Pipeline  │  │  Model Manager  │  │
│  └─────────────────┘  └─────────────────┘  └─────────────────┘  │
├─────────────────────────────────────────────────────────────────┤
│                     Hardware Layer                              │
│  ┌─────────────────┐  ┌─────────────────┐  ┌─────────────────┐  │
│  │  WiFi Routers   │  │ Processing Unit │  │   GPU Cluster   │  │
│  │   (CSI Data)    │  │   (CPU/Memory)  │  │  (Neural Net)   │  │
│  └─────────────────┘  └─────────────────┘  └─────────────────┘  │
└─────────────────────────────────────────────────────────────────┘

各层之间的协作可以用下面的组件交互图概括:路由器网络采集 CSI → CSI Processor 清洗 → Feature Extractor 特征化 → Neural Network 推理 → Pose Tracker 跟踪 → Analytics Engine 产生事件与告警,最终推送给 Client Applications,同时由 Model Manager 负责模型版本供给。

┌─────────────┐    CSI Data    ┌─────────────┐    Features    ┌─────────────┐
│   Router    │ ──────────────▶│ CSI         │ ──────────────▶│ Feature     │
│   Network   │                │ Processor   │                │ Extractor   │
└─────────────┘                └─────────────┘                └─────────────┘
                                       │                              │
                                       ▼                              ▼
┌─────────────┐    Poses       ┌─────────────┐    Inference   ┌─────────────┐
│   Client    │ ◀──────────────│ Pose        │ ◀──────────────│ Neural      │
│ Applications│                │ Tracker     │                │ Network     │
└─────────────┘                └─────────────┘                └─────────────┘
       │                               │                              │
       ▼                               ▼                              ▼
┌─────────────┐    Events      ┌─────────────┐    Models      ┌─────────────┐
│ Alert       │ ◀──────────────│ Analytics   │ ◀──────────────│ Model       │
│ System      │                │ Engine      │                │ Manager     │
└─────────────┘                └─────────────┘                └─────────────┘

从架构蓝图到真实代码的映射

在动手阅读后续组件之前,先给出蓝图与 archive/v1 实际源码的对应关系,便于读者带着映射去读代码:

架构文档中的组件 蓝图路径(原文档) archive/v1 中真实落点
CSI Data Processor src/hardware/csi_processor.py core/csi_processor.pyhardware/csi_extractor.py
Neural Network Service src/neural_network/inference.py services/pose_service.pymodels/densepose_head.pymodels/modality_translation.py
API Gateway src/api/main.py api/main.py
Tracking Service src/tracking/tracker.py v1 未单独成模块(见下文说明)
Analytics Engine src/analytics/engine.py v1 中体现为 pose_service 的 activity 分类与统计逻辑
WebSocket api/websocket/connection_manager.pyapi/websocket/pose_stream.py

注意:tracking/neural_network/analytics/ 等目录在 archive/v1/src 中并不存在。该文档部分代码属于设计示意archive/v1 的参考实现将 CSI 处理、模型推理与活动分类统一收敛在 PoseService 之中。阅读源码时以此为准。

核心组件逐个拆解

1. CSI Data Processor:信号入口与底层净化

文档中定义的 CSI Processor 承担四项职责:多路由器实时 CSI 采集、信号预处理与降噪、相位清洗与幅度归一化、多天线数据融合。其设计示意如下:

class CSIProcessor:
    """Processes raw CSI data from WiFi routers."""

    def __init__(self, config: CSIConfig):
        self.routers = self._initialize_routers(config.routers)
        self.buffer = CircularBuffer(config.buffer_size)
        self.preprocessor = CSIPreprocessor()

    async def process_stream(self) -> AsyncGenerator[CSIData, None]:
        """Process continuous CSI data stream."""
        async for raw_data in self._receive_csi_data():
            processed_data = self.preprocessor.process(raw_data)
            yield processed_data

仓库中的真实实现

archive/v1 中真正承载该职责的是 core/csi_processor.py,其构造函数通过 _validate_config 强制校验 4 个必填参数:sampling_rate(>0)、window_size(>0)、overlap(0≤x<1)、noise_threshold,并支持一系列可选参数:

参数 默认值 含义
human_detection_threshold 0.8 判定“有人”的平滑置信度阈值
smoothing_factor 0.9 置信度指数滑动平均(EMA)系数
max_history_size 500 CSI 历史环形缓冲长度
doppler_window 64 多普勒计算回溯帧数上限
enable_preprocessing / enable_feature_extraction / enable_human_detection True 三段管线开关

它的核心链路 process_csi_data 依次执行三个步骤:preprocess_csi_dataextract_featuresdetect_human_presence

async def process_csi_data(self, csi_data: CSIData) -> HumanDetectionResult:
    # Preprocess the data
    preprocessed_data = self.preprocess_csi_data(csi_data)
    # Extract features
    features = self.extract_features(preprocessed_data)
    # Detect human presence
    detection_result = self.detect_human_presence(features)
    # Add to history
    self.add_to_history(csi_data)
    return detection_result

预处理阶段包含三个可独立观察的子步骤(均可通过 metadata 标注验证执行轨迹):

  1. _remove_noise:把幅度转为 dB 后按 noise_threshold 生成噪声掩码,低于阈值的子载波幅度置零;
  2. _apply_windowing:用 scipy.signal.windows.hamming(num_subcarriers) 汉明窗抑制频谱泄漏;
  3. _normalize_amplitude:除以 std + 1e-12 归一化到单位方差。

特征提取阶段输出一个结构化的 CSIFeatures dataclass(源码定义),包含:幅度均值/方差、相邻子载波相位差、天线间相关系数矩阵、多普勒谱(对缓存的历史相位做 FFT,归一化到 64 个 bin)、当前帧功率谱密度(128 点 FFT)。源码特意用 deque 缓存每帧平均相位,保证多普勒提取为 O(1) 追加、窗口计算有界。

人体存在性检测综合三类指标加权:幅度指标(40%)、相位标准差指标(30%)、运动分数指标(30%),其中运动分数 motion_score = 0.6 * variance_score + 0.4 * correlation_score,再经 _apply_temporal_smoothing 做 EMA 平滑后与阈值比较。

从硬件角度看,hardware/csi_extractor.py 提供了 CSIData 数据类(timestamp/amplitude/phase/frequency/bandwidth/num_subcarriers/num_antennas/snr/metadata)与 ESP32CSIParser,用于解析 CSI_DATA: 前缀的 ESP32 固件上报格式;当数据字段不足或包含非数值时会抛出 CSIExtractionError,绝不静默返回占位数据。

2. Neural Network Service:姿态估计推理

文档把该组件定位为深度学习推理层,核心特性包括 DensePose 推理、批量处理、GPU 加速、模型版本化与热替换:

class PoseEstimationService:
    """Neural network service for pose estimation."""

    def __init__(self, model_config: ModelConfig):
        self.model = self._load_model(model_config.model_path)
        self.device = torch.device('cuda' if torch.cuda.is_available() else 'cpu')
        self.batch_processor = BatchProcessor(model_config.batch_size)

    async def estimate_poses(self, csi_features: CSIFeatures) -> List[PoseEstimation]:
        """Estimate human poses from CSI features."""
        with torch.no_grad():
            predictions = self.model(csi_features.to(self.device))
            return self._postprocess_predictions(predictions)

仓库中的真实实现

真实的推理编排在 services/pose_service.pyPoseService 中,与文档的“独立服务”描述不同,它是单体编排器:初始化时同时构造 CSI Processor、PhaseSanitizer,并依据 mock_pose_data 开关决定是否加载真实模型(见 _initialize_models)。

真实推理链路 _estimate_poses 可分为四步:

  1. 将 CSI 数据转为 torch.Tensor,必要时补 batch 维;
  2. 交给模态翻译网络(modality translator):从 64 通道 CSI 特征翻译到 256 通道“类视觉”特征,隐藏层为 [128, 256, 512],开启 attention;
  3. 将翻译后的视觉特征送入 DensePoseHead 输出原始预测张量;
  4. _parse_pose_outputs 中解析出 17 个关键点(按 nose、左右眼/耳/肩/肘/腕/髋/膝/踝的顺序)、边界框与活动类别,并按 pose_confidence_threshold 过滤、按 pose_max_persons 截断。

两个真实模型都位于 models/

  • densepose_head.pyDensePoseHead 需要 input_channels(与翻译器输出对齐,即 256)、num_body_parts(24,标准 DensePose body parts 数)、num_uv_coordinates(2);
  • modality_translation.pyModalityTranslationNetwork 完成 CSI→视觉特征跨模态翻译。

需要强调(DEPRECATED.md 中明确说明):archive/v1/src/models/densepose_head.py 只定义了网络结构,没有附带任何训练权重,权重初始化仅使用 kaiming_normal_ 随机初始化,目录下没有任何 .pth/.onnx/.safetensors 文件。也就是说,该树是一个“架构已被定义但不可直接产出真实姿态精度”的研究归档;真实训练权重位于 v2 工作区与 Hugging Face 仓库。

3. Tracking Service:跨帧身份保持

文档对跟踪服务的规划是维护时序一致性与人物身份,包含多目标卡尔曼滤波跟踪、行人重识别(ReID)、跟踪生命周期管理与轨迹平滑:

class PersonTracker:
    """Tracks multiple persons across time."""

    def __init__(self, tracking_config: TrackingConfig):
        self.tracks = {}
        self.track_id_counter = 0
        self.kalman_filter = KalmanFilter()
        self.reid_model = ReIDModel()

    def update(self, detections: List[PoseDetection]) -> List[TrackedPose]:
        """Update tracks with new detections."""
        matched_tracks, unmatched_detections = self._associate_detections(detections)
        self._update_matched_tracks(matched_tracks)
        self._create_new_tracks(unmatched_detections)
        return self._get_active_tracks()

archive/v1/src 目录结构中并未找到独立的 tracking/ 模块,因此这段应被理解为设计蓝图中的职责边界;若需更完整的跟踪/身份语义实现,可参考仓库中与之相关的演进文档,例如 ADR-082-pose-tracker-confirmed-output-filterADR-026-survivor-track-lifecycle 以及 v2 工作区。

4. API Gateway:统一接入面

文档规划 API Gateway 统一提供 REST 与 WebSocket 接入,负责认证鉴权、限流、路由与负载均衡、API 版本管理。仓库真实实现位于 api/main.py,是一个 FastAPI 应用,其启动逻辑(lifespan)清晰体现了“服务编排 + 后台任务 + 优雅关闭”三件事:

@asynccontextmanager
async def lifespan(app: FastAPI):
    # 初始化 hardware / pose / stream 三个服务,以及 PoseStreamHandler
    await initialize_services(app)
    await start_background_tasks(app)   # 启动 pose service 与实时流
    yield
    await cleanup_services(app)          # 逆序关闭流、连接管理器与各服务

中间件按功能开关挂载(对应 settings.py 中的特性开关):

if settings.enable_rate_limiting:
    app.add_middleware(RateLimitMiddleware)
if settings.enable_authentication:
    app.add_middleware(AuthMiddleware)
app.add_middleware(CORSMiddleware, **cors_config)
if settings.is_production:
    app.add_middleware(TrustedHostMiddleware, allowed_hosts=settings.allowed_hosts)

路由挂载同样遵循 /api/v1 前缀(settings.api_prefix),分为健康检查、姿态估计、流媒体、认证四大组,与文档 REST API 树中的 health/system/pose/config/analytics 布局基本吻合(详见下文 API 小节)。

5. Analytics Engine:洞察与告警

文档把 Analytics Engine 定义为处理姿态数据、生成洞察并触发告警的组件,支持跌倒/入侵等实时事件检测、统计分析、面向医疗/零售/安全等行业的领域分析,以及基于机器学习的模式识别:

class AnalyticsEngine:
    """Processes pose data for insights and alerts."""

    def __init__(self, domain_config: DomainConfig):
        self.domain = domain_config.domain
        self.event_detectors = self._load_event_detectors(domain_config)
        self.alert_manager = AlertManager(domain_config.alerts)

    async def process_poses(self, poses: List[TrackedPose]) -> AnalyticsResult:
        events = []
        for detector in self.event_detectors:
            detected_events = await detector.detect(poses)
            events.extend(detected_events)
        await self.alert_manager.process_events(events)
        return AnalyticsResult(events=events, metrics=self._calculate_metrics(poses))

在 v1 参考实现中,这部分能力被折叠进 PoseService_classify_activity 依据输出特征范数做阈值式活动判定(>2.0 为 walking,>1.0 为 standing,>0.5 为 sitting,>0.1 为 lying),estimate_poses 返回带 zone_summary(按区域统计人数)的结构化结果;而“区域/领域”维度的配置(ZoneType、ActivityType、区域物理边界、主备路由器分配等)位于 config/domains.py

实时数据流与数据模型

五级实时处理管线

文档给出了从 CSI 采集到输出的完整实时管线,是理解全系统时序关系的关键:

1. CSI Data Acquisition
   ┌─────────────┐
   │   Router 1  │ ──┐
   └─────────────┘   │
   ...  Router N ────┼───▶ CSI Buffer ──▶ Preprocessor
                     ▼
2. Feature Extraction
   Phase Sanitizer ◀── Feature Extractor
        │                     │
        ▼                     ▼
   Amplitude Processor   Frequency Analyzer
        └──────────┬───────┘
                   ▼
3. Neural Network Inference
   DensePose Model ──▶ Pose Decoder
                   ▼
4. Tracking and Analytics
   Raw Pose Detections ──▶ Person Tracker ──▶ Analytics Engine
                   ▼
5. Output and Storage
   Tracked Poses ──▶ WebSocket Streams ──▶ Client Applications
        │
        ▼
   Database Storage

CSI 数据结构

文档定义了两级 dataclass:CSIData 描述一次 CSI 测量(含天线对、子载波与元数据),SubcarrierData 描述单个子载波的频点、复数幅度、相位与 SNR:

@dataclass
class CSIData:
    """Channel State Information data structure."""
    timestamp: datetime
    router_id: str
    antenna_pairs: List[AntennaPair]
    subcarriers: List[SubcarrierData]
    metadata: CSIMetadata

@dataclass
class SubcarrierData:
    """Individual subcarrier information."""
    frequency: float
    amplitude: complex
    phase: float
    snr: float

v1 参考实现中 hardware.csi_extractor.CSIData 采用 (num_antennas, num_subcarriers) 的 numpy 数组承载幅度/相位,与文档结构并不完全一致——这是“文档示意”与“代码落地”差异的又一例证。

姿态数据结构

推理结果与跟踪结果的建模如下(TrackedPose 额外携带 track_id、速度矢量、轨迹年龄与轨迹置信度):

@dataclass
class PoseEstimation:
    """Human pose estimation result."""
    person_id: Optional[int]
    confidence: float
    bounding_box: BoundingBox
    keypoints: List[Keypoint]
    dense_pose: Optional[DensePoseResult]
    timestamp: datetime

@dataclass
class TrackedPose:
    """Tracked pose with temporal information."""
    track_id: int
    pose: PoseEstimation
    velocity: Vector2D
    track_age: int
    track_confidence: float

处理管线:从原始 CSI 到可推理特征

CSI 预处理

预处理在参考实现中的对应方法即上文所述 _remove_noise → _apply_windowing → _normalize_amplitude。文档给出的示意类结构如下,可作为阅读真实代码时的“分形结构”参照:

class CSIPreprocessor:
    def __init__(self, config: PreprocessingConfig):
        self.phase_sanitizer = PhaseSanitizer()
        self.amplitude_normalizer = AmplitudeNormalizer()
        self.noise_filter = NoiseFilter(config.filter_params)

    def process(self, raw_csi: RawCSIData) -> ProcessedCSIData:
        # Phase unwrapping and sanitization
        sanitized_phase = self.phase_sanitizer.sanitize(raw_csi.phase)
        # Amplitude normalization
        normalized_amplitude = self.amplitude_normalizer.normalize(raw_csi.amplitude)
        # Noise filtering
        filtered_data = self.noise_filter.filter(sanitized_phase, normalized_amplitude)
        return ProcessedCSIData(...)

v1 中相位处理独立为 core/phase_sanitizer.py,其配置项(unwrapping_method=numpyoutlier_threshold=3.0smoothing_window=5 及三个 enable 开关)在 PoseService._initialize 阶段从 settings 派生(见 pose_service.py)。

特征提取

特征提取器按配置选择要提取的特征族(amplitude / phase / doppler),再经 PCA 降维输出 CSIFeatures

class FeatureExtractor:
    def __init__(self, config: FeatureConfig):
        self.window_size = config.window_size
        self.feature_types = config.feature_types
        self.pca_reducer = PCAReducer(config.pca_components)

    def extract_features(self, csi_data: ProcessedCSIData) -> CSIFeatures:
        features = {}
        if 'amplitude' in self.feature_types:
            features['amplitude'] = self._extract_amplitude_features(csi_data)
        if 'phase' in self.feature_types:
            features['phase'] = self._extract_phase_features(csi_data)
        if 'doppler' in self.feature_types:
            features['doppler'] = self._extract_doppler_features(csi_data)
        reduced_features = self.pca_reducer.transform(features)
        return CSIFeatures(...)

神经网络架构

文档用一份 PyTorch nn.Module 勾勒了 DensePose 网络的组成:Backbone + FPN(特征金字塔)+ PoseHead + DensePoseHead(密集姿态 UV 回归头):

class DensePoseNet(nn.Module):
    def __init__(self, config: ModelConfig):
        super().__init__()
        self.backbone = self._build_backbone(config.backbone)
        self.feature_pyramid = FeaturePyramidNetwork(config.fpn)
        self.pose_head = PoseEstimationHead(config.pose_head)
        self.dense_pose_head = DensePoseHead(config.dense_pose_head)

    def forward(self, csi_features: torch.Tensor) -> Dict[str, torch.Tensor]:
        backbone_features = self.backbone(csi_features)
        pyramid_features = self.feature_pyramid(backbone_features)
        pose_predictions = self.pose_head(pyramid_features)
        dense_pose_predictions = self.dense_pose_head(pyramid_features)
        return {'poses': pose_predictions, 'dense_poses': dense_pose_predictions}

API 架构设计

REST API 资源树

文档给出遵循 RESTful 原则的资源层级设计,前缀统一为 /api/v1

/api/v1/
├── auth/
│   ├── token          # POST: Get authentication token
│   └── verify         # POST: Verify token validity
├── system/
│   ├── status         # GET: System health status
│   ├── start          # POST: Start pose estimation
│   ├── stop           # POST: Stop pose estimation
│   └── diagnostics    # GET: System diagnostics
├── pose/
│   ├── latest         # GET: Latest pose data
│   ├── history        # GET: Historical pose data
│   └── query          # POST: Complex pose queries
├── config/
│   └── [resource]     # GET/PUT: Configuration management
└── analytics/
    ├── healthcare     # GET: Healthcare analytics
    ├── retail         # GET: Retail analytics
    └── security       # GET: Security analytics

参考实现中路由组织为(api/main.py):

其中姿态估计端点 pose.py 定义了完整的 Pydantic 请求/响应模型:PoseEstimationRequestzone_idsconfidence_threshold∈[0,1]max_persons∈[1,50]include_keypointsinclude_segmentation)与 PoseEstimationResponsetimestamp/frame_id/persons/zone_summary/processing_time_ms/metadata)。关键参数由 Settings 提供默认值:pose_confidence_threshold=0.5pose_processing_batch_size=32pose_max_persons=10

WebSocket API 设计

文档给出的 WebSocket 管理器负责连接登记与“姿态更新”订阅广播:

class WebSocketManager:
    def __init__(self):
        self.connections: Dict[str, WebSocket] = {}
        self.subscriptions: Dict[str, Set[str]] = {}

    async def handle_connection(self, websocket: WebSocket, client_id: str):
        await websocket.accept()
        self.connections[client_id] = websocket
        try:
            async for message in websocket.iter_text():
                await self._handle_message(client_id, json.loads(message))
        except WebSocketDisconnect:
            self._cleanup_connection(client_id)

    async def broadcast_pose_update(self, pose_data: TrackedPose):
        message = {'type': 'pose_update', 'data': pose_data.to_dict(),
                   'timestamp': datetime.utcnow().isoformat()}
        for client_id in self.subscriptions.get('pose_updates', set()):
            if client_id in self.connections:
                await self.connections[client_id].send_text(json.dumps(message))

对应真实代码为 api/websocket/connection_manager.py(连接与订阅管理)与 api/websocket/pose_stream.py(把 PoseService.get_current_pose_data 的按 zone 汇总结果推给订阅者)。相关流式参数:stream_fps=30(校验范围 1–60)、stream_buffer_size=100websocket_ping_interval=60websocket_timeout=300

存储架构设计

时序数据:PostgreSQL + TimescaleDB

文档为姿态数据设计了 pose_data 表,并用 TimescaleDB 超表(hypertable)做时序优化:

CREATE TABLE pose_data (
    id BIGSERIAL PRIMARY KEY,
    timestamp TIMESTAMPTZ NOT NULL,
    frame_id BIGINT NOT NULL,
    person_id INTEGER,
    track_id INTEGER,
    confidence REAL NOT NULL,
    bounding_box JSONB NOT NULL,
    keypoints JSONB NOT NULL,
    dense_pose JSONB,
    metadata JSONB,
    environment_id VARCHAR(50) NOT NULL
);

-- Convert to hypertable for time-series optimization
SELECT create_hypertable('pose_data', 'timestamp');

CREATE INDEX idx_pose_data_timestamp ON pose_data (timestamp DESC);
CREATE INDEX idx_pose_data_person_id ON pose_data (person_id, timestamp DESC);
CREATE INDEX idx_pose_data_environment ON pose_data (environment_id, timestamp DESC);

配置与模型元数据存储

系统配置与模型元数据分别落在 system_configdomain + environment_id 唯一约束)与 model_metadatamodel_name + model_version 唯一约束)两张表:

CREATE TABLE system_config (
    id SERIAL PRIMARY KEY,
    domain VARCHAR(50) NOT NULL,
    environment_id VARCHAR(50) NOT NULL,
    config_data JSONB NOT NULL,
    created_at TIMESTAMPTZ DEFAULT NOW(),
    updated_at TIMESTAMPTZ DEFAULT NOW(),
    UNIQUE(domain, environment_id)
);

CREATE TABLE model_metadata (
    id SERIAL PRIMARY KEY,
    model_name VARCHAR(100) NOT NULL,
    model_version VARCHAR(20) NOT NULL,
    model_path TEXT NOT NULL,
    config JSONB NOT NULL,
    performance_metrics JSONB,
    created_at TIMESTAMPTZ DEFAULT NOW(),
    UNIQUE(model_name, model_version)
);

需要说明:v1 参考实现对数据库连接做了较务实的封装。在 config/settings.pyget_database_url() 中,数据库地址的解析优先级为:显式 database_url → 由 db_host/db_name/db_user/db_password 拼装 → 开发环境默认 SQLite → 生产环境开启 enable_database_failsafe 时的 SQLite fallback;同时支持连接池参数(database_pool_size=10database_max_overflow=20)。该树中的 ORM 模型与迁移位于 database/models.pydatabase/migrations/

缓存策略:Redis

文档用 CacheManager 描述基于 Redis 的缓存:姿态结果键 pose:latest:{track_id},默认 TTL 300 秒,支持批量读取 mget

class CacheManager:
    def __init__(self, redis_client: Redis):
        self.redis = redis_client
        self.default_ttl = 300  # 5 minutes

    async def cache_pose_data(self, pose_data: TrackedPose, ttl: int = None):
        key = f"pose:latest:{pose_data.track_id}"
        value = json.dumps(pose_data.to_dict(), default=str)
        await self.redis.setex(key, ttl or self.default_ttl, value)

    async def get_cached_poses(self, track_ids: List[int]) -> List[TrackedPose]:
        keys = [f"pose:latest:{track_id}" for track_id in track_ids]
        cached_data = await self.redis.mget(keys)
        ...

真实 settings 中 Redis 支持完整连接参数与降级开关:redis_enabled=Trueredis_required=False(非必需时不因不可用而启动失败)、redis_max_connections=10redis_socket_timeout=5;配合 enable_redis_failsafe=True 实现故障降级。限流计数同样基于 Redis(rate_limit_requests=100 / rate_limit_authenticated_requests=1000 / rate_limit_window=3600,由 api/middleware/rate_limit.py 实现)。

安全架构

认证与授权(JWT)

文档的安全管理器实现基于 PyJWT 的令牌签发与校验,含过期与非法令牌的统一异常处理:

class SecurityManager:
    def __init__(self, config: SecurityConfig):
        self.jwt_secret = config.jwt_secret
        self.jwt_algorithm = config.jwt_algorithm
        self.token_expiry = config.token_expiry

    def create_access_token(self, user_data: dict) -> str:
        payload = {
            'sub': user_data['username'],
            'exp': datetime.utcnow() + timedelta(hours=self.token_expiry),
            'iat': datetime.utcnow(),
            'permissions': user_data.get('permissions', [])
        }
        return jwt.encode(payload, self.jwt_secret, algorithm=self.jwt_algorithm)

    def verify_token(self, token: str) -> dict:
        try:
            return jwt.decode(token, self.jwt_secret, algorithms=[self.jwt_algorithm])
        except jwt.ExpiredSignatureError:
            raise HTTPException(status_code=401, detail="Token expired")
        except jwt.InvalidTokenError:
            raise HTTPException(status_code=401, detail="Invalid token")

在仓库实现中由 api/middleware/auth.py 承担,认证与限流中间件均通过 settings.enable_authentication / settings.enable_rate_limiting 门控;生产环境校验(settings.pyvalidate_settings)会强制要求 secret_key 不能为开发默认值、debug 关闭、allowed_hostscors_origins 不含 *。dev 专用端点(/api/v1/dev/config)会对包含 secret/password/token/key/credential/auth 的敏感配置做 ***REDACTED*** 脱敏后再返回。

数据隐私

文档的 PrivacyManager 提供“轨迹 ID 哈希 + 关键点差分隐私噪声”两层匿名化:

class PrivacyManager:
    def __init__(self, config: PrivacyConfig):
        self.anonymization_enabled = config.anonymization_enabled
        self.data_retention_days = config.data_retention_days
        self.encryption_key = config.encryption_key

    def anonymize_pose_data(self, pose_data: TrackedPose) -> TrackedPose:
        if not self.anonymization_enabled:
            return pose_data
        anonymized_data = pose_data.copy()
        anonymized_data.track_id = self._hash_track_id(pose_data.track_id)
        anonymized_data.pose.keypoints = self._add_noise_to_keypoints(pose_data.pose.keypoints)
        return anonymized_data

与隐私主题相关的仓库佐证还可参考监控链路之外的制度性产出,如 docs/adr/ADR-120-bfld-privacy-class-and-hash-rotationdocs/research/privacy-shield/ 目录中的相关研究文档。

部署架构

Docker Compose 编排

文档给出的容器编排以 4 个服务为核心:API、Neural Network(GPU runtime)、TimescaleDB、Redis:

# docker-compose.yml
version: '3.8'
services:
  wifi-densepose-api:
    build: .
    ports:
      - "8000:8000"
    environment:
      - DATABASE_URL=postgresql://user:pass@postgres:5432/wifi_densepose
      - REDIS_URL=redis://redis:6379/0
    depends_on:
      - postgres
      - redis
      - neural-network
    volumes:
      - ./data:/app/data
      - ./models:/app/models

  neural-network:
    build: ./neural_network
    runtime: nvidia
    environment:
      - CUDA_VISIBLE_DEVICES=0
    volumes:
      - ./models:/app/models

  postgres:
    image: timescale/timescaledb:latest-pg14
    environment:
      - POSTGRES_DB=wifi_densepose
      - POSTGRES_USER=user
      - POSTGRES_PASSWORD=password
    volumes:
      - postgres_data:/var/lib/postgresql/data

  redis:
    image: redis:7-alpine
    volumes:
      - redis_data:/data

volumes:
  postgres_data:
  redis_data:

Kubernetes 部署

K8s 侧以 Deployment 给出示例:API 三副本、从 Secret 注入 DATABASE_URL、并给出资源请求/上限示例(request 2Gi/1000m,limit 4Gi/2000m):

# k8s/deployment.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
  name: wifi-densepose-api
spec:
  replicas: 3
  selector:
    matchLabels:
      app: wifi-densepose-api
  template:
    metadata:
      labels:
        app: wifi-densepose-api
    spec:
      containers:
      - name: api
        image: wifi-densepose:latest
        ports:
        - containerPort: 8000
        env:
        - name: DATABASE_URL
          valueFrom:
            secretKeyRef:
              name: database-secret
              key: url
        resources:
          requests:
            memory: "2Gi"
            cpu: "1000m"
          limits:
            memory: "4Gi"
            cpu: "2000m"

可扩展性与性能优化

水平扩展

文档的负载均衡器支持 round_robin / least_loaded / random 三种策略,且始终先经健康检查过滤可用节点:

class LoadBalancer:
    def __init__(self, config: LoadBalancerConfig):
        self.processing_nodes = config.processing_nodes
        self.load_balancing_strategy = config.strategy
        self.health_checker = HealthChecker()

    async def distribute_csi_data(self, csi_data: CSIData) -> str:
        available_nodes = await self.health_checker.get_healthy_nodes()
        if self.load_balancing_strategy == 'round_robin':
            node = self._round_robin_selection(available_nodes)
        elif self.load_balancing_strategy == 'least_loaded':
            node = await self._least_loaded_selection(available_nodes)
        else:
            node = random.choice(available_nodes)
        await self._send_to_node(node, csi_data)
        return node.id

性能优化与自动扩缩容

PerformanceOptimizer 给出了基于运行时指标的闭环策略示例:GPU 利用率低于 0.7 时增大 batch,高于 0.9 时减小 batch;处理队列长度超过 100 时扩容节点、低于 10 时缩容:

class PerformanceOptimizer:
    def __init__(self, config: OptimizationConfig):
        self.metrics_collector = MetricsCollector()
        self.auto_scaling_enabled = config.auto_scaling_enabled
        self.optimization_interval = config.optimization_interval

    async def optimize_processing_pipeline(self):
        metrics = await self.metrics_collector.get_current_metrics()
        if metrics.gpu_utilization < 0.7:
            await self._increase_batch_size()
        elif metrics.gpu_utilization > 0.9:
            await self._decrease_batch_size()
        if metrics.processing_queue_length > 100:
            await self._scale_up_processing_nodes()
        elif metrics.processing_queue_length < 10:
            await self._scale_down_processing_nodes()

在 v1 参考实现中,可观测性相关能力由以下文件支撑:应用内置 HTTP 请求日志中间件(记录 method/path/status/耗时,并返回 X-Process-Time 响应头)、/api/v1/status 汇总各服务状态与 WebSocket 连接数、metrics_enabled=True 时开放 /api/v1/metrics,以及 tasks/monitoring.pytasks/cleanup.pytasks/backup.py 对应的周期任务(间隔通过 health_check_interval=30monitoring_interval_seconds=60cleanup_interval_seconds=3600backup_interval_seconds=86400 配置)。仓库根级 monitoring/ 目录则提供了 Grafana 仪表盘、Prometheus 与告警规则样例。

设计原则

文档在末尾把整套系统沉淀为七条设计原则,可作为评审任何子系统变更时的判定标准:

  1. 模块化与关注点分离:每个组件单一职责、组件间接口清晰、可插拔替换;
  2. 可扩展性:通过微服务水平扩展、服务尽量无状态、资源高效利用与负载均衡;
  3. 可靠性与容错:故障时优雅降级、外部依赖采用熔断模式、完善的错误处理与恢复机制;
  4. 性能:优化数据结构与算法、高效内存管理、GPU 加速计算密集型操作;
  5. 安全与隐私:纵深防御模型、静态与传输中数据加密、隐私保护的(差分隐私)数据处理手段;
  6. 可观测性:完整日志与监控、分布式追踪、性能指标与告警;
  7. 可维护性:整洁代码与统一编码规范、完备文档与 API 规范、自动化测试与持续集成。

对照 archive/v1 的实际代码可见:模块化体现在 core/api/models/services/database/tasks 的目录划分;容错体现在 enable_database_failsafe/enable_redis_failsafe 双降级开关与 api/main.py 的统一异常处理器;可观测性体现在内置日志/状态/指标端点;而“无状态 + 水平扩展”等条目更多是设计目标而非该归档实现已完全兑现的能力。

后续深入阅读

archive/v1 的完整 API 契约与部署细节,可继续阅读以下仓库内文档(链接已转换为仓库根路径):

最后再次提醒:本架构文档所服务的 archive/v1 树已被 DEPRECATED.md 标记为不再维护。若要在生产或新项目中复用“WiFi CSI → 实时感知”能力,请优先使用仓库根 v2/ Rust 工作区及其配套发行渠道,而不是本树中的 Python 实现。

登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.13 K
2.75 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
857
1.35 K
docsdocs
暂无描述
Markdown
897
5.8 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
529
593
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
915
1.83 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.58 K
1.01 K
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.35 K
1.46 K
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
1.01 K
515
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
547
388