Files
msgexchange-v2/docs/message-lifecycle.md
T

16 KiB
Raw Blame History

上游消息生命周期设计

本文范围与读者

本文定义 msgexchange-v2(下称"本系统")对共享 MySQL 信箱中报文的完整生命周期:从上游写入、本系统处理、处理标记回信箱,到信箱数据最终清除。各阶段的输入输出、幂等方式、故障恢复方法,以及本系统与信箱库管理方之间的分工,都在本文界定。

  • 模块协作、状态机与配置参数见 design.md;系统边界与"不建表、不改结构"的红线见 architecture.md;验收口径见 user-stories.md
  • 读者包括本系统开发与运维人员、以及负责共享 MySQL 的库方接口人;§6、§7 中标注 Q 编号的条款即需要与库方书面确认的开放项。

角色与术语约定(全文沿用):

  • 上游:向信箱写入报文的源头系统(CIIMS、AODB 等)。
  • 信箱:共享 MySQL 中的入站信箱表 CMINMSGS;出站方向为 COUTMSGS(§7)。
  • 库方:共享 MySQL 的管理方;表结构变更与数据清除只能由库方执行或书面授权。
  • 处理标记:信箱行上表示"本系统已处理"的约定字段;逻辑名 DATE_PROCESSED / STATUS,实际列名以库方契约为准(legacy 为 CMINMSGS_DATE_PROCESSED 等真实列)。
  • 自有 PG:本系统唯一的业务数据库 PostgreSQL;信箱与自有 PG 之间不存在跨库事务。

1. 五个独立事实

消息的认定分为五个环节,每个环节都有独立的证据,互不替代:上游写入了不等于本系统接手了,本系统处理完了也不等于信箱标记已写、下游已收到。

环节 认定依据 执行方
落信:报文进入信箱 CMINMSGS 中存在该行 上游
入队:本系统开始处理 自有 PG 建立 PROC_STATE 记录 InboxPoller
处理完成:业务有结论 PROC_STATE 到达终态(SUCCEEDED / SKIPPED / DEAD 主泵
已回填:信箱写入处理标记 信箱行持有处理标记 回填与补偿通道
下游确认:下游已接收 MSG_EVENTSENT deliveryDispatcher / flushSchd

分环节的原因是存储边界:信箱与自有 PG 之间的写入不共享事务。前三个环节以自有 PG 为准;"已回填""下游确认"两个环节各自独立重试,都可能单独失败。

2. 状态总纲

信箱行    未标记 ────────────────→ 已处理标记 ────────→ 已归档 ────────→ 已清除
                                     ↑ 回填/超期补写     §6(库方执行)
PROC_STATE  PENDING / FAILED ──→ SUCCEEDED(成功)
                                ├→ SKIPPED(业务重复)
                                └→ DEADMALFORMED / PROTOCOL / EXHAUSTED,人工处置)
MSG_EVENT   PENDING ──→ SENT
        └──→ PENDING(退避重试)──→ DEADDLQ)
  • PENDINGFAILED 是处理中的状态,FAILED 继续退避重试并占用 FIFO 队头;SUCCEEDED / SKIPPED / DEAD 是终态,到达后队列方可推进。
  • DEADFAILED 不是不可逆:经人工批准,指定错误类别的记录可重新置回 PENDING 处理(放行范围见 §9 / design.md §6.1)。重放处理的是当前状态,不恢复历史处理顺序。
  • 航班变更、待发事件、处理终态与回填待预登记四类写入位于同一个 PG 事务,一起提交或一起回滚;信箱处理标记写在该事务之后,依赖 BACKFILL_TODO 补偿记录最终写入。

3. 处理阶段

