# CDC 与复制槽治理

LLMS 索引： [llms.txt](/llms.txt)

---

内建 logical replication 的目标是 PostgreSQL table，CDC 的目标可能是 Kafka、
对象存储、搜索引擎、缓存、数据仓库或另一个业务服务。两者都从 logical decoding
出发，却在“谁保存位点、何时确认、如何处理重复、怎样表达 schema”上走向不同系统。

本节不绑定某一个 connector，而是建立任何 CDC 实现都绕不开的合同。

## 29.2.1 逻辑解码、插件与消费位点 {#item-29-2-1}

### 从物理 WAL 到业务事件

PostgreSQL WAL 记录的是恢复数据库所需的物理/逻辑底层信息，并不是一个稳定的 JSON
业务接口。logical decoding 负责把一个 database 的持久表变化重组为按事务提交边界
可解释的 change stream：

```text
WAL records
  -> decode relation and tuple changes
      -> reorder by transaction
          -> output plugin
              -> logical replication protocol or SQL decoding API
                  -> consumer
```

output plugin 决定编码，不决定下游业务语义：

| plugin / path | 典型用途 | 重要边界 |
|---|---|---|
| `pgoutput` | 内建 publication/subscription | 二进制 replication protocol，不是给人直接读的文本 |
| `test_decoding` | 测试和理解 logical decoding | 测试插件，不是生产业务协议 |
| 第三方 JSON/connector plugin | 通用 CDC | 版本、扩展、schema envelope、升级兼容需独立负责 |

不要把“输出是 JSON”误认为 schema compatibility 已解决。事件仍需定义：

```text
source identity
database / schema / table
transaction and commit position
operation
replica identity before image
new column values
schema version
event id
producer version
```

缺少 before image 时，DELETE 怎样定位？TOAST 列未变化时，connector 如何表示？列
rename 是新字段还是同一字段？这些都不是 JSON 格式本身能回答的。

### SQL 解码 API 与 streaming protocol

教学或诊断可用：

```sql
SELECT *
FROM pg_create_logical_replication_slot(
  'demo_slot',
  'test_decoding'
);

SELECT *
FROM pg_logical_slot_peek_changes(
  'demo_slot',
  NULL,
  100
);

SELECT *
FROM pg_logical_slot_get_changes(
  'demo_slot',
  NULL,
  100
);
```

`peek` 不消费，适合观察；`get` 会推进消费位置。对 `pgoutput` 应使用 binary API 或
replication protocol，而不是强行调用文本函数。

生产 connector 通常通过 replication protocol：

```text
START_REPLICATION SLOT ...
  -> server sends changes
      -> client sends standby status updates
          -> slot confirmed_flush_lsn advances
```

“客户端已经读到”“写进本地 buffer”“写到外部系统”“外部事务提交”“向 PostgreSQL
确认”是五个不同瞬间。connector 必须声明哪个瞬间触发 acknowledgement。

### snapshot 与 stream 必须无缝衔接

新建 logical slot 时可导出一个 snapshot。这个 snapshot 表示：

```text
snapshot sees database state at boundary B
slot emits all committed changes after boundary B
```

于是通用 bootstrap 可以：

1. 创建 slot 并取得 exported snapshot；
2. 在另一个事务 `SET TRANSACTION SNAPSHOT`；
3. 全量导出 snapshot 中的基线；
4. 从 slot 持续消费 B 之后的变化；
5. 在下游把全量与增量合并。

内建 subscription 将这套协调封装在 table sync workers 中。自行开发 CDC 时，如果先
普通 `SELECT` 全表、很久以后才创建 slot，会在两者间漏变化；若先创建 slot却不消费也
不预算磁盘，会让 WAL 无界增长。bootstrap 顺序必须成为协议，不是运气。

### 四个位置不要混为一个“offset”

一个成熟 CDC 管道至少有：

```text
restart_lsn
  publisher still needs WAL from here

confirmed_flush_lsn
  source slot believes consumer confirmed through here

consumer durable checkpoint
  connector can restart from here

sink commit / business watermark
  downstream effects are durable through here
```

