从盘点快照到业务事件:MQTT 上报与 WMS 联动的完整协议
把 KLM9700 的盘点快照变成 WMS 听得懂的业务事件——从状态机抽象到 MQTT topic 设计,从 JSON payload 规格到 WMS 幂等接口。本文含完整协议规格——把 YAML 和 JSON 喂给 AI,它应该能直接写出对接代码。
承接上篇:我们把 KLM9700 的原始二进制帧洗成了干净的盘点快照——<code>{"epc": "...", "present": true, "confidence": 0.92}</code>。
然后呢?
一个仓库管理员不会盯着 JSON 看。他需要知道:这箱货是不是该上架了?那个托盘是不是被移走了?WMS 系统需要收到一个"入库事件",而不是每秒 10 次的"标签还在"。
这篇讲怎么把盘点快照变成业务系统听得懂的事件——从 MQTT topic 设计到 WMS 接口对接,含完整协议规格。把 YAML 和 JSON 喂给 AI,它应该能直接写出对接代码。
0. 先看一段真实的烂事件
盘点快照是连续的。业务事件是离散的。中间的鸿沟比你想的大。
这是一台 KLM9700 在入库口跑 30 秒吐出来的盘点快照序列(截段):
<pre><code class="lang-json">{"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}</code></pre>
看出问题了吗?
<strong>抖动:</strong> 同一个标签,confidence 从 0.95 掉到 0.71 再掉到 0.15——是标签被移走了,还是多径干扰导致信号变差?
<strong>延迟:</strong> 业务系统需要的是"入库事件",但盘点快照每秒推 2 次"还在"。你不能让 WMS 每秒处理 2 次"还在",它会疯。
<strong>缺失:</strong> 没有"从哪来、到哪去"。入库口读到标签,然后标签消失了——是入库了,还是被拿走了?
我们需要一个<strong>事件抽象层</strong>:把连续的"在场/不在场"变成离散的"进入/离开"事件。
1. 事件抽象:从状态到跃迁
人话版
盘点快照是"状态":这个标签现在在不在。业务事件是"跃迁":这个标签从无到有(进入),或从有到无(离开)。
状态是每秒 2 次的流水账。跃迁是一天几十次的有意义事件。
状态机定义
<pre><code class="lang-yaml"># 盘点事件状态机 # 输入:盘点快照流 (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 # 秒,防抖动</code></pre>
实现
<pre><code class="lang-python">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</code></pre>
两个反直觉点:
<strong>进入阈值要高于离开阈值。</strong> 这是迟滞(hysteresis)设计。如果进出用同一个阈值,信号在阈值附近抖动时你会收到一连串 enter-leave-enter-leave。进入要"确认再确认",离开要"怀疑再怀疑"。
<strong>离开需要持续时间。</strong> 3 秒的 leave_duration 是防抖。多径干扰偶尔会让信号掉到阈值以下,但 3 秒后通常会恢复。硬编码 1 秒是给自己挖坑——仓库里金属环境复杂,3 秒是经验值。
2. MQTT 上报:Topic 设计与 Payload 规格
人话版
事件抽象出来后,要推给业务系统。MQTT 是物联网的事实标准——轻量、支持 QoS、断网可缓存。
但 MQTT 不是"发个 JSON 就完事"。Topic 怎么设计?Payload 什么格式?QoS 选 0 还是 1?断网了怎么办?
Topic 设计
<pre><code class="lang-yaml"># 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 爆炸)"</code></pre>
Payload 规格
<pre><code class="lang-yaml"># 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</code></pre>
QoS 选择
<pre><code class="lang-yaml"># 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)在应用层去重,比协议层去重更灵活"</code></pre>
断网缓存
<pre><code class="lang-python">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()</code></pre>
一个坑:
<strong>缓存要有上限。</strong> 如果断网 24 小时,缓存了几十万条事件,恢复连接后一股脑 flush 会打爆 broker。加个 <code>MAX_CACHE_SIZE</code>(默认 10000),超过就丢弃最老的。业务上可以接受"丢失 24 小时前的入库事件",但不能接受"broker 被打挂导致所有仓库瘫痪"。
3. WMS 联动:接口协议与幂等性
人话版
WMS(仓储管理系统)是仓库的大脑。它不关心"哪个 EPC 进入了哪个天线",它关心"哪箱货进了哪个库位"。
MQTT 上报的是技术事件(tag_enter)。WMS 需要的是业务事件(入库单 +1)。中间的映射关系,是"物品注册表":EPC → SKU → 入库单号。
接口协议
<pre><code class="lang-yaml"># 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"</code></pre>
幂等性设计
<pre><code class="lang-python">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"}</code></pre>
一个坑:
<strong>event_id 必须在源头生成。</strong> 如果让 WMS 生成 event_id,断网重传时你会收到"新事件"(因为 event_id 是 WMS 生成的,源头不知道)。event_id 必须由事件源头(读写器/边缘控制器)生成,保证同一次事件无论重传多少次,event_id 都相同。
4. 完整数据流:从读头到 WMS
把前几篇和这篇串起来:
<pre><code>KLM9700 二进制帧 ↓ (TCP 解析) TagRead 流 ↓ (清洗:去重、抗噪、解缠) 干净 TagRead 流 ↓ (盘点:滑窗统计) 盘点快照流 (epc, present, confidence) ↓ (事件抽象:状态机跃迁) 业务事件流 (enter/leave) ↓ (MQTT 上报) Broker ↓ (WMS 订阅) WMS 接口 ↓ (幂等处理) 物品状态更新</code></pre>
每一层都有明确的输入输出,每一层都可以独立测试。MockReader 可以模拟读头,MockBroker 可以模拟 MQTT,MockWMS 可以模拟 WMS。
这就是工程化的好处:不是"一个脚本从读头直连 WMS",而是分层解耦,每层可替换。
*本系列下一篇:多读头融合——从"谁在场"到"在哪里"。*