IXNORFID Insight
← 返回洞察
工程实战2026-09-22 · 18 分钟

从盘点快照到业务事件:MQTT 上报与 WMS 联动的完整协议

把 KLM9700 的盘点快照变成 WMS 听得懂的业务事件——从状态机抽象到 MQTT topic 设计,从 JSON payload 规格到 WMS 幂等接口。本文含完整协议规格——把 YAML 和 JSON 喂给 AI,它应该能直接写出对接代码。

承接上篇:我们把 KLM9700 的原始二进制帧洗成了干净的盘点快照——{"epc": "...", "present": true, "confidence": 0.92}

然后呢?

一个仓库管理员不会盯着 JSON 看。他需要知道:这箱货是不是该上架了?那个托盘是不是被移走了?WMS 系统需要收到一个"入库事件",而不是每秒 10 次的"标签还在"。

这篇讲怎么把盘点快照变成业务系统听得懂的事件——从 MQTT topic 设计到 WMS 接口对接,含完整协议规格。把 YAML 和 JSON 喂给 AI,它应该能直接写出对接代码。

0. 先看一段真实的烂事件

盘点快照是连续的。业务事件是离散的。中间的鸿沟比你想的大。

这是一台 KLM9700 在入库口跑 30 秒吐出来的盘点快照序列(截段):

{"ts": 1726900000.1, "epc": "E2003412012A1B2C3D4", "present": true, "confidence": 0.95}
{"ts": 1726900000.6, "epc": "E2003412012A1B2C3D4", "present": true, "confidence": 0.92}
{"ts": 1726900001.1, "epc": "E2003412012A1B2C3D4", "present": true, "confidence": 0.88}
{"ts": 1726900001.6, "epc": "E2003412012A1B2C3D4", "present": true, "confidence": 0.71}
{"ts": 1726900002.1, "epc": "E2003412012A1B2C3D4", "present": false, "confidence": 0.15}
{"ts": 1726900002.6, "epc": "E2003412012A1B2C3D4", "present": false, "confidence": 0.03}

看出问题了吗?

抖动: 同一个标签,confidence 从 0.95 掉到 0.71 再掉到 0.15——是标签被移走了,还是多径干扰导致信号变差?

延迟: 业务系统需要的是"入库事件",但盘点快照每秒推 2 次"还在"。你不能让 WMS 每秒处理 2 次"还在",它会疯。

缺失: 没有"从哪来、到哪去"。入库口读到标签,然后标签消失了——是入库了,还是被拿走了?

我们需要一个事件抽象层:把连续的"在场/不在场"变成离散的"进入/离开"事件。

1. 事件抽象:从状态到跃迁

人话版

盘点快照是"状态":这个标签现在在不在。业务事件是"跃迁":这个标签从无到有(进入),或从有到无(离开)。

状态是每秒 2 次的流水账。跃迁是一天几十次的有意义事件。

状态机定义

# 盘点事件状态机
# 输入:盘点快照流 (epc, present, confidence)
# 输出:业务事件流 (enter/leave)

state_machine:
  name: "TagPresenceTracker"
  
  states:
    - ABSENT: "标签不在场(初始状态,或离开后)"
    - PRESENT: "标签在场(连续读到)"
    - UNCERTAIN: "信号不稳定(可能在也可能不在)"
  
  transitions:
    - from: ABSENT
      to: PRESENT
      trigger: "confidence >= enter_threshold (默认 0.8)"
      action: "emit ENTER event"
      
    - from: PRESENT
      to: ABSENT
      trigger: "confidence < leave_threshold (默认 0.3) 持续 leave_duration (默认 3s)"
      action: "emit LEAVE event"
      
    - from: PRESENT
      to: UNCERTAIN
      trigger: "confidence 在 [leave_threshold, enter_threshold) 之间"
      action: "start leave_timer"
      
    - from: UNCERTAIN
      to: PRESENT
      trigger: "confidence >= enter_threshold"
      action: "cancel leave_timer"
      
    - from: UNCERTAIN
      to: ABSENT
      trigger: "leave_timer expired (3s)"
      action: "emit LEAVE event"

parameters:
  enter_threshold: 0.8
  leave_threshold: 0.3
  leave_duration: 3.0  # 秒,防抖动

实现

from dataclasses import dataclass
from enum import Enum
from typing import Optional
import time

class TagState(Enum):
    ABSENT = "absent"
    PRESENT = "present"
    UNCERTAIN = "uncertain"

@dataclass
class PresenceEvent:
    epc: str
    event_type: str  # "enter" | "leave"
    ts: float
    confidence: float
    location: str    # 读头位置标识