它们应满足与具体协议一致的单调关系。最危险的情况是 connector 先确认 source，再把
buffer 异步写 sink：

```text
confirmed_flush_lsn advances
  -> connector crashes before sink commit
      -> source is allowed to recycle old changes
          -> permanent downstream gap
```

反过来，sink 已提交而 checkpoint 尚未推进会造成 replay；这通常可以用幂等处理，永久
缺口却无法凭空修复。因此多数 CDC 更愿意接受 at-least-once，而不是冒险提前确认。

### plugin、slot 与版本都要进 inventory

建议为每个 consumer 记录：

```yaml
consumer_id: search-orders-v3
source_cluster: pg-prod
database: shop
slot: cdc_search_orders_v3
plugin: pgoutput
publication: cdc_search_orders
owner: search-platform
schema_contract: order-event-v7
checkpoint_store: kafka-connect-offsets
ack_after: sink_transaction_commit
max_replay_window: 15m
max_slot_retained_wal: 80GiB
rebootstrap_method: snapshot-plus-stream
retirement_ticket: null
```

slot 本身不知道 consumer owner、SLO 或 sink 状态。没有外部 inventory，inactive slot
只能告诉你“现在没人连”，不能告诉你“可以删”。

## 29.2.2 至少一次、重复事件与下游幂等 {#item-29-2-2}

### PostgreSQL 已经明确允许 replay

logical slot 是 crash-safe 的，但它的当前位置只在 checkpoint 时持久化。服务器崩溃后，
slot 可能回到较早 LSN，于是最近变化再次发送。网络断开、consumer 在 sink commit 后但
ack 前崩溃，也会产生相同结果。

因此正确假设是：

```text
event may be delivered more than once
transaction order is meaningful within a stream
acknowledged history can no longer be requested from that slot
```

“测试十次没重复”不能升级为 exactly-once 保证。

### exactly-once 是端到端属性

若 sink 也是 PostgreSQL，可以在同一个目标事务里同时写业务状态与 dedup ledger：

```sql
BEGIN;

INSERT INTO cdc_applied(event_id, source_lsn, applied_at)
VALUES (:event_id, :source_lsn, clock_timestamp())
ON CONFLICT (event_id) DO NOTHING;

-- 只有上一条确实插入时，才应用业务变化。
UPDATE search_projection
SET ...
WHERE ...
  AND :new_event_was_inserted;

COMMIT;
```

随后才向 source 确认。若事务已提交但 ack 丢失，replay 会命中唯一键，不重复副作用。

但如果副作用跨多个系统：

```text
write database
send email
charge payment
publish Kafka
advance offset
```

没有一个 PostgreSQL transaction 能原子覆盖全部。应使用：

- transactional outbox；
- sink 自身的幂等 key / conditional write；
- 去重 ledger；
- 可重建投影；
- saga/补偿；
- 对不可重复副作用的业务级 request identity。

不要用“处理成功后更新 offset”一句话掩盖多个 commit 之间的崩溃窗口。

### 事件身份怎么设计

只用 `txid` 不够：32-bit XID 会回卷，跨 cluster/database 也会重复。只用表主键也不够：
同一行可以变化很多次。一个实用 envelope 可包含：

```text
source system identifier
database identity
slot / publication contract
commit LSN
transaction identity
change ordinal within transaction
schema contract version
```

具体 connector 能提供哪些字段取决于协议。核心要求是：

```text
same logical change -> same event id on replay
different logical changes -> different event ids
```

不要把消费时间戳或随机 UUID 当 event id；replay 时它们会变。

### upsert 不自动等于幂等

下面的 sink 写法：

```sql
INSERT ... ON CONFLICT (id) DO UPDATE ...
```

只有在“最后写入覆盖即可”且顺序严格时才近似幂等。它处理不了：

- 较旧事件晚到，覆盖较新状态；
- `amount = amount + delta` 被重复执行；
- DELETE/tombstone 后旧 UPDATE 复活；
- 多个 source 同写一个 key；
- 事件 schema 变化导致部分字段保留旧值；
- 外部副作用已经发生。

