LabHub
学习 学习路径 课程

系统间对接 (EAI)

实现文件队列消费者与幂等重处理

在 LabHub 中继续学习

目标

使用基于目录的队列构建消费者,隔离解析失败消息,通过数据库约束保证幂等性,并实现积压监控和重新处理。

为什么重要

在异步集成中,重复不是例外,而是默认情况。“exactly-once”并非传输层的魔法,而是“至少传输一次 + 接收端去重”的结果。因此,消费者必须始终假设会收到重复消息。另一方面,如果把解析失败的消息放回重试队列,该消息会在队首不断失败,阻塞后面的正常消息(毒消息)。队列一旦阻塞,业务很快就会停摆。亲手实现这两项机制后,就能建立对异步设计的直观认识。

步骤

  1. 创建 /root/q/inbox, /root/q/processing, /root/q/done, /root/q/error,并把 /opt/lab/fixtures/eai/queue/inbox/ 中的所有 .json 复制到 /root/q/inbox。inbox 中必须有 20 个文件。
  2. 分析消息结构,创建 /root/q/schema.csv。第一行为 field,type,role。写入正常消息中的所有字段;在 role 列中,为充当幂等键的字段写入 idempotency-key,为用于判断顺序的字段写入 sequence
  3. 创建并运行 /root/q/consume.py。其行为如下。
    • inbox 中的文件逐个移动到 processing 后再读取。
    • JSON 解析失败,或缺少必填字段(msg_id, order_no, seq, amount)时,移动到 error
    • 正常时,加载到 sqlite 数据库 /root/q/ledger.dbprocessed 表,再移动到 done
    • processed 表中的 msg_id 必须具有PRIMARY KEY 或 UNIQUE约束,并包含 order_no, seq, amount, processed_at 列。
    • 运行结束时,processing 必须为空。
  4. 运行结果必须如下。
    • done 18 个,error 2 个
    • processed 表行数为去重后的 15 条
    • 将因重复而忽略的 3 个 msg_id 按升序保存到 /root/q/dup.txt
  5. 把移入 error 的消息原因整理到 /root/q/error.csv。第一行为 file,reasonreasonparsemissing-field
  6. 编写 /root/q/order-check.sql。该查询在 processed 表中查找同一个 order_noseq 重复或缺失的情况。把结果(没有问题时为 0 行)保存到 /root/q/order-result.txt。没有问题时,文件第一行必须为 OK
  7. 创建 /root/q/lag.sh。接收两个参数(큐디렉터리 임계치),输出一行 inbox=<n> processing=<n> error=<n>;当 inbox 数量超过阈值时,以非 0 退出码结束。
  8. 创建 /root/q/replay.sh。接收一个参数(文件名),把 error 中的该文件移回 inbox。如果该文件不在 error 中,则完全不改变队列状态,并以非 0 退出码结束。

参考

配置队列目录

创建 /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。其行为如下。

解析失败的消息即使重试也会以同样方式失败。必须立即发送到 error,避免阻塞后续正常消息。处理历史应与业务处理在同一事务中记录。

确认幂等性

运行结果必须如下。

用数据库约束阻止重复,比在代码中筛除重复更安全。即使应用程序有缺陷,约束也不会被绕过。

分类失败原因

把移入 error 的消息原因整理到 /root/q/error.csv。第一行为 file,reasonreasonparsemissing-field

请把原因分为“格式错误”和“业务错误”。前者需要修改源数据,后者补充数据后可以重新投入。

验证顺序保证

编写 /root/q/order-check.sql。该查询在 processed 表中查找同一个 order_noseq 重复或缺失的情况。把结果(没有问题时为 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 时,必须不执行任何操作并失败。