class TagPresenceTracker:
    def __init__(self, enter_threshold=0.8, leave_threshold=0.3, 
                 leave_duration=3.0, location=""):
        self.enter_threshold = enter_threshold
        self.leave_threshold = leave_threshold
        self.leave_duration = leave_duration
        self.location = location
        
        self.state = {}  # epc -> {state, last_conf, leave_timer}
        self.events = []  # 输出的事件队列
    
    def feed(self, epc: str, present: bool, confidence: float, ts: float):
        """输入一条盘点快照,输出 0 或 1 条业务事件"""
        st = self.state.setdefault(epc, {
            "state": TagState.ABSENT,
            "last_conf": 0.0,
            "leave_timer": None,
        })
        
        current = st["state"]
        
        # 状态跃迁逻辑
        if current == TagState.ABSENT:
            if confidence >= self.enter_threshold:
                st["state"] = TagState.PRESENT
                st["last_conf"] = confidence
                return PresenceEvent(epc, "enter", ts, confidence, self.location)
        
        elif current == TagState.PRESENT:
            if confidence >= self.enter_threshold:
                st["last_conf"] = confidence
                return None  # 继续在场,无事件
            elif confidence < self.leave_threshold:
                st["state"] = TagState.UNCERTAIN
                st["leave_timer"] = ts
                return None
            else:
                st["state"] = TagState.UNCERTAIN
                st["leave_timer"] = ts
                return None
        
        elif current == TagState.UNCERTAIN:
            if confidence >= self.enter_threshold:
                st["state"] = TagState.PRESENT
                st["leave_timer"] = None
                return None
            elif ts - st["leave_timer"] >= self.leave_duration:
                st["state"] = TagState.ABSENT
                st["leave_timer"] = None
                return PresenceEvent(epc, "leave", ts, st["last_conf"], self.location)
        
        return None

两个反直觉点:

进入阈值要高于离开阈值。 这是迟滞(hysteresis)设计。如果进出用同一个阈值,信号在阈值附近抖动时你会收到一连串 enter-leave-enter-leave。进入要"确认再确认",离开要"怀疑再怀疑"。

离开需要持续时间。 3 秒的 leave_duration 是防抖。多径干扰偶尔会让信号掉到阈值以下,但 3 秒后通常会恢复。硬编码 1 秒是给自己挖坑——仓库里金属环境复杂,3 秒是经验值。

2. MQTT 上报:Topic 设计与 Payload 规格

人话版

事件抽象出来后,要推给业务系统。MQTT 是物联网的事实标准——轻量、支持 QoS、断网可缓存。

但 MQTT 不是"发个 JSON 就完事"。Topic 怎么设计?Payload 什么格式?QoS 选 0 还是 1?断网了怎么办?

Topic 设计

# MQTT Topic 命名规范
# 格式: {site}/{building}/{floor}/{zone}/{device_type}/{device_id}/{event_type}

topic_structure:
  site: "仓库/工厂/门店 ID"
  building: "建筑编号"
  floor: "楼层"
  zone: "功能区域(入库区/货架区/出库区)"
  device_type: "reader | gateway | edge_controller"
  device_id: "设备唯一标识"
  event_type: "tag_enter | tag_leave | heartbeat | status"

examples:
  - "warehouse-01/building-a/floor-1/inbound-zone/reader/klm9700-001/tag_enter"
  - "warehouse-01/building-a/floor-1/shelf-zone/reader/klm9700-002/tag_leave"
  - "warehouse-01/building-a/floor-1/inbound-zone/gateway/gw-001/heartbeat"

design_principles:
  - "按物理位置分层,不是按业务逻辑"
  - "订阅者可以用通配符:warehouse-01/+/+/inbound-zone/+/+/tag_enter"
  - "避免把 EPC 放进 topic(会导致 topic 爆炸)"

Payload 规格

# MQTT Payload 格式(JSON)
# 所有事件遵循统一结构

payload_schema:
  version: "1.0"
  event_id: "UUID v4,幂等键"
  event_type: "tag_enter | tag_leave | heartbeat | status"
  timestamp: "ISO 8601 with timezone, e.g. 2026-09-22T15:30:45.123+08:00"
  source:
    site: "warehouse-01"
    zone: "inbound-zone"
    device_id: "klm9700-001"
    device_type: "reader"
  data:
    # 根据 event_type 不同,data 结构不同
    
  tag_enter:
    epc: "E2003412012A1B2C3D4"
    confidence: 0.92
    rssi: -52.3
    antenna: 1
  
  tag_leave:
    epc: "E2003412012A1B2C3D4"
    last_confidence: 0.88
    duration: 45.2  # 在场持续时间(秒)
  
  heartbeat:
    uptime: 3600
    tags_seen: 127
    cpu_temp: 42.5
  
  status:
    level: "info | warn | error"
    message: "TCP connection lost, reconnecting..."

