Outbox、CDC 与领域事件:把数据库提交送入事件流
订单已经写入数据库,发往消息系统的通知却尚未完成。进程在这时退出,消费者就可能一直不知道这笔订单。Outbox 将“需要发布哪一个事件”也保存为数据库行,与订单一起提交;后续发布器可以根据这行记录继续工作。
业务事实与事件记录共同提交
两种直接双写各留下什么窗口
先写数据库
订单 commit → 进程退出 → 消息尚未发送
结果:订单存在,消费者缺少事件
先发消息
消息可消费 → 订单事务 rollback
结果:消费者收到最终未成立的订单
Outbox
同一事务:INSERT order + INSERT outbox → commit
结果:订单与待发布事件同时存在,发布器稍后读取同一 PostgreSQL 事务里的两条写入共同提交或共同回滚。应用仍然需要处理本地 commit 结果未知,但可以按订单 ID、事件 ID 查询持久记录,决定后续动作。PostgreSQL 事务
Outbox 改变的是生产侧的提交方式。消费者更新另一套数据库时,还要自行处理幂等与确认;相关失败窗口见 投递、幂等与顺序。
保存一份完整、稳定的事实
常用 Outbox 行包含事件 ID、聚合类型、聚合 ID、事件类型、业务 payload、契约版本和发生时间。event ID 在业务事务中确定,Relay 重试时复用;aggregate ID 可作为 Kafka key,让同一业务对象的事件采用相同分区选择规则。
payload 应保存此次变更已经确定的内容。例如 OrderCreated 保存创建时金额,而非只留下订单号,让发布器稍后查询最新金额。后续修改订单不应偷偷改变原事件的含义。
CREATE TABLE lab.outbox(
id uuid PRIMARY KEY,
aggregatetype text NOT NULL,
aggregateid text NOT NULL,
type text NOT NULL,
payload jsonb NOT NULL,
occurred_at timestamptz NOT NULL DEFAULT now()
);这张表是 CDC 使用的追加式事件表。发布后不对原行写 published=true;轮询发布器使用另一张 polling_outbox 表,避免两种发布流程争用同一行的语义。
表变更事件与集成事件
直接对订单表做 CDC,通常得到的是行字段变化,可能包含内部列、存储枚举和数据库删除语义。领域事件由业务决定何时成立、叫什么,以及消费者需要什么字段。Outbox 让应用先形成这份事件,再由 CDC 搬运。
Debezium 的 Event Router 根据 Outbox 列设置目标 Topic、key、header 和 value,不替应用判断订单是否真的创建成功。其标准处理面向 INSERT;UPDATE 有额外处理策略,DELETE 被过滤,用于清理旧行。Outbox Event Router
运行 PostgreSQL、Debezium 与 Kafka
版本与操作环境
下载完整源码工程,解压进入 message-event-lab。使用 Linux shell、Docker Engine、Compose v2+、jq,以及有 Docker 权限的普通宿主用户。Maven 3.9.12 使用 Java 25,源码编译到 Java 17。
实验固定 PostgreSQL 18.6、Kafka Broker/Java client 4.3.1、Debezium Connect 镜像 3.6.2.Final、pgJDBC 42.7.13 与 Jackson 2.21.4。Debezium 3.6 的支持资料覆盖 PostgreSQL 14–18;镜像内的 Connect runtime 是其自带版本,不能用外部 Broker 版本冒充内部 runtime 版本。Debezium 3.6 支持范围
compose.kafka.yaml
└── 一个 controller + 三个 Broker,数据 Topic 可使用 RF=3
compose.cdc.yaml(叠加文件)
├── PostgreSQL:wal_level=logical
└── Connect:Debezium PostgreSQL connector + Outbox SMT
config/init-cdc.sql
├── app:业务表和 Outbox 的应用拥有者
├── cdc_reader:LOGIN + REPLICATION,只读 lab.outbox
└── me12_outbox_pub:管理初始化时显式创建
kafka/src/main/java/example/
├── KafkaLabSupport.java
└── OutboxLab.java这套 Compose 约需要 4.5 GiB 可用容器内存。一个 controller 没有元数据高可用,数据库也只有一个实例。所有服务仅在 me12-cdc 的 Docker 网络通信,不将 Connect 管理端口、数据库或 Kafka 暴露到宿主。
PostgreSQL 逻辑复制读取 WAL,需要 logical 设置、复制槽与足够的发送进程。Connector 角色不使用超级用户,也不负责自动创建 publication;初始化文件由数据库管理身份完成这些准备。Debezium PostgreSQL connector
启动与构建
test "$(id -u)" -ne 0 || { echo '请使用普通宿主用户'; exit 1; }
command -v jq >/dev/null || { echo '请先安装 jq'; exit 1; }
mkdir -p .m2
dc() { docker compose -f compose.kafka.yaml -f compose.cdc.yaml "$@"; }
dc up -d
docker run --rm --user "$(id -u):$(id -g)" \
-e MAVEN_CONFIG=/maven-cache \
-v "$PWD:/workspace" -v "$PWD/.m2:/maven-cache" -w /workspace \
maven:3.9.12-eclipse-temurin-25 \
mvn -B -ntp '-Dmaven.repo.local=/maven-cache' '-Duser.home=/tmp' \
-pl kafka -am clean package dependency:copy-dependencies
connect_api() {
docker run --rm --user "$(id -u):$(id -g)" --network me12-cdc_default \
-v "$PWD:/workspace" -w /workspace curlimages/curl:8.21.0 \
-q --noproxy '*' --connect-timeout 3 --max-time 20 "$@"
}
ready=0
for attempt in $(seq 1 60); do
if connect_api --fail-with-body -sS http://connect:8083/connector-plugins \
-o .lab-plugins.json &&
jq -e '.[] | select(.class=="io.debezium.connector.postgresql.PostgresConnector")' \
.lab-plugins.json >/dev/null; then
ready=1
break
fi
sleep 2
done
test "$ready" -eq 1 || { dc logs --tail=60 connect postgres; exit 1; }BUILD SUCCESS 说明 Java 实验已构建。HTTP 检查还要求插件列表实际包含 PostgreSQL connector;容器处于 running 状态时,Connect 可能仍在初始化 Kafka 内部 Topic。
Connector 配置与字段路由
先显式建立本实验的业务 Topic,避免自动建 Topic 使用不符合预期的副本数:
dc exec -T kafka1 /opt/kafka/bin/kafka-topics.sh --bootstrap-server kafka1:9092 \
--create --if-not-exists --topic me12.events.orders --partitions 1 \
--replication-factor 3 --config min.insync.replicas=2
dc exec -T kafka1 /opt/kafka/bin/kafka-topics.sh --bootstrap-server kafka1:9092 \
--describe --topic me12.events.orders描述应显示 ReplicationFactor=3、min.insync.replicas=2,且三副本进入 ISR。已有同名 Topic 时 --if-not-exists 不会修改其配置,应先检查实际状态,不把命令成功视为已经完成副本调整。
config/outbox-connector.json 的完整内容:
{
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"tasks.max": "1",
"database.hostname": "postgres",
"database.port": "5432",
"database.user": "cdc_reader",
"database.password": "cdc-lab-only",
"database.dbname": "lab",
"topic.prefix": "me12_source",
"plugin.name": "pgoutput",
"slot.name": "me12_outbox_slot",
"publication.name": "me12_outbox_pub",
"publication.autocreate.mode": "disabled",
"table.include.list": "lab.outbox",
"snapshot.mode": "initial",
"heartbeat.interval.ms": "1000",
"key.converter": "org.apache.kafka.connect.storage.StringConverter",
"value.converter": "org.apache.kafka.connect.json.JsonConverter",
"value.converter.schemas.enable": "false",
"transforms": "outbox",
"transforms.outbox.type": "io.debezium.transforms.outbox.EventRouter",
"transforms.outbox.table.expand.json.payload": "true",
"transforms.outbox.table.fields.additional.placement": "type:header:eventType",
"transforms.outbox.route.topic.replacement": "me12.events.${routedByValue}",
"transforms.outbox.predicate": "isOutbox",
"predicates": "isOutbox",
"predicates.isOutbox.type": "org.apache.kafka.connect.transforms.predicates.TopicNameMatches",
"predicates.isOutbox.pattern": "me12_source[.]lab[.]outbox"
}Outbox 的 aggregatetype=orders 生成 me12.events.orders;aggregateid 成为字符串 key;id 保存在 id header,type 额外放入 eventType header。payload 展开为 JSON 对象,配合关闭 schemas 的 JsonConverter 输出业务 JSON。
predicate 只匹配原始 lab.outbox 数据 Topic,使心跳和其他元数据记录不误入 Outbox SMT。表中的 occurred_at 留作业务审计;当前配置没有把它指定为 Kafka 时间戳。若需要覆盖事件时间,应先确认对应列的 Debezium 类型和转换器要求,不能把 timestamptz 的字符串映射直接交给期望 INT64 的事件时间字段。
status=$(connect_api --fail-with-body -sS \
-X PUT -H 'Content-Type: application/json' \
--data-binary @config/outbox-connector.json \
-o .lab-config-response.json -w '%{http_code}' \
http://connect:8083/connectors/me12-outbox/config) || exit 1
case "$status" in 200|201) ;; *) echo "unexpected HTTP $status"; exit 1;; esac
jq -e '.name=="me12-outbox"' .lab-config-response.json >/dev/null
running=0
for attempt in $(seq 1 60); do
if connect_api --fail-with-body -sS -o .lab-status.json \
http://connect:8083/connectors/me12-outbox/status &&
jq -e '.connector.state=="RUNNING" and (.tasks|length)==1 and
all(.tasks[]; .state=="RUNNING")' .lab-status.json >/dev/null; then
running=1
break
fi
sleep 2
done
test "$running" -eq 1 || { jq . .lab-status.json; exit 1; }PUT 成功只表示配置已接收。Connector 和任务是两个状态:connector=RUNNING、task=FAILED 时,业务变更仍无法输出。失败 trace 应继续查看账号权限、publication、slot 和 SMT 字段类型。Connect 管理 API 与任务状态
首次提交、回滚和消费结果
run_outbox() {
docker run --rm --user "$(id -u):$(id -g)" --network me12-cdc_default \
-v "$PWD:/workspace:ro" -w /workspace eclipse-temurin:25.0.4_7-jdk \
java '-Dorg.slf4j.simpleLogger.defaultLogLevel=warn' \
-cp 'kafka/target/classes:kafka/target/dependency/*' example.OutboxLab "$@"
}
run_outbox seed
run_outbox observedb: committed=1 rolledBack=1
cdc: uniqueOrders=1 records=1 rollbackInvisible=true keyHeaderPayloadMatch=trueseed 在本地事务中分别写入订单与 Outbox:committed-42 提交,rolled-back-42 回滚。observe 通过真实 KafkaConsumer 读取事件,核对聚合 key、id/eventType header、JSON 金额和订单 ID,收到预期事件后继续观察三秒以检查额外可见记录。
seed 使用固定业务主键,应在此实验第一次初始化时执行一次。重复执行会得到唯一键冲突,不应把异常改成忽略后继续声称新事件已写入。observe 不提交消费位点,每次从 earliest 观察仍保留的事件;重放产生重复时 records 可以增加,但 uniqueOrders 与事件身份必须正确。
完整 Java 入口
package example;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import java.nio.charset.StandardCharsets;
import java.sql.DriverManager;
import java.time.Duration;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.TimeUnit;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.TopicPartition;
import static example.KafkaLabSupport.*;
public final class OutboxLab {
static final ObjectMapper JSON=new ObjectMapper();
static java.sql.Connection database() throws Exception {
return DriverManager.getConnection("jdbc:postgresql://postgres:5432/lab","app","app-lab-only");
}
static String eventId(String name) {return UUID.nameUUIDFromBytes(name.getBytes(StandardCharsets.UTF_8)).toString();}
static void write(String order,boolean commit) throws Exception {
try(var c=database()) {
c.setAutoCommit(false);
try(var business=c.prepareStatement("INSERT INTO cdc_order VALUES (?,1250)");
var event=c.prepareStatement("INSERT INTO outbox(id,aggregatetype,aggregateid,type,payload) VALUES (?::uuid,'orders',?,'OrderCreated',?::jsonb)")) {
business.setString(1,order);business.executeUpdate();
event.setString(1,eventId(order));event.setString(2,order);
event.setString(3,JSON.writeValueAsString(Map.of("orderId",order,"amountMinor",1250,"currency","CNY","schemaVersion",1)));
event.executeUpdate();if(commit)c.commit();else c.rollback();
} catch (Exception | Error failure) {
try { c.rollback(); }
catch (java.sql.SQLException rollbackFailure) { failure.addSuppressed(rollbackFailure); }
throw failure;
}
}
}
static void observe(Set<String> expected) throws Exception {
try(var c=consumer("me12-outbox-reader","read_committed")) {
c.subscribe(List.of("me12.events.orders"));
Set<String> seen=new HashSet<>();int records=0;
long deadline=System.nanoTime()+TimeUnit.SECONDS.toNanos(60),completeAt=Long.MAX_VALUE;
while(System.nanoTime()<deadline&&System.nanoTime()<completeAt) {
for(var r:c.poll(Duration.ofMillis(200))) {
JsonNode payload=JSON.readTree(r.value());String order=payload.path("orderId").asText();
check(expected.contains(order),"unexpected event (including rollback): "+order);
check(order.equals(r.key()),"aggregate ID not used as Kafka key");
check(payload.path("amountMinor").asInt()==1250,"payload changed");
var id=r.headers().lastHeader("id");var type=r.headers().lastHeader("eventType");
check(id!=null&&eventId(order).equals(new String(id.value(),StandardCharsets.UTF_8)),"event ID header missing");
check(type!=null&&"OrderCreated".equals(new String(type.value(),StandardCharsets.UTF_8)),"event type header missing");
seen.add(order);records++;
}
if(seen.equals(expected)&&completeAt==Long.MAX_VALUE)completeAt=System.nanoTime()+TimeUnit.SECONDS.toNanos(3);
}
check(seen.equals(expected),"missing events: expected="+expected+" seen="+seen);
System.out.println("cdc: uniqueOrders="+seen.size()+" records="+records+" rollbackInvisible=true keyHeaderPayloadMatch=true");
}
}
static void relay() throws Exception {
resetTopic("me12.polling.orders",1);
try(var c=database();var s=c.createStatement()) {
s.execute("CREATE TABLE IF NOT EXISTS polling_outbox(event_id text PRIMARY KEY,payload text NOT NULL,published boolean NOT NULL DEFAULT false)");
s.execute("INSERT INTO polling_outbox VALUES('relay-42','OrderCreated',false) ON CONFLICT(event_id) DO UPDATE SET published=false");
}
try(var p=producer()) {
for(int pass=0;pass<2;pass++) {
try(var db=database();var s=db.createStatement()) {
db.setAutoCommit(false);
try(var row=s.executeQuery("SELECT event_id,payload FROM polling_outbox WHERE published=false ORDER BY event_id FOR UPDATE SKIP LOCKED LIMIT 1")) {
check(row.next(),"pending row absent");String id=row.getString(1),payload=row.getString(2);
p.send(new ProducerRecord<>("me12.polling.orders",id,payload)).get(15,TimeUnit.SECONDS);
if(pass==0) {db.rollback();} // Simulate losing the DB transaction after broker acknowledgement.
else {try(var mark=db.prepareStatement("UPDATE polling_outbox SET published=true WHERE event_id=?")) {mark.setString(1,id);check(mark.executeUpdate()==1,"mark failed");}db.commit();}
}
}
}
}
try(var c=consumer("me12-relay-reader","read_committed")) {
c.assign(List.of(new TopicPartition("me12.polling.orders",0)));c.seekToBeginning(c.assignment());
var rows=read(c,2).records();check(rows.stream().allMatch(r->r.key().equals("relay-42")),"event identity changed");
System.out.println("relay: acknowledgedCopies=2 stableEventId=true pendingRowRecovered=true");
}
}
record LegacyOrder(String orderId,long amountMinor){}
static LegacyOrder decodeV1(String json) throws Exception {
JsonNode node=JSON.readTree(json);
check(node.hasNonNull("orderId")&&node.path("amountMinor").isIntegralNumber()&&node.path("amountMinor").canConvertToLong(),"required field missing or invalid");
check(!node.has("currency")||node.path("currency").asText().equals("CNY"),"unsupported currency");
return new LegacyOrder(node.get("orderId").asText(),node.get("amountMinor").longValue());
}
static void schema() throws Exception {
var old=decodeV1("{\"orderId\":\"order-42\",\"amountMinor\":1250}");
var newer=decodeV1("{\"orderId\":\"order-42\",\"amountMinor\":1250,\"currency\":\"CNY\",\"source\":\"web\"}");
check(old.equals(newer),"additive fields changed old reader result");
boolean rejected=false;
try {decodeV1("{\"orderId\":\"order-42\",\"amount\":12.50}");}catch(IllegalStateException expected){rejected=true;}
check(rejected,"breaking rename accepted");
System.out.println("schema: oldAndNewPayloadRead=true additiveFieldsIgnored=true breakingRenameRejected=true");
}
public static void main(String[] args) throws Exception {
check(args.length>=1,"mode: seed|offline-write|observe|relay|schema");
switch(args[0]) {
case "seed" -> {write("committed-42",true);write("rolled-back-42",false);System.out.println("db: committed=1 rolledBack=1");}
case "offline-write" -> {write("offline-43",true);System.out.println("db: committedWhileConnectorStopped=true");}
case "observe" -> observe(args.length==2&&args[1].equals("recovered")?Set.of("committed-42","offline-43"):Set.of("committed-42"));
case "relay" -> relay();
case "schema" -> schema();
default -> throw new IllegalArgumentException("unknown mode");
}
}
}连接器恢复与轮询发布器
WAL、复制槽与 Connect offset 各记什么
PostgreSQL 提交记录进入 WAL,逻辑解码按 publication 输出目标表变更。初次运行通常先建立一致性快照,再继续流式读取。复制槽让 PostgreSQL 保留尚需读取的日志;Connect 的 offset Topic 持久化源位置,以便任务或进程重启后续读。PostgreSQL 逻辑解码与复制槽
三个位置应分别观察:数据库当前 WAL 位置、复制槽的 restart_lsn/confirmed_flush_lsn,以及 Connector 已保存的源 offset。它们的推进条件不同,也不能把一个 LSN 当作消费者已经完成业务的凭证。
停止连接器进程,保留其内部 Topic 和数据库卷:
dc stop connect
run_outbox offline-write
dc start connect预期 offline-write 输出 committedWhileConnectorStopped=true。等待前面的任务状态循环再次通过后执行:
run_outbox observe recoveredcdc: uniqueOrders=2 records=2 rollbackInvisible=true keyHeaderPayloadMatch=true新增 offline-43 在连接器停机期间已提交。恢复后读取到 committed-42 与 offline-43 两个事件,回滚订单仍未出现;整个过程没有由 Java 再次发送这两条业务消息。记录数可能因连接器恢复重复而增加,因此测试按 event ID 计唯一业务事件。
检查复制槽:
dc exec -T postgres psql -U labadmin -d lab -v ON_ERROR_STOP=1 \
-c "SELECT slot_name, active, restart_lsn, confirmed_flush_lsn, pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn) AS retained_wal_bytes FROM pg_replication_slots WHERE slot_name='me12_outbox_slot';"连接器停机时 active=false 是正常状态;retained_wal_bytes 持续增长说明数据库正在为槽保留更多日志。保留量还受其他数据库活动影响,不能直接换算成未发送事件条数。
“已发出但未记账”的重复窗口
run_outbox relaypolling_outbox 单独维护 published 标记。Relay 用 FOR UPDATE SKIP LOCKED 领取一行,向 Kafka 发布并等待确认。第一次故意回滚数据库事务,第二次重新领取同一行,发送后提交 published=true。
relay: acknowledgedCopies=2 stableEventId=true pendingRowRecovered=true两条 Kafka 记录均携带 relay-42。它们来自两次不同的 send;即使生产者启用了 Kafka 幂等生产,应用层重新读取并发送同一个事件仍可产生多条记录。消费者需要按业务身份去重。
SKIP LOCKED 适合多个发布实例各自领取待处理行,但本例为了直观,在等待 Broker 时一直持有数据库行锁。生产可使用短事务领取租约、批量发送和条件更新结束状态,代价是增加租约过期与重复领取处理。抢占方式与索引设计可查 PostgreSQL SELECT 锁定子句。
CDC 将发布工作放到日志读取和 Connect 运行时,减少应用的轮询代码;它依然需要管理复制槽、Connector 位点、重启重复和 schema 变化。没有必要同时让 Relay 和 CDC 对同一事件表发消息。
清理不能破坏恢复起点
Outbox 历史行、数据库 WAL、Connect 内部 Topic、Kafka 业务 Topic 和消费者 Inbox 各有自己的保存期。清理 Outbox 行不等于已释放复制槽保留的 WAL;删掉 slot 也不会删除 Kafka 中已经发布的事件。
Kafka 发布成功后立刻删行,会失去查询和修复素材。清理策略宜在连接器已越过相关提交、业务审计保留满足之后分批执行,并考虑数据库备份恢复重新出现旧 Outbox 行的情况。消费者去重保存期至少覆盖可能的重放范围。
契约演进、监控与故障处理
旧消费者实际读取新 JSON
run_outbox schemaschema: oldAndNewPayloadRead=true additiveFieldsIgnored=true breakingRenameRejected=true代码用 Jackson 解析 JSON,再按旧消费者要求校验 orderId 和 amountMinor。增加可忽略的 source 字段以及已支持的 currency=CNY 后,旧消费者仍读出相同结果;将 amountMinor 改名为 amount 则被拒绝。Jackson Databind 官方项目
这个通过条件来自明确的读取策略,不是所有 JSON 消费者的默认能力。严格字段映射可能拒绝未知字段,枚举新增也可能导致旧代码异常。新增字段先可选,重命名通常需要双字段过渡,金额单位和时区变化则属于语义变更,即使类型未变也需要迁移设计。
本例金额始终是 CNY 的最小货币单位。若将同一个整数解释为元,语法解析仍会通过,但业务结果已经改变。契约测试还应覆盖单位、空值、删除含义、历史版本和不认识的枚举值。
可查询的发布状态
CDC 没有对 Outbox 行设置 published=true。判断某事件是否已经进入 Kafka,可以通过 event ID 的消费记录、审计存储或专用查询索引,结合 Connector 位点判断;不能仅查询“Outbox 行存在”便宣布通知完成。
监控通常包括最老未发布/未观察事件年龄、Connector task 状态、复制槽保留 WAL 字节、重启重复率、Kafka 发送错误,以及消费者最终业务状态。集群 Topic 中有记录,只说明事件已进入消息系统,下游可能仍在等待数据库或外部服务。
| 现象 | 检查与下一步 |
|---|---|
| 数据库有订单,没有 Outbox 行 | 检查写入是否确实处于同一事务、是否存在绕过路径;补事件前核对历史发送 |
| connector=RUNNING、task=FAILED | 读取任务 trace;核对权限、字段类型和转换器,不只重启 worker |
| WAL 磁盘持续增长 | 查 slot 是否活跃、位点是否前进、Connect 是否有发送错误;先恢复消费者再规划清理 |
| 事件重复 | 比对稳定 id 与业务结果,检查重启或发送后 offset 保存窗口;不要直接生成新 ID |
| 事件到达顺序变化 | 查同聚合是否使用同一 key、是否扩分区、多个写入源是否有业务版本 |
| 新字段导致旧消费者失败 | 用保存的原始样本运行旧版本解码器,修复契约或实施兼容迁移 |
| 复制槽所需 WAL 已被移除 | 按可用备份和快照流程重建起点,并计划历史重放去重;不能凭空跳过丢失区间 |
权限和结束实验
连接器配置包含实验口令,只在隔离网络内使用。真实部署应使用密钥配置提供器,限制 Connect 管理端点访问,并分别限制数据库复制角色、Kafka 发布权限和应用业务权限。
dc stop 保留实验记录。若要停用此连接器但保留数据库,应先删连接器并确认任务停止,再用管理身份删除不再需要的 slot:
status=$(connect_api -sS -X DELETE -o .lab-delete-response.txt -w '%{http_code}' \
http://connect:8083/connectors/me12-outbox) || exit 1
test "$status" = 204 || { echo "unexpected HTTP $status"; exit 1; }
dc exec -T postgres psql -U labadmin -d lab -v ON_ERROR_STOP=1 \
-c "SELECT pg_drop_replication_slot('me12_outbox_slot');"slot 仍被使用时删除会失败,应等待任务真正退出,不能强行中断未知复制客户端。只有确认全部 me12-cdc 数据可丢弃时才执行 dc down --volumes,它删除该 Compose 项目的 Kafka 和数据库卷,也删除内部 offset/config/status Topic。源码与 Maven 缓存不受影响。
权威资料与规范地址
| 查阅内容 | 官方地址 |
|---|---|
| 数据库事务 | https://www.postgresql.org/docs/18/tutorial-transactions.html |
| Outbox 字段、路由与 SMT 限制 | https://debezium.io/documentation/reference/stable/transformations/outbox-event-router.html |
| Debezium 3.6 支持范围 | https://debezium.io/releases/3.6/ |
| PostgreSQL Connector 权限与配置 | https://debezium.io/documentation/reference/stable/connectors/postgresql.html |
| Connect 管理与任务状态 | https://kafka.apache.org/43/kafka-connect/administration/ |
| WAL 解码与复制槽 | https://www.postgresql.org/docs/18/logicaldecoding-explanation.html |
| SKIP LOCKED 与查询锁定 | https://www.postgresql.org/docs/18/sql-select.html |
| JSON 解析和对象映射 | https://github.com/FasterXML/jackson-databind |