更安全的投影常带 source version：

```sql
INSERT INTO order_projection (
  order_id, status, source_commit_lsn, payload
)
VALUES (...)
ON CONFLICT (order_id) DO UPDATE
SET status = excluded.status,
    source_commit_lsn = excluded.source_commit_lsn,
    payload = excluded.payload
WHERE order_projection.source_commit_lsn
    < excluded.source_commit_lsn;
```

LSN 只在同一 source timeline/合同内有序；跨 source 合并仍需业务 version/vector 或冲突
规则。

### transaction boundary 不能随意打散

源事务：

```text
debit account A
credit account B
append ledger
```

若 connector 把三条 row event 分别确认并让下游实时可见，中途失败可能暴露不平衡状态。
CDC envelope 应保留 begin/commit 或 transaction grouping；sink 要么原子应用整个事务，
要么明确只提供最终一致读模型并隐藏未完成 batch。

大事务还会带来另一组选择：

- `streaming = off`：源端完整解码后发送，内存/延迟风险；
- `streaming = on`：未提交变化先写 subscriber 临时文件，commit 后应用；
- `streaming = parallel`：有 worker 时直接并行 apply，否则回退临时文件。

吞吐设置不能改变“只有源端 commit 后才把结果视为提交”的语义。

### poison event 需要隔离，不是静默跳过

遇到无法解析或不满足 sink constraint 的事件：

```text
stop entire stream forever
skip and lose data silently
retry at full speed forever
```

都不是完整策略。应记录：

```yaml
event_id: ...
source_position: ...
schema_version: ...
error_class: conversion | constraint | permission | code
first_seen: ...
attempts: ...
raw_payload_ref: encrypted-private-location
owner: ...
decision: repair-and-replay | compensate | approved-skip
```

dead-letter queue 只是隔离区，不是数据正确性的垃圾桶。任何 approved skip 都要进入
reconciliation，并记录业务影响。

## 29.2.3 槽停滞、WAL 保留与磁盘风险 {#item-29-2-3}

### inactive 不等于无害

slot 与连接生命周期独立。consumer 下线后：

```text
active = false
confirmed_flush_lsn stops
source continues writing WAL
restart_lsn remains old
pg_wal retained bytes grow
catalog_xmin may also hold catalog tuples
```

基础查询：

```sql
SELECT slot_name,
       database,
       active,
       active_pid,
       inactive_since,
       xmin,
       catalog_xmin,
       restart_lsn,
       confirmed_flush_lsn,
       pg_size_pretty(
         pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)
       ) AS retained,
       wal_status,
       pg_size_pretty(safe_wal_size) AS safe_wal,
       invalidation_reason
FROM pg_replication_slots
WHERE slot_type = 'logical'
ORDER BY restart_lsn NULLS FIRST;
```

`pg_wal_lsn_diff(current, restart_lsn)` 是按当前时刻估算该 slot 的 WAL 保留距离，不等于
磁盘上所有 WAL 文件恰好这么大；checkpoint、archive、其他 slot 与 segment 粒度都会
影响实际 `pg_wal`。

### 把风险换成时间

若近期 WAL 产生率为 \(r\) bytes/s，slot 当前保留 \(L\)，可用于增长的安全空间为
\(F\)，则最粗略的时间预算：

$$
T_{\text{disk}}
\approx
\frac{F}{r}
$$

若配置 `max_slot_wal_keep_size = M`，距离 slot 可能在 checkpoint 后失去所需 WAL 的
预算：

$$
T_{\text{slot}}
\approx
\frac{M - L}{r}
$$

实际告警应使用变化率和低水位：

```text
retained bytes
safe_wal_size
wal_status / invalidation_reason
pg_wal filesystem free
archive health
consumer lag and error rate
catalog_xmin / XID age
```

只告警 `active = false` 会在计划维护时产生噪声；只告警磁盘使用率会在 WAL 已快速增长时
太晚。

### 上限保护的是 source，不保证 consumer 可恢复

