断网这件事,在工厂现场几乎是躲不掉的。车间里一台注塑机或者一条产线,数据全靠工控网关往上送,网关和平台之间的链路一断,前端 MES 大屏上要么出现一大段空洞,要么曲线直接断掉。这时候如果网关只是把数据缓存在内存里,重启一次就全没了;如果发消息和写库两件事没有原子保证,还可能出现“业务数据已经入库、消息却没发出去”这种让人抓狂的问题。今天想聊的 Outbox 模式,就是用来解决“断网期间数据不丢、恢复之后自动补齐”这套问题的成熟思路。它原本是微服务架构里的老面孔,用来解决本地事务和远程消息的一致性问题;我把它平移到了工控网关上,实测下来比传统的“断线重连 + 内存队列”方案稳得多。
这篇内容适合做边缘网关、物联网采集器、SCADA 前置机的朋友参考,也适合打算自己搭一套可靠采集链路的嵌入式或后端开发者。接下来我按“问题拆解、表结构设计、代码实现、实测避坑、方案取舍”这个顺序讲,尽量把为什么这么做说透。
1. 断网丢数据这件事,到底难在哪
1.1 现场的三种典型丢数据场景
很多工控网关向上对接的是 MQTT Broker 或平台 HTTP 接口,向下对接 PLC、仪表、传感器。最典型的丢数据场景其实就那么几类:
- 网络闪断。车间网络交换机重启、无线 AP 漫游、4G 信号瞬时丢失,断个几十秒或者几分钟。平台侧 MQTT 连接断开,网关本地队列如果只放在内存里,断点一过消息就没了。
- 长时离线。厂区停电检修、运营商基站故障,断网可能持续几小时甚至一整天。此时网关本地如果不做持久化,等网络恢复时平台拿到的就是从重连开始的新数据,断网期间的历史数据变成空白。
- 网关重启和掉电。很多网关部署在生产现场,电源并不稳定。如果断网期间恰好赶上掉电重启,内存队列里积压的数据直接清零;如果系统是边写边发的裸逻辑,还可能只发了一半,平台收到一条不完整或缺失的记录。
你可能说,MQTT 本身有 QoS 1 和 QoS 2 啊。确实,QoS 能保证客户端和 Broker 之间的送达语义,但它成立的默认前提是“网关这边网络是通的,Broker 是可访问的”。在断网的窗口里,客户端本地缓存策略就成了决定生死的东西。而且很多 MQTT 客户端库的本地队列是写在一段固定大小的内存 buffer 里,消息量一大就覆盖,覆盖掉的旧消息等于丢失。换句话说,协议层面的 QoS 解决不了“本地缓存不可靠”的问题,真正要下功夫的地方在网关内部。
1.2 传统缓存方案的三个软肋
我最早做网关的时候,用的也是比较朴素的办法:采集到数据之后塞进一个环形队列,后台线程连上网络就往平台推。这套方案在实验室里怎么测都没问题,一上现场就暴露软肋。
- 缓存容量有限,断网时间稍长就溢出。环形队列设计成多少条都不合适,条数少了大流量下不够用,条数多了内存吃不消。而且一旦溢出,旧数据被丢弃,还没有任何痕迹,排查问题的时候你根本不知道数据是在哪个环节丢的。
- 业务数据和上行消息分开处理,容易双写不一致。采集线程先把数据写进归档文件,同时又要发消息给平台。如果写文件成功、发消息失败,或者发消息成功、文件没写进去,两边就对不上了。等到要对账的时候,平台说收到一条带某个时间戳的数据,本地文件里却没有,很难解释。
- 进程崩溃或者掉电时,内存队列没法恢复。就算你做了队列快照,快照和业务数据也不是同一个原子操作,起不来一个“数据只丢最后一个瞬间”的保证。
其实这些软肋概括起来就一句话:缺少一个“本地事务 + 持久化队列”的机制,能同时保证业务数据落盘和上行消息的可靠性。而这个机制,在微服务社区早就有了标准答案——Transactional Outbox,事务发件箱模式。
1.3 Outbox 模式到底做了什么
这里简单把 Outbox 的原理说透。它解决的是经典的双写(dual-write)问题:当你既要更新本地数据库,又要调用远程接口发送消息时,这两个操作无法保证原子性。Outbox 的做法是,把“发送消息”这件事拆成两步。
- 第一步,在同一个本地事务里,把业务数据写好,同时往一张 outbox 表里插入一条“待发送事件”记录。因为这一步是同一个数据库事务,要么都成功,要么都失败,所以业务数据落盘的那一刻,待发送事件也就落盘了。
- 第二步,一个后台线程或者独立进程不断扫描 outbox 表,把状态为“待发送”的记录取出来,真正发送到 MQTT Broker 或 HTTP 服务端;发送成功后,把记录标记为“已发送”或直接删除。
这样整个系统就摆脱了对“网络一次性成功”的依赖。网络断了,记录就躺在 outbox 表里;网络恢复,扫描线程自然会把它发出去。消息的持久化、顺序、重试都顺势变成了数据库层面可以解决的问题。
这个模式原本是订单服务写入订单后,还要发消息通知其他服务的场景。我现在把它原封不动搬进工控网关,表面上是换了个运行环境,本质上想解决的是同一个问题:有一份数据必须落库,也有一份 message 必须送达,两边不能有任何一边单独失败。
2. 网关侧的整体设计和表结构
2.1 数据链路的四个环节
在讲具体实现之前,先明确一个工控网关里典型的数据链路:采集 → 解析 → 入库 → 上行。
- 采集:网关用 Modbus TCP、OPC UA、EtherNet/IP 等协议,按周期轮询下挂的设备,拿到原始报文。
- 解析:把原始报文翻译成统一的数据结构,通常是 device_id、point_id、timestamp、value、quality 这样的点位模型。
- 入库:把点位数据写入本地存储。这一步以前很多方案是“能写就写,不能写就丢”,但在 Outbox 模式下,它变成了可靠性的入口。
- 上行:把数据打包成 MQTT 消息或者 HTTP 请求,推送到 MES、SCADA 或者云平台。
引入 Outbox 后,链路并不会变复杂,只是把“上行”环节的发送动作改为“写入 outbox 表 + 异步发送线程消费”。实时数据和历史补传也会走同一条链路:实时数据进来,写入业务表,同时落一条 outbox 记录;发送线程实时消费,网络通则马上发出去,网络断则积压,恢复后继续发。
2.2 为什么选 SQLite 作为本地存储
工控网关的本地存储可选方案其实不少:普通文件、LevelDB、SQLite、内置的时序数据库等。我的经验是,除非你已经有一个非常成熟的时序存储模块,否则SQLite 是综合成本最低的选择。
- 事务支持。这是 Outbox 模式最硬的要求。业务数据和 outbox 记录必须在同一个事务里提交,SQLite 天然支持,而且它的嵌入式特性不需要额外起一个数据库服务进程。
- 崩溃恢复能力。SQLite 的 WAL(Write-Ahead Logging)模式配合合适的 synchronous 设置,可以在掉电后自动恢复到最近一次完整事务的状态。这一点对现场网关太重要了,电源不稳、意外断电几乎是常态。
- 省资源。一个采集网关通常只有几百 MB 到几 GB 的内存,跑一个完整的数据库服务太奢侈,SQLite 一个库文件就搞定。
- 部署简单。网关上不需要安装额外的数据库软件,升级维护也省心。
如果你选型的时候已经有 InfluxDB、TDengine 之类的时序库在跑,理论上也可以在时序库里加一张 outbox 表,原理一样。但嵌入式场景我还是推荐 SQLite:轻、稳、零维护。
在初始化 SQLite 的时候,建议打开 WAL 模式,并把 synchronous 设为 FULL 或者至少 NORMAL。WAL 的好处是读不阻塞写、写不阻塞读,非常适合“采集线程写入、发送线程读取”这种并发模型。synchronous=FULL 会保证每个事务的 WAL 日志都落盘到存储介质,掉电不会丢已提交的事务;代价是单事务写入延迟稍高,但对 1 秒甚至更慢的采集周期来说,完全不是问题。
2.3 outbox 表字段设计
业务表可以根据你自己的点位模型来定,这里重点说 outbox 表。我建议的最小字段集合如下:
CREATE TABLE t_outbox ( id INTEGER PRIMARY KEY AUTOINCREMENT, topic TEXT NOT NULL, payload BLOB NOT NULL, qos INTEGER DEFAULT 1, partition_key TEXT, session_id TEXT, seq_no INTEGER, status INTEGER DEFAULT 0, retry_count INTEGER DEFAULT 0, next_retry_time INTEGER, create_time INTEGER NOT NULL, send_time INTEGER ); CREATE INDEX idx_outbox_status_id ON t_outbox(status, id);逐个解释关键字段:
- topic:目标 MQTT 主题。这样发送线程拿到一条记录就知道往哪儿发,不需要再根据业务类型去映射。
- payload:序列化好的业务数据。工控网关里不同点位、不同协议的数据结构差别很大,直接存序列化后的字节流,发送线程不需要关心业务细节,把它当成不透明的数据投递出去就行。这也是 Outbox 模式的通用做法。
- partition_key:分区键,通常用 device_id。需要并发发送时,同一个设备的数据必须路由到同一个发送线程,保证该设备的顺序性。
- session_id 和 seq_no:用来做消息的唯一标识。网关每次启动生成一个新的 session_id,seq_no 是单调递增的序列号,两者合起来就是全局唯一消息 ID。平台侧可以拿它做去重。
- status:0 待发送,1 发送中,2 已发送。发送中的状态是为了防止多个发送线程重复取同一条记录。
- retry_count 和 next_retry_time:发送失败后的重试次数和下次重试时间,可以做指数退避。
- create_time:记录写入时间,用于清理策略。
关键索引一定要加在 (status, id) 上,因为发送线程的主要查询是“取状态为待发送的最老一批记录”。如果没有这个索引,扫描会随着表变大越来越慢。
2.4 存储容量与清理策略
引入 Outbox 后,一个必须提前想清楚的问题就是:断网时间长了,outbox 表会无限增长。
来算一笔账。假设你有个网关采集 50 个点位,1 秒钟一轮,一小时就是 18 万条记录。每条记录 payload 按 200 字节算,加上 SQLite 的行开销和索引,一小时大概产生 50 到 80 MB 的数据。断网一天,积压量就在 1.2 到 2 GB 左右。普通工业级 SD 卡或固态盘可以扛住,但你不可能让它无限累积下去。
我的做法是两层清理策略配合:
- 已发送记录:发送成功后,记录其实已经没有保留价值。可以设置一个保留窗口,比如“已发送成功且发送时间早于当前时间 24 小时”的记录统一删掉。这个窗口是为了给运维留一点排查时间。
- 未发送记录:设置一个硬上限,比如 outbox 表总行数达到 100 万条,或者文件大小达到 2 GB,触发保护策略。工控场景我建议优先采用“丢弃最老的待发送记录并打点告警”,因为实时采集业务不能被历史积压堵死。如果选择阻塞新的写入,有可能因为网络长期不恢复,导致后续所有采集数据都写不进去,影响反而更大。
清理任务不要和发送线程抢资源,最好放在低峰期,按 id 范围批量删除。比如:
DELETE FROM t_outbox WHERE status = 2 AND send_time < strftime('%s', 'now') - 86400 AND id <= (SELECT id FROM t_outbox WHERE status = 2 AND send_time < strftime('%s', 'now') - 86400 ORDER BY id DESC LIMIT 1) - 10000;这种按 id 范围删除比逐条 delete 高效得多,也减少了 SQLite 的碎片。
3. 把 Outbox 落到代码实现里
3.1 同一事务写入业务表和 outbox 表
这是整个模式最核心的一段。我用 Python 伪代码演示,实际上用 Go、C++ 也是一样的思想,关键是必须用同一个数据库连接、同一个事务来完成两次插入。
def on_data_received(device_id, point_id, ts, value, quality): payload = json.dumps({ "device_id": device_id, "point_id": point_id, "ts": ts, "value": value, "quality": quality, "msg_id": f"{session_id}:{next_seq()}", }, separators=(',', ':')) topic = f"factory/{device_id}/data" try: conn.execute("BEGIN IMMEDIATE") conn.execute( "INSERT INTO t_measure(device_id, point_id, ts, value, quality) VALUES (?,?,?,?,?)", (device_id, point_id, ts, value, quality), ) conn.execute( "INSERT INTO t_outbox(topic, payload, qos, partition_key, session_id, seq_no, status, create_time) " "VALUES (?,?,?,?,?,?,0,?)", (topic, payload, 1, device_id, session_id, seq_no, int(time.time())), ) conn.execute("COMMIT") except Exception: conn.execute("ROLLBACK") # 这里要告警,业务数据落库失败意味着采集链路出了问题几个关键细节:
- BEGIN IMMEDIATE:SQLite 的默认事务是 deferred,可能在第一次写操作时才拿写锁。用 IMMEDIATE 可以在事务一开始就获得写锁,避免两个事务交错时出现 SQLITE_BUSY。
- 不要在事务里做网络请求。有些同事会把 MQTT publish 直接塞进事务里,这是典型的反向耦合。网络一卡,事务挂住,SQLite 写锁不释放,采集线程全部堵死。Outbox 模式的意义正在于把“网络发送”从“数据写入”这条关键路径上剥离开。
- msg_id 的生成:session_id 在网关启动时生成一次,seq_no 用一个全局原子计数器递增。这样即使网关重启,因为 session_id 变了,平台侧也能区分重启前后的数据。
3.2 发送线程:轮询、重试与退避
发送线程是 Outbox 模式的心脏。它要做的很简单:反复查询出 status=0 的记录,挨个发送,成功后更新状态。真正要仔细写的是重试和退避逻辑。
def outbox_sender_loop(): while True: rows = fetch_pending_outbox(batch_size=50) if not rows: time.sleep(0.5) continue for row in rows: # 先改为发送中,防止并发线程重复处理 mark_sending(row["id"]) try: mqtt_client.publish( row["topic"], row["payload"], qos=row["qos"], ) # 等待 PUBACK,确认消息已经到达 Broker # 根据客户端库不同,publish 可以是同步确认或异步回调 mark_sent(row["id"]) except Exception as e: retry_count = row["retry_count"] + 1 backoff = min(2 ** retry_count, 60) # 指数退避,封顶 60 秒 update_retry(row["id"], retry_count, time.time() + backoff)重试退避的计算要注意:如果每失败一次就立刻重试,断网期间会形成一个疯狂重试的循环,白白消耗 CPU 和电量。指数退避是基本操作:第一次失败等 2 秒,第二次 4 秒,第三次 8 秒……封顶到 60 秒或更长,网络恢复后自然会把积压的数据发出去。
如果你用 MQTT 客户端,要注意 QoS 1 的确认机制。以 Paho 为例,publish()方法本身只负责把消息交给客户端库的发送队列,真正确认消息到达 Broker 是在回调里。所以正确做法是:
# 使用回调等待 PUBACK,而不是调用 publish() 后立刻置为已发送 def on_publish(client, userdata, mid): mark_sent_by_mid(mid)如果忽略了 PUBACK 回调,直接把状态改成已发送,会有一批消息其实没到 Broker 就被标记成功了,严格说就失去了 at least once 的可靠性。
3.3 顺序与并发怎么平衡
很多工控场景对数据顺序敏感。同一个设备的温度、压力、流量,如果到达平台的顺序颠倒了,趋势曲线和报警判断都会出问题。Outbox 模式天然支持顺序性:只要发送线程是单线程,并且按 id 升序发送,记录就一定按写入顺序发出。
但单线程发送吞吐有限,断网恢复后如果需要快速补传几十万条数据,单线程可能不够快。这时候可以用 partition_key 做分区并发:
- 先把待发送记录按 partition_key 分组;
- 同一个 partition_key 的记录,始终由同一个发送线程处理;
- 不同分区之间的顺序互相不影响。
简化理解就是:把“工厂/3号线/注塑机”的数据都扔给线程 A,“工厂/3号线/压铸机”的数据扔给线程 B,每个线程内部严格按 id 升序发送。全局顺序不保证,但单设备顺序绝对不乱。
实际做的时候要注意,标记 status=1(发送中)时,记录被哪个线程领走,就必须由哪个线程负责完成,否则会出现一条记录被两个线程同时拿到的竞态。最简单的方式是为每个发送线程指定一个固定的 worker_id,在领取记录时按 worker_id 分组:
UPDATE t_outbox SET status = 1, send_time = ? WHERE id IN ( SELECT id FROM t_outbox WHERE status = 0 AND partition_key IN (SELECT DISTINCT partition_key FROM t_outbox WHERE status = 0) AND id % worker_num = ? ORDER BY id LIMIT 50 ) RETURNING *;这个写法取决于你用的数据库,SQLite 对 RETURNING 的支持要看版本。更通用的做法是在应用层先按 partition_key 路由,再单线程处理每个分区的记录。
3.4 幂等设计:为什么必须做
Outbox 模式给到平台的语义是“至少一次”(at least once),不是“恰好一次”。网络超时、发送线程崩溃重试,都有可能导致平台重复收到同一条数据。所以在网关侧和平台侧的协议设计里,幂等是必须的。
网关侧要做的很简单:在 payload 里带上msg_id,这个 ID 全局唯一。平台侧收到数据后,根据 msg_id 做去重,或者按 (device_id, point_id, ts) 做 upsert,重复插入时直接忽略或覆盖。
{ "device_id": "injection_03", "point_id": "temp_mold", "ts": 1739000000, "value": 45.6, "quality": 0, "msg_id": "a3f1c2e0-8f90-4b1e-9c2d-000000001234" }我见过有的项目没有做幂等,断网恢复后,平台上一小时的数据量翻了几倍,聚合报表全部算错。排查了半天才发现是网关重发导致的重复数据。这个坑一定要提前堵上。
3.5 已发送记录的清理任务
清理逻辑可以放在发送线程的循环里,每隔一段时间触发一次,或者在低峰期单独跑一个定时任务。核心就是限制表体积:
def outbox_cleanup_loop(): while True: # 每小时清理一次 time.sleep(3600) execute(""" DELETE FROM t_outbox WHERE status = 2 AND send_time < strftime('%s', 'now') - 86400 AND id % 10 = 0 """) # 如果表仍然过大,触发保护策略 size = query_value("SELECT COUNT(*) FROM t_outbox") if size > MAX_OUTBOX_ROWS: notify_alert("outbox table overflow, dropping oldest pending records") execute(""" DELETE FROM t_outbox WHERE id IN ( SELECT id FROM t_outbox WHERE status = 0 ORDER BY id LIMIT 10000 ) """)实际项目里清理策略会因为业务不同而调整,但原则是一致的:不能让 outbox 表无限膨胀,也不能因为清数据把未发送的历史数据误删。建议删数据前先统计、先告警。
4. 实测结果与避坑指南
4.1 模拟网络故障的实测表现
我在一个 X86 工控机上跑过一轮测试,环境是这样的:一台 Modbus TCP 模拟器下挂了 50 个点位,网关 1 秒采集一轮,上行 MQTT 到本机搭的一个 Broker,网关和 Broker 之间用一个可以手动断网的交换机隔开。
- 正常状态:采集产生的 outbox 记录基本发一条清一条,表里常年保持在个位数的行数。
- 断网 30 分钟:outbox 表积累了 9 万条左右的待发送记录,占用磁盘约 60 MB。网关内存占用几乎没有变化,采集线程照常工作,业务表正常写入。
- 恢复网络:发送线程按指数退避逻辑探测,网络恢复后 1 秒内开始补传。补传速度可以达到每秒 2000 到 5000 条,大约 30 分钟积压的数据在几分钟内全部追平。
- 断网期间重启网关:重启后 outbox 表里的记录依然在,发送线程启动后继续补传,一条都没有丢。
测试下来最直观的感受是:断网时长不再是问题,数据量才是。只要磁盘容量规划好,断一小时和断十天,行为模式完全一致——积压、恢复、追平。
4.2 实测中容易踩的坑
下面这些坑全是我在实际运行里踩过的,列出来,帮你省掉排查的时间:
- SQLite 锁冲突。采集线程写入频繁,发送线程又在读取和更新,如果没开 WAL,会出现 SQLITE_BUSY。开了 WAL 还不够,写入端要避免长事务,发送端更新状态不要和读取放在同一个事务里。
- MQTT 客户端库自带的上行队列和 outbox 叠加导致内存暴涨。很多 MQTT 客户端在断线时会内部缓存消息,如果你已经在 outbox 表里积压,客户端又在内存里攒一份,网络恢复瞬间会同时出现两份数据释放,内存可能直接被打满。建议把客户端库的本地队列关掉,或者设到很小的值,让 outbox 成为唯一的缓存层。
- 清理任务把发送中的记录删掉。status=1 的记录正在发送,如果清理任务按“早于某个时间”无脑删除,可能导致某条数据永远丢失。清理条件一定要带上 status=2。
- 时间戳不一致。断网期间如果 NTP 失效,网关本地时钟会漂移,恢复后发送的数据时间戳和真实时间可能差几十秒甚至几分钟。平台按时间排序时,数据会显得错乱。对策是网关启动时从平台校准一次时间,平时尽量依赖 RTC 晶振,并且把“采集时间”和“网关发送时间”两个字段在 payload 里分开。
- 每次事务只插一条数据,写入放大严重。采集周期是 1 秒、点位 50,可能还好;如果采集频率高,比如 100ms 一轮,每秒就有 1000 条记录。每次写一个事务,SQLite 的 fsync 会把磁盘 IO 打满。正确做法是攒一批,比如每 500ms 或每 500 条批量提交一个事务。
- 发送线程轮询间隔太短,空转耗电。网络正常时,outbox 表经常是空的,发送线程还在 10ms 轮询一次,CPU 白白浪费。可以做成动态间隔:连续几次查询不到记录,就把 sleep 从 0.1 秒逐步增加到 2 秒;一旦查到记录,马上改为密集发送。
4.3 常见问题速查表
| 症状 | 可能原因 | 解决方案 |
|---|---|---|
| 断网恢复后数据少了一部分 | 清理策略误删了 status=1 的记录;MQTT QoS 设置成了 0 | 清理条件加上 status=2;上行 QoS 改成 1 以上 |
| 平台收到大量重复数据 | 没有做幂等;发送线程重复领取记录 | payload 加 msg_id;平台按 msg_id 或 (device_id, point_id, ts) 去重 |
| 断网期间内存持续上涨 | MQTT 客户端库内部缓存叠加 outbox | 关闭或限制客户端库的本地队列 |
| 采集线程偶发延迟很大的尖刺 | SQLite 写锁被发送事务阻塞;每次提交次数太频繁 | 开启 WAL;批量提交;发送线程避免长事务 |
| 补传速度慢,平台端数据迟到太久 | 单线程发送吞吐有限;退避参数太长 | 按 partition_key 分区并发;降低重试退避的封顶时间 |
| outbox 表文件越来越大 | 没有清理任务,或者清理任务没生效 | 设置保留窗口 + 定时批量删除 |
5. 设计取舍与后续扩展
5.1 有损和无损,取决于业务怎么定位
Outbox 不是银弹,它保证的是“只要 outbox 表里还有记录,断网数据就还能补传”。但如果你因为存储空间限制,主动丢弃了最老的待发送记录,那本质上还是有损的。所以方案设计的第一件事,是和业务方对齐一个问题:断网期间的监控数据到底能不能丢?
- 对大多数过程量监控来说,几分钟的数据空洞会带来报警延迟和趋势断层,最好不丢。
- 对某些极端场景,比如超高速采集、振动分析,数据量极大,网关本地存储根本放不下,那就要在设计上默认“只保留最近一段时间的窗口,超出窗口直接丢弃并打点”。
我的建议是:网关侧把“丢弃行为”也变成一条告警数据。比如 outbox 表溢出时,主动往平台发一条t_outbox_overflow的告警消息,内容包括丢了多少条、最早丢弃的时间范围。这样就算数据有损,运维也能知道损在了哪儿,而不是靠猜。
5.2 批量融合发送,减少连接开销
断网补传时,一条一条地发 MQTT 消息,在数据量大时效率不高。可以做一个批量融合逻辑:把同一个 topic 下连续的多条 outbox 记录合并成一条批量消息,一次发布。比如:
{ "type": "batch", "device_id": "injection_03", "records": [ {"ts": 1739000000, "value": 45.6, "quality": 0}, {"ts": 1739000001, "value": 45.7, "quality": 0} ] }平台侧按 batch type 解析,逐条 upsert 进去,仍然以单条 msg_id 去重。这个优化在正常运行时可以不启用,只在发现 outbox 积压量超过阈值时动态开启,降低 broker 压力和网络包数量。
5.3 反向指令场景也用得上 Outbox,但要谨慎
Outbox 不仅可以用于数据上行,网关做控制指令下发时,同样会遇到“命令必须下达到设备、同时要记录下发结果”的一致性需求。把控制指令写入 outbox、后台线程异步下发,可以在网关重启后避免指令状态不一致的问题。
但工控下行和数据上行不一样,控制指令有实时性要求,而且涉及安全。不能简单地把指令排在历史数据的后面慢慢发。我的经验是:下行指令要单独一个 outbox 表,或者至少用 topic 前缀区分,并且发送线程要有更高优先级和独立的连接。同时,指令下发必须配套超时确认和急停机制,不能无限重试,否则现场设备可能被重复下发的启动指令搞出事故。
最后再说一点实操感受
用 Outbox 模式做网关的上行通道,我最大的体会是:它把“网络不可靠”这个心腹大患,变成了一张可以度量、可以管理的表。以前排查掉数据,要在日志里翻半天;现在只需要看 t_outbox 表里的 status 分布和积压量,就知道问题是出在网络、平台还是本地存储。排障从“猜”变成了“看数据”。
如果你准备在新网关项目里用这个模式,我建议一开始就把幂等、清理、告警这三件事做进去,而不是先把通道跑通再补。很多项目就是在“能发就行”的阶段上线,一旦遇到真正长时间断网,再回来补这些能力,成本比一开始写高好几倍。
另外一个实用的小技巧是把 outbox 表的积压量作为一个指标上报到你自己的运维监控里。比如“当前待发送记录数”超过某个阈值就告警,网络恢复后看到积压曲线降下来,比什么消息队列的监控都直观。工控网关这个场景,稳定压倒一切,设计上多花一天,现场少熬三个夜。