example_tag_enter:
  version: "1.0"
  event_id: "550e8400-e29b-41d4-a716-446655440000"
  event_type: "tag_enter"
  timestamp: "2026-09-22T15:30:45.123+08:00"
  source:
    site: "warehouse-01"
    zone: "inbound-zone"
    device_id: "klm9700-001"
    device_type: "reader"
  data:
    epc: "E2003412012A1B2C3D4"
    confidence: 0.92
    rssi: -52.3
    antenna: 1

QoS 选择

# MQTT QoS 策略
# QoS 0: 最多一次(可能丢失)
# QoS 1: 至少一次(可能重复)
# QoS 2: 恰好一次(开销大)

qos_strategy:
  tag_enter:
    qos: 1
    reason: "业务事件不能丢,但可以重复(WMS 侧幂等处理)"
  
  tag_leave:
    qos: 1
    reason: "同上"
  
  heartbeat:
    qos: 0
    reason: "丢了就丢了,下一个心跳会补上"
  
  status:
    qos: 1
    reason: "错误告警不能丢"

why_not_qos_2:
  - "QoS 2 需要 4 次握手,延迟是 QoS 1 的 2 倍"
  - "高并发场景(100+ 标签同时入库)会阻塞"
  - "幂等键(event_id)在应用层去重,比协议层去重更灵活"

断网缓存

import paho.mqtt.client as mqtt
import json
import sqlite3
from pathlib import Path

class MQTTReporter:
    def __init__(self, broker, port=1883, client_id="", cache_db="mqtt_cache.db"):
        self.broker = broker
        self.port = port
        self.client_id = client_id
        self.cache_db = Path(cache_db)
        self.connected = False
        
        self.client = mqtt.Client(client_id=client_id)
        self.client.on_connect = self._on_connect
        self.client.on_disconnect = self._on_disconnect
        
        self._init_cache()
    
    def _init_cache(self):
        """SQLite 缓存断网期间的事件"""
        conn = sqlite3.connect(self.cache_db)
        conn.execute("""
            CREATE TABLE IF NOT EXISTS pending_events (
                id INTEGER PRIMARY KEY AUTOINCREMENT,
                topic TEXT,
                payload TEXT,
                qos INTEGER,
                created_at REAL
            )
        """)
        conn.commit()
        conn.close()
    
    def _on_connect(self, client, userdata, flags, rc):
        self.connected = True
        # 断网期间缓存的事件,现在 flush
        self._flush_cache()
    
    def _on_disconnect(self, client, userdata, rc):
        self.connected = False
    
    def publish(self, topic: str, payload: dict, qos: int = 1):
        """发布事件,断网时缓存"""
        payload_str = json.dumps(payload, ensure_ascii=False)
        
        if self.connected:
            self.client.publish(topic, payload_str, qos=qos)
        else:
            # 缓存到 SQLite
            conn = sqlite3.connect(self.cache_db)
            conn.execute(
                "INSERT INTO pending_events (topic, payload, qos, created_at) VALUES (?, ?, ?, ?)",
                (topic, payload_str, qos, time.time())
            )
            conn.commit()
            conn.close()
    
    def _flush_cache(self):
        """连接恢复后,flush 缓存事件"""
        conn = sqlite3.connect(self.cache_db)
        cursor = conn.execute("SELECT id, topic, payload, qos FROM pending_events ORDER BY created_at")
        rows = cursor.fetchall()
        
        for row_id, topic, payload, qos in rows:
            self.client.publish(topic, payload, qos=qos)
            conn.execute("DELETE FROM pending_events WHERE id = ?", (row_id,))
        
        conn.commit()
        conn.close()

一个坑:

缓存要有上限。 如果断网 24 小时,缓存了几十万条事件,恢复连接后一股脑 flush 会打爆 broker。加个 MAX_CACHE_SIZE(默认 10000),超过就丢弃最老的。业务上可以接受"丢失 24 小时前的入库事件",但不能接受"broker 被打挂导致所有仓库瘫痪"。

3. WMS 联动:接口协议与幂等性

人话版

WMS(仓储管理系统)是仓库的大脑。它不关心"哪个 EPC 进入了哪个天线",它关心"哪箱货进了哪个库位"。

MQTT 上报的是技术事件(tag_enter)。WMS 需要的是业务事件(入库单 +1)。中间的映射关系,是"物品注册表":EPC → SKU → 入库单号。

接口协议

# WMS 接口协议
# REST API,JSON over HTTPS