阶段 输入 输出 幂等依据 失败处理
收报 DATE_PROCESSED IS NULL 的信箱行 PROC_STATE(PENDING) MSG_ID 主键 下轮扫描补入队
取队头 最小未完成 MSG_ID 本批处理消息 表状态即队列 无队头则等待
解码与身份 信箱原文 DecodedMessageIDENTITY_KEY 身份唯一约束 报文非法→DEAD;处理能力不足→FAILED 退避
事务处理 当前完整状态 + 报文载荷 航班变更、终态、事件、回填待办 消息 ID + 业务身份 事务整体回滚重试
提交后回填 终态消息 信箱处理标记 标记单调(§11 记入 BACKFILL_TODO 重试
投递 MSG_EVENT 下游确认 EVENT_ID 退避至 DEAD;至少一次
补偿作业 回填待办 / 到期批次 追平标记 待办主键 指数退避,独立线程执行

处理器只在持有 PIPELINE_LOCK 的事务内写入自有 PG 状态,不操作信箱与 Kafka;全部外部副作用发生在事务提交之后。

4. 故障恢复

恢复的唯一依据是各存储中已持久化的记录,不依赖任何进程内存中的状态:

中断位置 重启后的判定 恢复动作
已落信、未入队 信箱无标记且 PG 无记录 重扫补建入队记录
事务执行中 PG 无该消息终态 事务整体回滚,按 PENDING 重新处理
事务已提交、标记未写 BACKFILL_TODO 存在待办 仅补写标记;业务处理结果保持不变
回填待办登记失败(与业务同事务) 事务未提交 同"事务执行中",不构成独立窗口
标记写入中途 标记仍为空 重新写入;重复写入同一值无副作用
下游已接收、SENT 未置 事件仍 PENDING 允许重发,下游按事件身份去重

两个缺口在闭环前不能宣称恢复完整:

  1. 事务提交后、补偿开始前仍存在崩溃窗口;
  2. 补偿重试尚未证明"最终一定完成"。

两项登记于 design.md §10(事务与外部副作用)与 user-stories.md US-09。

5. 消费水位、超期补写与历史积压

5.1 消费水位 W

水位 W 定义为信箱 ID 的连续上界:从最小 ID 到 W 的区间已全部读入自有 PG,无空洞;遇到空洞即停止推进。W 只随新 ID 的成功入队推进,不依赖处理完成或标记回写。

水位的有效性依赖两条需要库方确认的承诺(合计为 Q2):

  1. ID 单调:信箱 ID 按提交顺序分配,晚提交的较小 ID 不参与。承诺缺失时,W 只能作为快路径的提示,不能证明该区间收齐。
  2. 最大提交时延:上游提交到 ID 可见的最长时间。补偿扫描窗口(NOW 最大提交时延)的宽度由此确定。

日常执行方式:快路径从 ID > W 起按升序有限批次扫描;另按周期对窗口内可能迟到或空洞的行做补偿扫描。需要区分:水位表示"读取进度",与"已处理标记"是两个事实,不能互相替代。

5.2 超期标记补写

两类信箱行无法通过正常回填获得处理标记:报文残缺且缺少元数据的死信,以及回填通道长期失败的行。处理规则为:

已达 PROC_STATE 终态、且接收时间超期(DATE_RECEIVED < NOW R)仍无标记的信箱行,由回填通道补写一个库方认可的"已处理"类标记。

  • 期限 R 必须不小于(人工重放期限 + 人工处置期限)之和。提前补写会使仍可重放的消息失去原件资格。
  • 中间态(PENDING / FAILED)不适用本规则:处理未完成时不打标,也不允许被清除。
  • 补写值仅限于库方认可的 legacy 值集(Q7);死信、业务重复等内部原因记录在 PROC_STATE 与审计日志,不在信箱新增枚举。
  • 补写只针对空标记;已有值不回撤、不覆盖,重复执行无副作用。

该规则同时保证三件事成立:收报表不被永不回填的记录占满(Q2 中"积压挡批"场景)、§6 的清除边界可以达到、重放期限与信箱保留期具备可核验的下限关系(§9)。

5.3 历史积压消息的处理方案

"历史积压"指信箱中成规模的未处理存量:上线前遗留、停机期间累积、或消量未完成的批次。处理方案分三步,并约束四条行为红线。

第一步:摸底。 处理动作开始前,先确定积压的范围与构成:ID 区间、条数、时间跨度、报文类型分布(SCHD / FLOP / FDEL / ADFT 等),并与库方确认其中哪些仍需业务处理、哪些按约定跳过(跳过基调只在该批确认时成立,值集见 Q7)。

第二步:先入队,后处理。 两阶段执行,禁止边收边处理一辆长队:

  1. 消化阶段只做入队:InboxPoller 按升序有限批次将历史行全部建为 PROC_STATE(PENDING),水位随之推到积压末端。此阶段只写自有 PG,不触碰信箱标记。
  2. 入队完成后交给主泵按最小未完成 ID 顺序消化。顺序与阈限与日常完全相同:不加速、不分流、不走旁路。

这样做的理由:入队是便宜操作(可大批量),处理是昂贵操作(解码 + 事务),把两者分开后,入库阶段的中断恢复只涉及重扫(§4 第一行),不会把半处理的事务复杂化;同时水位可以尽早到达完整上界,快路径与补扫窗口立即生效。

第三步:等待过程中的四个边界。

  • 不插队:主泵按 FIFO 消化,积压期间到达的实时消息排在积压之后。本设计不允许并行队头,也不允许实时通道跳过积压。
  • 不失控:队头滞留上限(head-deadline,当前 10 分钟)对积压同样生效;队头长期失败按既有规则转 DEAD(EXHAUSTED) 并告警,不会因积压而延长容忍。
  • 报文类型的处理方式不变:旧的全量日计划(SCHD-DNLD)按最近一份覆盖即可正确收敛,但仍逐条执行;动态增量(FLOP / FDEL / ADFT)持有时序语义,必须逐条。
  • 可放弃但必须留痕:摸底批中经库方确认"不再处理"的行,处置方式为——PROC_STATESKIPPED 并记录跳过原因,到达终态后走 §5.2 通道补写标记;本方案不支持任何"整段 DELETE"的快速通道。

验收口径:积压消化期间持续输出三项指标——剩余积压条数、最老未处理信龄、预计消化时长;期间不允许出现 FIFO 越序、身份去重失效或头行滞留超时未告警。红线依据:architecture.md §5(消息严格 FIFO)、单写者约束。

6. 信箱数据清除

表结构变更与物理清除由库方执行或书面授权执行;本系统对共享 MySQL 不建表、不改结构(红线见 architecture.md §6)。本系统在清除事务中的义务只有一项:为处理完成的行及时写入处理标记,使可清除范围存在明确边界。

具体方案由库方选择(Q9)。以下两种方案均基于 MySQL 自身能力:支持 RANGE 分区与 TRUNCATE / DROP PARTITION,不支持 EXCHANGE PARTITION

方案 A:按日分区(首选),适用于库方可以为表增加分区的场合。

  1. CMINMSGSDATE_RECEIVED 建立日粒度 RANGE 分区;
  2. 某分区到达保留期时,确认该分区全部行已持有处理标记(§5.2 保证该条件在有限时间内满足);
  3. 将该分区复制入历史表:INSERT INTO CMINMSGS_HST SELECTNOT EXISTS 判重);
  4. TRUNCATE / DROP PARTITION 执行清除:DDL 级操作,无行锁竞争,页外大字段(CMINMSGS_CLOB_MSG 列)整块释放,不产生碎片与 binlog 压力。

