从盘点快照到业务事件: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: 1QoS 选择
# 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",而是分层解耦,每层可替换。
*本系列下一篇:多读头融合——从"谁在场"到"在哪里"。*