GitHub user WangzJi edited 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 而言,RM 侧二阶段提交主要是删除 undo_log;TC 
侧还要完成解锁、删除分支行和全局会话等收尾工作,这些操作没必要让用户线程等待。

因此 `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 存储的用户。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/>(beginTime + threshold)"]
    B --> C["② 排水速率<br/>有资格之后, 清理得多快"]
    C --> D["记录删除"]

    style B fill:#fff3cd,stroke:#856404
    style C fill:#d4edda,stroke:#155724
```

| | 是什么 | 处理方式 |
|---|---|---|
| **① 准入条件** | Committing/Rollbacking 等事务年龄达到 70s、终态事务年龄达到 10s 后才有资格被处理;时间从事务 
`beginTime` 起算,不是从进入该状态起算 | **不动** —— 先保持当前语义,其设计意图是否属于防抢跑边界仍需社区确认 |
| **② 排水速率** | 有资格之后的清理吞吐 | **本方案的目标** |

占绝大多数的 happy path(AsyncCommitting)**不受上述准入条件约束**,可以做到秒级;70s 
条件只作用于故障/重试态。因此优化空间主要在 happy path。

终态记录的剩余等待时间实际是 `max(0, beginTime + 10s - now)`。事务本身已运行超过 10s 
时,进入终态后可以立即处理,因此不存在固定的 10s 存活下限;稳态积压能否接近零取决于生产速率、扫描周期和实际排水速率。

---

## 3. 为什么"加线程"没有解决问题

### 3.1 事实基线

#7334 中指出的"异步线程只有 1 个",是对 **1.4 版本**的社区诊断。1.6.1 已经会用 parallel stream 并行处理一批 
`GlobalSession`;但单个 `GlobalSession` 内的多个分支仍是串行处理。#5118 在 2.0.0 
中增加了分支二阶段并发能力,当前由 `server.enableParallelHandleBranch` 控制且默认关闭。#5226 是 Raft 
集群支持,不能作为二阶段并行化的依据。

因此 2.x 的事实基线是:**批次内多个全局会话默认并行,单会话内多个分支默认串行**。分析 #7362 的结果前,还需要确认其压测是否开启了分支并行配置。

### 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 
下降,没有公开线程配置、数据库等待或连接池指标。后台与前台共享连接池和索引资源是一个合理解释,但目前仍是**待压测验证的假设**,不能作为既定根因。可以从代码确认的是:每次调度只查询
 `queryLimit` 限制的一页数据,因此单纯增加下游线程不会突破上游单轮拉取上限。

多线程方向是否有效取决于瓶颈位置;在补齐分阶段吞吐、连接池等待、数据库锁等待和前台 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. 安全性

- **崩溃恢复基础已存在,仍需验证新故障窗口**:后台队列丢失时事务记录仍在持久化存储中,现有兜底扫描可以重新发现记录;但必须通过入队前后、RPC 
前后、分支删除后和全局行删除前的 kill -9 测试确认没有遗漏窗口。
- **重复处理预期可收敛,尚未证明完全无害**:RM 二阶段幂等和 DB 
状态条件更新提供了基础,但推模式与扫描跨节点并发、重复删分支/全局行的行为仍需并发回归测试。
- **默认关闭、可回退是实现约束**:只有在开关关闭路径与现状完全一致,并且实现确实不引入 DDL、数据迁移或协议变更时,才能成立。
- **滚动升级安全是验收目标**:新老节点混跑时,老节点扫描与新节点推模式可能重叠,需要覆盖该场景的集群测试后确认。
- **Saga/TCC 隔离取决于入口约束**:若推模式只在 `AsyncCommitting` 的 AT 路径启用,其他模式路径可以保持不变;实现后仍需 
Saga/TCC 回归测试。

---

## 6. 当前进展与验收标准

当前阶段完成的是**代码分析与设计**,尚无压测数据。

考虑到 #7362 正是卡在数据环节,推进方式是:先把"为什么会变慢"解释清楚,再动手实现;拆成若干个独立可合的小改动,**每个改动带自己的 
benchmark 数据**,而非一个大 PR 后再验证效果。

验收口径:

| 目标 | 标准 |
|------|------|
| 正常事务秒级清理 | 稳态积压不随时间单调增长 |
| 不拖累前台 | 开启前后前台 TPS 差异 < 3% |
| 故障积压清理提速 | 纯兜底扫描的清理速率 ≥ 现状 5 倍 |

工具:仓库现有 `test-suite/seata-benchmark-cli`,外加 kill -9 混沌验证数据不丢失。

---

## 7. 代码分析中发现的两个问题

这两点独立于本方案,可能对社区有参考价值:

1. **`doGlobalCommit(session, retrying=false)` 会跳过所有 `canBeCommittedAsync()` 
的分支** —— 而这正是 AsyncCommitting 会话的全部内容。任何试图"提前处理异步提交"的实现,若传 
`false`,会静默地什么都不做且不报错。
2. **MySQL DDL 中 `global_table.gmt_modified` 是 `DATETIME`(秒级),`branch_table` 是 
`DATETIME(6)`** —— 两表精度不一致。若只使用 `gmt_modified` 作为 keyset 
游标,高并发同秒多行时会漏扫;实现时应使用复合游标 `(gmt_modified, xid)`,并为其他数据库方言分别验证排序和索引行为。是否统一 DDL 
精度值得单独讨论。

---

## 8. 待讨论的问题

1. **方向**:单节点流水线化(不碰分布式),是否符合 #6615 的既定结论?
2. **推模式**:happy path 绕开轮询,是本方案唯一引入新概念的部分,是否可接受?
3. **准入条件**:当前 70s/10s 从事务 `beginTime` 起算,而不是从状态变更时间起算;这一语义及其安全意图是否需要调整?
4. **目标版本**:2.x 还是更后?
5. **协作**:是否有 committer 可以担任 reviewer?是否邀请 YongGoose 一同推进(#7362 为其提案)?
6. **DDL 精度**:维持现状用复合 keyset 绕开,还是提议统一为 `DATETIME(6)`?

---

## 附:与已有方案的关系

| 方案 | 出处 | 状态 |
|------|------|------|
| 调大 queryLimit | slievrly @ #6615 | 提高单轮上限;能否消除积压取决于调整后的消费速率是否超过生产速率 |
| 1.6.1 批次并行、2.0.0 分支并行 | funky-eyes @ #7334、#5118 | 批次内 GlobalSession 
并行已落地;分支并行能力已实现但默认关闭 |
| ShardingSphere 分库分表 | YongGoose @ #7362 | 社区因复杂度和潜在递归依赖否决;ShardingSphere 
并非所有事务模式都依赖 Seata |
| 每节点不同表名(简易分片) | 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]

Reply via email to