前提:库方确认现场 MySQL 版本支持分区 DDL,并已授权执行。

方案 B:整表轮换,适用于库方拒绝增加分区的场合。

  1. CREATE TABLE CMINMSGS_NEW LIKE CMINMSGS,并将 AUTO_INCREMENT 种子设为 max(ID) + 1
  2. 将保留窗内行(DATE_RECEIVED ≥ NOW R_keep)复制至新表;
  3. 以单语句原子 RENAME TABLE 完成换名;
  4. 对账:换名与复制之间新写入及新标记的行,从旧表幂等补回;
  5. 旧表中早于保留窗的部分追加写入 CMINMSGS_HSTNOT EXISTS 判重);
  6. DROP TABLE 旧表:秒级完成,碎片与大字段一并释放。

两种中断均安全:换名前失败则废弃新表重新执行;换名后失败则旧表仍完整,对账与归档语句均可重跑。

方案共同前提:

  • 执行清除时,边界内不存在未打标记的行;未达终态的行顺延至处理完成后清除(§5.2 仅对终态行补标);
  • 时间比较与换算统一采用机场时区 Asia/Shanghai 及明确的类型转换(口径同 user-stories.md §6);
  • 清除保留期 R_keep 不小于 max(回填重试上限、重放期限、审计期限);且 §5.2 的 R 不大于 R_keep,否则尚在重放窗口内的消息会先于重放被清除;
  • 信箱 ID 全程不断链:方案 A 天然满足;方案 B 依赖 AUTO_INCREMENT 种子,种子缺失时新 ID 与旧记录主键冲突,水位随之失效。

