GitHub user WangzJi created a discussion: TC 事务清理优化:二阶段处理流水线化
> 关联 > issue:[#6615](https://github.com/apache/incubator-seata/issues/6615)(tracking,OPEN)、[#7362](https://github.com/apache/incubator-seata/discussions/7362)、[#7334](https://github.com/apache/incubator-seata/issues/7334) --- ## 1. 背景 ### 1.1 一行记录的生命周期 TC 用 `global_table` 记录每个全局事务的状态,**一行 = 一个全局事务**: ```mermaid flowchart LR TM1["TM 开启<br/>全局事务"] --> B["Begin(1)<br/><b>INSERT 一行</b>"] B --> R["分支注册<br/>branch_table<br/>+ lock_table"] R --> TM2["TM 提交"] TM2 --> AC["AsyncCommitting(8)<br/>AT 模式二阶段<br/>可异步执行"] AC -.->|"TC 后台处理"| N["通知 RM<br/>删 undo_log"] N --> D["<b>DELETE 该行</b><br/>记录消失"] TM2 -.->|"异常 / 超时"| RT["Committing(2)<br/>Rollbacking(4)<br/>CommitRetrying(3) ..."] RT -.->|"重试直至成功"| D style AC fill:#d4edda,stroke:#155724 style RT fill:#fff3cd,stroke:#856404 ``` ### 1.2 删除滞后是特性,不是缺陷 AT 模式的高性能来源于:**一阶段提交本地事务后立即返回,二阶段异步做**。对 AT 而言二阶段提交只是删掉 undo_log,没必要让用户线程等待。 因此 `AsyncCommitting(8)` 这个状态的存在就意味着——记录被删除的时刻,必然晚于用户拿到提交成功的时刻。 问题在于滞后是否**有界**。当业务 TPS(生产速率)持续超过 TC 后台清理(消费速率),有界滞后就变成无界积压。 ### 1.3 积压会雪崩,而非线性劣化 ```mermaid flowchart LR A["表变大"] --> B["按状态+时间排序的<br/>扫描查询变慢"] B --> C["清理速率进一步下降"] C --> A style A fill:#f8d7da,stroke:#721c24 ``` 清理用的查询本身要扫这张表,表越大扫得越慢,清得越慢表又更大。所以现象往往是"平时正常,一到高峰断崖",而不是缓慢劣化——#6615 中的生产监控图就是这个形状。这也解释了为什么调大 `queryLimit` 只能缓解、不能根治。 ### 1.4 影响范围 | 存储模式 | 是否受影响 | |---|---| | **db** | ✅ 主要场景(#6615) | | **redis** | ✅ 同样存在(#7334) | | file / raft | ❌ 内存态 + 日志压缩,无此问题 | 受影响的是高并发 + db/redis 存储的用户,即生产环境的主流配置。 ### 1.5 社区已有工作与结论 | 时间 | 来源 | 内容 | |---|---|---| | 2024-06 | #6615 · PeppaO | 真实生产环境积压问题,附监控图。**至今 OPEN** | | 2024-06 | #6615 · slievrly | 方向建议:调大 queryLimit、异步清理任务、数据分片 | | 2024-06 | #6615 · funky-eyes | 提出分片、raft/redis/db 任务分发三类方案 | | 2024-06 | #6615 · 结论 | **1. 先解决单节点的问题,使用多线程的方式 2. 先不往复杂分布式的情况考虑** | | 2025-05 | #7334 · funky-eyes | 根因诊断:1.4 版本下异步线程只有 1 个、业务线程 50 个 | | 2025-05 | #7362 · YongGoose | 提出多线程清理 + ShardingSphere 分片;收敛到多线程后压测,**TPS 反而下降** | 本方案接续 #7362 的多线程方向的结论:**单节点内优化,不引入分布式协调,不引入外部组件**。 --- ## 2. 两种不同的"慢" 用户看到的"记录堆积",混合了两个性质不同的原因: ```mermaid flowchart LR A["记录产生"] --> B["① 准入延迟<br/>等够 70s / 10s 才有资格被处理<br/>(retryDeadThreshold)"] B --> C["② 排水速率<br/>有资格之后, 清理得多快"] C --> D["记录删除"] style B fill:#fff3cd,stroke:#856404 style C fill:#d4edda,stroke:#155724 ``` | | 是什么 | 处理方式 | |---|---|---| | **① 准入延迟** | Committing/Rollbacking 等 70s、终态残留等 10s 才有资格被扫描 | **不动** —— 防止与属主节点在途处理抢跑的安全边界,是有意设计 | | **② 排水速率** | 有资格之后的清理吞吐 | **本方案的目标** | 占绝大多数的 happy path(AsyncCommitting)**不受准入延迟约束**,可以做到秒级;需要等 70s 的只是故障/重试态。因此优化空间主要在 happy path。 需要明确预期:即使排水速率无限大,终态记录仍有约 10s 的存活下限,**积压不会归零**,而是收敛为一个有界的小常数。 --- ## 3. 为什么"加线程"没有解决问题 ### 3.1 事实基线 #7334 中指出的"异步线程只有 1 个",是 **1.4 版本**的情况。**1.6.1+ 已经并行化**(#5118 两阶段并发通知、#5226 raft):当前 `SessionHelper.forEach` 默认走 parallel stream,单会话处理早已是并行的。 因此在 2.x 上,瓶颈已经换了一个。 ### 3.2 当前的四个瓶颈 ```mermaid flowchart TD B1["① 消费节奏<br/>每周期只拉一页<br/>上限 = 批量 ÷ 周期<br/><b>加线程喂不饱</b>"] B2["② 异质工作耦合<br/>网络 RPC 与数据库删除<br/>绑在同一批线程<br/><b>前者要并发, 后者要批量</b>"] B3["③ 资源共享<br/>后台清理与前台事务<br/>抢同一连接池、同一段索引<br/><b>越并发抢得越凶</b>"] B4["④ 单消费者<br/>分布式锁把清理<br/>钉死为单节点"] ``` **这解释了 #7362 中压测 TPS 下降的现象**:单会话处理本来就是并行的,再叠加线程,增量并发全部落在瓶颈 ③ 上——后台抢走连接池和索引资源,前台事务拿不到,于是前台 TPS 下降。同时不管加多少线程,每周期一页的节奏(瓶颈 ①)仍是天花板。 多线程方向本身没有问题,但"加线程"这个动作治不了这四个病。 --- ## 4. 设计:一个瓶颈一个解法 不引入外部组件,不改协议,不改表结构,全部开关控制、**默认关闭**、可随时回退。 ### 4.1 总览:谁提交、谁处理 核心是把原本挤在一起的工作,拆成**三种线程各司其职**,用队列解耦: ```mermaid flowchart TD subgraph P["生产者(投递工作)"] direction LR FG["<b>前台 Netty 业务线程</b><br/>处理 GlobalCommitRequest<br/>投递后立即响应 TM"] SW["<b>兜底扫描线程</b> ×1<br/>持分布式锁<br/>连续翻页捞故障/漏网记录"] end Q(["<b>二阶段待办队列</b><br/>有界内存队列<br/>满则丢弃, 降级由扫描兜底"]) subgraph C1["消费者一:二阶段执行"] RPC["<b>recovery-rpc 线程池</b> ×N<br/>doGlobalCommit(session, true)<br/>通知 RM 删 undo_log<br/>同步解锁 + 删分支行"] end D(["<b>待删除缓冲</b><br/>攒批: 100 条 或 50ms"]) subgraph C2["消费者二:回收"] CL["<b>cleaner 线程</b> ×1<br/>批量 DELETE ... IN (...)<br/>单写者, 无锁竞争"] end RG{{"<b>资源信号量</b><br/>后台 DB 并发上限<br/>前台不受限"}} DB[("global_table<br/>branch_table")] FG -- "推模式:正常事务" --> Q SW -- "兜底:仅积压部分" --> Q Q --> RPC RPC --> D D --> CL CL --> RG SW -.-> RG RG --> DB style Q fill:#d4edda,stroke:#155724 style D fill:#d4edda,stroke:#155724 style RG fill:#fff3cd,stroke:#856404 ``` **两个生产者、两级队列、两类消费者**,推模式与兜底扫描**共用下游**——不是两套并行实现。 ### 4.2 正常事务的完整时序 ```mermaid sequenceDiagram autonumber participant TM participant FG as 前台 Netty 业务线程 participant Q as 待办队列 participant W as recovery-rpc worker participant RM participant B as 待删除缓冲 participant CL as cleaner 线程 ×1 participant DB as global_table TM->>FG: GlobalCommitRequest FG->>DB: 状态改为 AsyncCommitting FG->>Q: offer(session) 非阻塞 FG-->>TM: 立即返回提交成功 Note over FG: 前台线程到此结束<br/>不参与后续任何处理 W->>Q: take() W->>W: 登记 in-flight 租约 W->>RM: branchCommit(分支二阶段) RM-->>W: 删除 undo_log 完成 W->>DB: 同步解锁 + 删除分支行 W->>B: 投递终态会话 Note over W: worker 到此结束<br/>不执行全局行删除 CL->>B: 攒够 100 条 或 等满 50ms CL->>DB: 批量 DELETE ... WHERE xid IN (...) Note over CL: 单写者<br/>删除之间无锁竞争 ``` 关键点:**前台线程只负责"投递",不负责"处理";worker 只负责"执行二阶段",不负责"删行";删行统一交给单写者。** 每一段职责单一,才能分别按各自的性质调优。 ### 4.3 兜底扫描做什么 扫描线程**不再是主力**,只负责推模式覆盖不到的部分:节点崩溃后遗留的、队列溢出被丢弃的、以及故障重试态的记录。 它与推模式的分工: | | 推模式 | 兜底扫描 | |---|---|---| | 触发 | 事务提交时,由接收节点即时投递 | 定时扫描,持分布式锁单节点执行 | | 覆盖 | happy path(AsyncCommitting) | 故障态 + 漏网记录 | | 拉取方式 | 无需拉取,直接持有会话对象 | 一轮内连续翻页,而非每周期一页 | | 去重 | — | 跳过 in-flight 租约未过期、以及刚变更的记录 | | 下游 | **同一个 rpc 池、同一个 cleaner** | **同上** | ### 4.4 五个组件的线程模型 | 组件 | 执行主体 | 职责 | 治的瓶颈 | |------|---------|------|---------| | **① 推模式入口** | 前台 Netty 业务线程(不新增线程) | 只做一次非阻塞入队,随即返回 | ④ / ① | | **② recovery-rpc 池** | 独立线程池 ×N,每 RM 资源限流 | 执行二阶段 RPC、解锁、删分支行 | ② / ③ | | **③ cleaner** | 独立线程 ×1 | 攒批删除全局行 | ②(删除侧) | | **④ 兜底扫描** | 复用现有调度线程 ×1,持分布式锁 | 连续翻页捞积压,投递给 ② | ① | | **⑤ 资源信号量** | 无独立线程 | 限制后台 DB 并发,前台不受限 | ③ | 设计上的两条原则: - **RPC 要并发、删除要批量** —— 所以 ② 是多线程池,③ 是单线程攒批。这正是瓶颈 ② 的解法:两类工作性质相反,不能共用一套并发策略。 - **后台永远让位前台** —— 所有后台 DB 访问过 ⑤,前台不过。这是"加线程反而变慢"的直接修复。 ### 4.5 为什么不需要分布式协调 推模式的分发规则是**"谁接收谁处理"**:事务由哪个 TC 节点接收提交请求,就由哪个节点执行二阶段。流量本身经过负载均衡,处理压力也就自动按流量分布到各节点——**不需要 leader、不需要任务表、不需要节点间通信**。 #6615 中提出的 raft/redis/db 任务分发方案思路是成立的,但被社区暂缓为"先不往复杂分布式考虑"。本方案没有走那条路,而是用一个更简单的方式获得了同样的负载分散效果。 --- ## 5. 安全性 - **崩溃不丢事务**:后台队列丢失时记录仍在库里,兜底扫描自动接管。这套兜底机制**代码中已存在**(现有 db 模式的延迟删除已依赖它),本方案只是复用。 - **重复处理无害**:二阶段幂等 + 状态原子更新,推模式与扫描偶尔重叠属无害重复。 - **默认关闭、随时回退**:开关关闭时行为与现状完全一致;**无 DDL 变更、无数据迁移、无协议变更**。 - **滚动升级安全**:新老节点可混跑。 - **不影响 Saga/TCC**:推模式只覆盖 AsyncCommitting(AT 分支),其他模式路径不变。 --- ## 6. 当前进展与验收标准 当前阶段完成的是**代码分析与设计**,尚无压测数据。 考虑到 #7362 正是卡在数据环节,推进方式是:先把"为什么会变慢"解释清楚,再动手实现;拆成若干个独立可合的小改动,**每个改动带自己的 benchmark 数据**,而非一个大 PR 后再验证效果。 验收口径: | 目标 | 标准 | |------|------| | 正常事务秒级清理 | 稳态积压不随时间单调增长 | | 不拖累前台 | 开启前后前台 TPS 差异 < 3% | | 故障积压清理提速 | 纯兜底扫描的清理速率 ≥ 现状 5 倍 | 工具:仓库现有 `test-suite/seata-benchmark-cli`,外加 kill -9 混沌验证数据不丢失。 --- ## 7. 代码分析中发现的两个问题 这两点独立于本方案,可能对社区有参考价值: 1. **`doGlobalCommit(session, retrying=false)` 会跳过所有 `canBeCommittedAsync()` 的分支** —— 而这正是 AsyncCommitting 会话的全部内容。任何试图"提前处理异步提交"的实现,若传 `false`,会静默地什么都不做且不报错。 2. **`global_table.gmt_modified` 是 `DATETIME`(秒级),`branch_table` 是 `DATETIME(6)`** —— 两表精度不一致。基于 `gmt_modified` 做 keyset 分页的实现,在高并发同秒多行时会整段漏扫。本方案用复合 keyset `(gmt_modified, xid)` 绕开,但是否统一 DDL 精度值得单独讨论。 --- ## 8. 待讨论的问题 1. **方向**:单节点流水线化(不碰分布式),是否符合 #6615 的既定结论? 2. **推模式**:happy path 绕开轮询,是本方案唯一引入新概念的部分,是否可接受? 3. **准入延迟**:70s/10s 作为有意设计的安全边界不做改动,这一理解是否准确? 4. **目标版本**:2.x 还是更后? 5. **协作**:是否有 committer 可以担任 reviewer?是否邀请 YongGoose 一同推进(#7362 为其提案)? 6. **DDL 精度**:维持现状用复合 keyset 绕开,还是提议统一为 `DATETIME(6)`? --- ## 附:与已有方案的关系 | 方案 | 出处 | 状态 | |------|------|------| | 调大 queryLimit | slievrly @ #6615 | 缓解手段,已在使用 | | 升级 1.6.1+ 获得并行化 | funky-eyes @ #7334 | **已落地**,本方案在此基础上继续 | | ShardingSphere 分库分表 | YongGoose @ #7362 | 循环依赖,已否 | | 每节点不同表名(简易分片) | funky-eyes @ #6615 | 节点故障无法接管,有丢数据风险 | | Leader + 任务表分发(raft/redis/db) | funky-eyes @ #6615 | 社区暂缓("先不往复杂分布式考虑") | | **本方案:单机流水线 + 推模式** | — | 不引组件、不搞协调、对症四个瓶颈、可回退 | GitHub link: https://github.com/apache/incubator-seata/discussions/8211 ---- This is an automatically sent email for [email protected]. To unsubscribe, please send an email to: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