`max_slot_wal_keep_size` 为非负值时，checkpoint 可以允许回收超出上限的 WAL，slot
可能进入 `unreserved` 乃至 `lost`。这避免一个遗忘 consumer 无限填满源盘，但代价是
consumer 必须 rebootstrap。

PostgreSQL 18 还提供 `idle_replication_slot_timeout`：slot inactive 超过时限后可在
checkpoint 被 invalidated。它同样不是“暂停后自动保存到对象存储”。启用前必须明确：

- 哪类 slot 允许因空闲失效；
- 谁收到即将失效告警；
- snapshot/rebootstrap 需要多久；
- 同步到 standby 的 slot 等豁免/特殊语义；
- checkpoint 周期带来的执行延迟。

不要为了保护磁盘把上限设得小于正常故障恢复窗口，然后把频繁 `lost` 当 consumer 问题。

### 本章停滞实验

正式实验先让 20,000 初始订单和 500/200/100 增量收敛，再执行：

```sql
ALTER SUBSCRIPTION pg36_shop_sub DISABLE;
```

确认目标 `subenabled = false`、worker 为零、源端 slot `active = false` 后，在源端插入
3,000 行有界 payload。

结果：

| 字段 | 写入前 | 写入后 |
|---|---:|---:|
| `active` | false | false |
| `confirmed_flush_lsn` | `1/9FDBBF30` | `1/9FDBBF30` |
| retained bytes | 227,008 | 2,867,128 |
| `wal_status` | `reserved` | `reserved` |

这里能确认因 consumer 停滞，确认位点没有推进且保留距离增长。不能用这 2.87 MB 推导
生产每 3,000 行的固定 WAL 成本：其他 database、full-page write、索引、payload、
checkpoint 和并发都会改变 WAL。

恢复：

```sql
ALTER SUBSCRIPTION pg36_shop_sub ENABLE;
```

验收同时要求：

```text
confirmed_flush_lsn >= source marker LSN
both relation states = r
source logical manifest = target logical manifest
```

“slot active 又变 true”仍不够。

### 停滞处置顺序

```text
1. identify exact slot and owner
2. confirm source database / plugin / consumer contract
3. capture restart, confirmed, xmin, wal_status, safe_wal_size
4. inspect consumer, network, auth, schema and apply conflicts
5. calculate disk and slot-loss deadlines
6. choose resume, repair, rebootstrap or retire
7. verify acknowledgement plus semantic convergence
8. only then close or drop the slot
```

优先恢复 consumer，而不是先 drop slot。若 slot 已 lost，继续重试同一位点没有意义；
冻结下游写入，按 snapshot + new slot 的协议重建。若 consumer 已永久退役，保留审批与
最后消费位置后精确 drop。

### 与 Pigsty 观测对齐

Pigsty 的 PGSQL Replication / Persist / Instance / Alert 看板可把：

```text
slot retention
WAL production
archive
disk free
physical replica lag
logical pub/sub
host I/O
```

放在同一时间轴。原生 SQL 则确认 slot、subscription、table state 与 conflict identity。
平台看板用于发现趋势，不能替代 consumer owner 和业务 reconciliation。

进一步阅读：

- [PostgreSQL 18：Logical Decoding Concepts](https://www.postgresql.org/docs/18/logicaldecoding-explanation.html)
- [PostgreSQL 18：Replication Settings](https://www.postgresql.org/docs/18/runtime-config-replication.html)
- [PostgreSQL 18：`pg_replication_slots`](https://www.postgresql.org/docs/18/view-pg-replication-slots.html)
- [PostgreSQL 18：Logical Replication Monitoring](https://www.postgresql.org/docs/18/logical-replication-monitoring.html)
- [Pigsty：PGSQL Dashboards](https://pigsty.io/docs/pgsql/dashboard/)

---

[上一节：逻辑复制原语](../01/) · [返回本章目录](../) · [下一节：批量装载与数据校验](../03/) ·
[查看全书目录](/toc/) · [查看索引中心](/indexes/)