7. 出站信箱(COUTMSGS

本系统一侧的规则(实现中:出站适配与请求协调尚未交付,机制描述见 design.md §4.2):

出站前在 REQ_TRACK 登记 REGISTERED;写入 COUTMSGS 并确认落信后置 SENT 并关联出站记录 ID。"落信"指写入信箱成功,不等于下游已读取或已发送,交付承诺仅到落信为止(US-08)。

本系统写入的行由谁消费、按什么顺序消费,当前未知。以下事项在 Q10 关闭前按未知处理,不得作为已具备能力描述:

  • 消费方与消费顺序;
  • COUTMSGS_ACK_DATE_RECV / ACK_RESEND_TIMES / DATE_SENT / ERROR 各列的语义与写入责任;
  • 出站行的清除责任与保留期;
  • 落信成功但本地未置 SENT 时的重复写入风险及下游去重契约。

参考事实:legacy 的 SIS 接口规范记载各系统持有各自的 COUTMSGS 副本,由 JDBC Adapter 读取后发送至 CIIMS。该机制是否在现场沿用,属于 Q10 的确认范围,不作为现状依据。

8. 自有记录归档

  • PROC_STATE:终态记录归档至 PROC_STATE_HST(尚未建表)。未达终态的记录不归档;归档不得使"同一业务身份只能绑定一条有效处理记录"的去重能力失效(US-11)。
  • SCHD_SNAP_LOG:保留 90 天,清理窗口与判据见 design.md §6.2。
  • 航班当前态的清理规则唯一归属 flight-state.md §6;其判据不依赖信箱原文。信箱原文的可用性仅影响重放能力(§9)。

9. 重放与原文可用性

  • 可重放的错误类别为 CODEC_ERROR / UNSUPPORTED / INFRA / EXHAUSTED 四类(以 design.md §6.1 为准);重放处理当前状态,不恢复历史顺序。
  • 重放的前提是信箱原文仍可读取:R_keep 必须覆盖重放期限(§6),属于硬约束。
  • 处理时无法读取原文,一律按 MALFORMED → DEAD 处理。若原文缺失源于库方违反保留契约提前清除,按契约违例走运维追责通道;该消息的死信处置本身不变。

10. 开放问题索引

编号 待确认事项 本文相关章节
Q2 信箱 ID 单调承诺;最大提交时延;空洞与迟到处理 §5.1
Q7 处理标记值集与写权限;原文保留期;处理时间语义 §5.2、§6
Q9(新增) 清除执行方与 DDL 授权;方案 A / B 选型 §6
Q10(新增) 出站消费方、ACK 列语义、出站清理与去重契约 §7
Q11(新增) 上游 SEQN 重置规则与业务身份的日期边界 design.md §2.2

Q9Q11 的完整登记见 user-stories.md §6 的 Q 表。

11. 不变量

  • 五个事实互不替代:入队不引用信箱标记,回填不引用投递,投递不引用回填。
  • 处理终态不可逆:已提交的 SUCCEEDED 不因回填或投递失败回改。
  • 处理标记单调:任何路径只将空标写为已处理,不回撤、不覆盖。
  • 水位不越过未入队的 ID;遇空洞即停。
  • 执行清除前,边界内全部行已持有处理标记;物理删除仅发生在归档成功之后(追加写入与分区留档均构成归档成功)。
  • 信箱 ID 全局单调、不断链,在任何清除方案下成立。
  • 对外投递按至少一次设计;端到端恰好一次不在交付范围内。