实现文件队列消费者与幂等重处理
目标
使用基于目录的队列构建消费者,隔离解析失败消息,通过数据库约束保证幂等性,并实现积压监控和重新处理。
为什么重要
在异步集成中,重复不是例外,而是默认情况。“exactly-once”并非传输层的魔法,而是“至少传输一次 + 接收端去重”的结果。因此,消费者必须始终假设会收到重复消息。另一方面,如果把解析失败的消息放回重试队列,该消息会在队首不断失败,阻塞后面的正常消息(毒消息)。队列一旦阻塞,业务很快就会停摆。亲手实现这两项机制后,就能建立对异步设计的直观认识。
步骤
- 创建
/root/q/inbox,/root/q/processing,/root/q/done,/root/q/error,并把/opt/lab/fixtures/eai/queue/inbox/中的所有.json复制到/root/q/inbox。inbox 中必须有 20 个文件。 - 分析消息结构,创建
/root/q/schema.csv。第一行为field,type,role。写入正常消息中的所有字段;在role列中,为充当幂等键的字段写入idempotency-key,为用于判断顺序的字段写入sequence。 - 创建并运行
/root/q/consume.py。其行为如下。- 把
inbox中的文件逐个移动到processing后再读取。 - JSON 解析失败,或缺少必填字段(
msg_id,order_no,seq,amount)时,移动到error。 - 正常时,加载到 sqlite 数据库
/root/q/ledger.db的processed表,再移动到done。 processed表中的msg_id必须具有PRIMARY KEY 或 UNIQUE约束,并包含order_no,seq,amount,processed_at列。- 运行结束时,
processing必须为空。
- 把
- 运行结果必须如下。
done18 个,error2 个processed表行数为去重后的 15 条- 将因重复而忽略的 3 个
msg_id按升序保存到/root/q/dup.txt
- 把移入
error的消息原因整理到/root/q/error.csv。第一行为file,reason。reason为parse或missing-field。 - 编写
/root/q/order-check.sql。该查询在processed表中查找同一个order_no内seq重复或缺失的情况。把结果(没有问题时为 0 行)保存到/root/q/order-result.txt。没有问题时,文件第一行必须为OK。 - 创建
/root/q/lag.sh。接收两个参数(큐디렉터리 임계치),输出一行inbox=<n> processing=<n> error=<n>;当inbox数量超过阈值时,以非 0 退出码结束。 - 创建
/root/q/replay.sh。接收一个参数(文件名),把error中的该文件移回inbox。如果该文件不在error中,则完全不改变队列状态,并以非 0 退出码结束。
参考
- 查看 sqlite 模式:
sqlite3 /root/q/ledger.db '.schema processed' - 忽略重复插入:
INSERT OR IGNORE或ON CONFLICT DO NOTHING - 原子移动:同一文件系统内使用
os.rename/mv - 常见错误 1:把解析失败消息放回 inbox,造成无限循环。
- 常见错误 2:只用应用程序条件判断防止重复。 并发执行时会失守。数据库约束是最后一道防线。
- 常见错误 3:不使用
processing,直接从inbox读取。 启动两个消费者时,两者都会处理同一消息。
配置队列目录
创建 /root/q/inbox, /root/q/processing, /root/q/done, /root/q/error,并把 /opt/lab/fixtures/eai/queue/inbox/ 中的所有 .json 复制到 /root/q/inbox。inbox 中必须有 20 个文件。
创建 inbox/processing/done/error 四个阶段。请记住,只有放在同一文件系统中,移动操作才是原子的。
分析消息结构
分析消息结构,创建 /root/q/schema.csv。第一行为 field,type,role。写入正常消息中的所有字段;在 role 列中,为充当幂等键的字段写入 idempotency-key,为用于判断顺序的字段写入 sequence。
打开消息样本并整理字段。请特别留意哪个字段可以充当幂等键。
实现并运行消费者
创建并运行 /root/q/consume.py。其行为如下。
- 把
inbox中的文件逐个移动到processing后再读取。 - JSON 解析失败,或缺少必填字段(
msg_id,order_no,seq,amount)时,移动到error。 - 正常时,加载到 sqlite 数据库
/root/q/ledger.db的processed表,再移动到done。 processed表中的msg_id必须具有PRIMARY KEY 或 UNIQUE约束,并包含order_no,seq,amount,processed_at列。- 运行结束时,
processing必须为空。
解析失败的消息即使重试也会以同样方式失败。必须立即发送到 error,避免阻塞后续正常消息。处理历史应与业务处理在同一事务中记录。
确认幂等性
运行结果必须如下。
done18 个,error2 个processed表行数为去重后的 15 条- 将因重复而忽略的 3 个
msg_id按升序保存到/root/q/dup.txt
用数据库约束阻止重复,比在代码中筛除重复更安全。即使应用程序有缺陷,约束也不会被绕过。
分类失败原因
把移入 error 的消息原因整理到 /root/q/error.csv。第一行为 file,reason。reason 为 parse 或 missing-field。
请把原因分为“格式错误”和“业务错误”。前者需要修改源数据,后者补充数据后可以重新投入。
验证顺序保证
编写 /root/q/order-check.sql。该查询在 processed 表中查找同一个 order_no 内 seq 重复或缺失的情况。把结果(没有问题时为 0 行)保存到 /root/q/order-result.txt。没有问题时,文件第一行必须为 OK。
使用 SQL 确认相同键的消息是否按序号处理。应比较保存的序号,而不是处理时间。
积压监控脚本
创建 /root/q/lag.sh。接收两个参数(큐디렉터리 임계치),输出一行 inbox=<n> processing=<n> error=<n>;当 inbox 数量超过阈值时,以非 0 退出码结束。
通过参数接收队列目录,脚本才能复用。超过阈值时必须用退出码通知,才能接入 cron 或监控工具。
重新处理脚本
创建 /root/q/replay.sh。接收一个参数(文件名),把 error 中的该文件移回 inbox。如果该文件不在 error 中,则完全不改变队列状态,并以非 0 退出码结束。
重新处理并不意味着随意把任何内容放回队列。收到不存在的消息 ID 时,必须不执行任何操作并失败。