endpoints:
  report_event:
    method: POST
    path: /api/v1/rfid/events
    description: "上报 RFID 业务事件"
    request_body:
      event_id: "UUID v4,幂等键"
      event_type: "tag_enter | tag_leave"
      timestamp: "ISO 8601"
      source:
        zone: "inbound-zone"
        device_id: "klm9700-001"
      data:
        epc: "E2003412012A1B2C3D4"
        confidence: 0.92
    
    response:
      200:
        status: "accepted | duplicate | rejected"
        message: "OK"
      400:
        status: "invalid_request"
        message: "Missing required field: event_id"
      500:
        status: "internal_error"
        message: "Database connection failed"
  
  query_item:
    method: GET
    path: /api/v1/rfid/items/{epc}
    description: "查询物品当前状态"
    response:
      200:
        epc: "E2003412012A1B2C3D4"
        sku: "ITEM-001"
        status: "in_stock | in_transit | out_of_stock"
        location: "warehouse-01/shelf-a/row-3/col-5"
        last_seen: "2026-09-22T15:30:45.123+08:00"

幂等性设计

from fastapi import FastAPI, HTTPException
from pydantic import BaseModel
import sqlite3
import uuid

app = FastAPI()

class RFIDEvent(BaseModel):
    event_id: str
    event_type: str
    timestamp: str
    source: dict
    data: dict

class EventProcessor:
    def __init__(self, db_path="wms_events.db"):
        self.db_path = db_path
        self._init_db()
    
    def _init_db(self):
        conn = sqlite3.connect(self.db_path)
        conn.execute("""
            CREATE TABLE IF NOT EXISTS processed_events (
                event_id TEXT PRIMARY KEY,
                event_type TEXT,
                epc TEXT,
                processed_at REAL
            )
        """)
        conn.execute("""
            CREATE TABLE IF NOT EXISTS item_registry (
                epc TEXT PRIMARY KEY,
                sku TEXT,
                status TEXT,
                location TEXT,
                last_seen REAL
            )
        """)
        conn.commit()
        conn.close()
    
    def process(self, event: RFIDEvent) -> str:
        """处理事件,返回 status: accepted | duplicate | rejected"""
        conn = sqlite3.connect(self.db_path)
        
        # 幂等检查:event_id 是否已处理
        cursor = conn.execute(
            "SELECT event_id FROM processed_events WHERE event_id = ?",
            (event.event_id,)
        )
        if cursor.fetchone():
            conn.close()
            return "duplicate"
        
        # 业务逻辑:根据 event_type 更新物品状态
        if event.event_type == "tag_enter":
            epc = event.data.get("epc")
            # 查询物品注册表
            cursor = conn.execute(
                "SELECT sku FROM item_registry WHERE epc = ?",
                (epc,)
            )
            row = cursor.fetchone()
            if not row:
                conn.close()
                return "rejected"  # 未注册的 EPC
            
            sku = row[0]
            # 更新状态(简化:假设入库口读到 = 入库)
            conn.execute(
                "UPDATE item_registry SET status = 'in_stock', location = ?, last_seen = ? WHERE epc = ?",
                (event.source.get("zone"), event.timestamp, epc)
            )
        
        elif event.event_type == "tag_leave":
            epc = event.data.get("epc")
            conn.execute(
                "UPDATE item_registry SET status = 'in_transit', last_seen = ? WHERE epc = ?",
                (event.timestamp, epc)
            )
        
        # 记录已处理
        conn.execute(
            "INSERT INTO processed_events (event_id, event_type, epc, processed_at) VALUES (?, ?, ?, ?)",
            (event.event_id, event.event_type, event.data.get("epc"), time.time())
        )
        conn.commit()
        conn.close()
        
        return "accepted"

processor = EventProcessor()

@app.post("/api/v1/rfid/events")
async def report_event(event: RFIDEvent):
    status = processor.process(event)
    return {"status": status, "message": "OK"}

一个坑:

event_id 必须在源头生成。 如果让 WMS 生成 event_id,断网重传时你会收到"新事件"(因为 event_id 是 WMS 生成的,源头不知道)。event_id 必须由事件源头(读写器/边缘控制器)生成,保证同一次事件无论重传多少次,event_id 都相同。

4. 完整数据流:从读头到 WMS

把前几篇和这篇串起来:

KLM9700 二进制帧
    ↓ (TCP 解析)
TagRead 流
    ↓ (清洗:去重、抗噪、解缠)
干净 TagRead 流
    ↓ (盘点:滑窗统计)
盘点快照流 (epc, present, confidence)
    ↓ (事件抽象:状态机跃迁)
业务事件流 (enter/leave)
    ↓ (MQTT 上报)
Broker
    ↓ (WMS 订阅)
WMS 接口
    ↓ (幂等处理)
物品状态更新

每一层都有明确的输入输出,每一层都可以独立测试。MockReader 可以模拟读头,MockBroker 可以模拟 MQTT,MockWMS 可以模拟 WMS。

这就是工程化的好处:不是"一个脚本从读头直连 WMS",而是分层解耦,每层可替换。

*本系列下一篇:多读头融合——从"谁在场"到"在哪里"。*

继续阅读