一、从WAL到逻辑复制:PostgreSQL CDC的底层机制
PostgreSQL的CDC能力建立在预写日志(Write-Ahead Logging,WAL)之上。每一次数据变更都会先写入WAL,再应用到数据文件,这保证了崩溃恢复。逻辑复制正是利用WAL中的信息,通过一个逻辑解码插件将物理变更还原成行级操作。默认的插件是pgoutput,它输出的消息包含事务边界、表的relation信息以及每一行的前后镜像。要使用逻辑复制,需要将wal_level参数设置为logical,并且创建复制槽。复制槽会保留尚未被消费者确认的WAL,防止数据丢失。

从使用角度看,逻辑复制流由一系列二进制消息组成,每条消息有类型标识。常见类型包括B(Begin,事务开始)、C(Commit,事务提交)、R(Relation,表结构信息)、I(Insert)、U(Update)、D(Delete)。这些消息通过复制协议的流式接口发送给客户端。你不需要关心WAL的物理格式,只需要按照协议解析消息即可。这也是Debezium等工具能够实现数据库变更捕获的基础。
需要特别注意的是,逻辑复制发送的是行级变化,而不是SQL语句。这带来两个好处:一是不会受到触发器或者某些DDL的影响,二是可以获取到每一列的新旧值,便于下游做精确处理。但代价是复制槽会持续占用WAL空间,如果消费者长时间不确认,可能导致磁盘膨胀。因此设计消费者时必须考虑确认策略和心跳机制。
二、方案选型:Debezium Server与自研逻辑复制消费者
把PostgreSQL CDC数据放进RabbitMQ,常见做法有两种。第一种是使用Debezium Server。Debezium通常与Kafka配合,但它也提供了独立的Server模式,可以配置RabbitMQ作为sink,直接把变更事件推送到指定交换机。这种方式无需编写任何Java代码,只需要调整配置文件和依赖即可,适合团队快速验证或统一数据管道。Debezium Server内部仍然使用PostgreSQL的逻辑复制,并将其封装成结构化的JSON事件,事件中包含schema、payload等信息。
第二种做法是基于PostgreSQL JDBC驱动的复制API,自己实现一个消费者。这样做的优点是轻量、可控,不需要引入Debezium的完整依赖,可以更灵活地定制消息格式、过滤规则和投递逻辑。对于只需要同步少量表、且团队有一定Java基础的场景,自研消费者可能更合适。缺点是需要处理pgoutput的二进制协议、事务边界以及恢复时的offset管理。
下面给出一段Debezium Server的RabbitMQ配置示例,帮助理解第一种方案。配置中需要指定数据库连接、复制槽、以及RabbitMQ的地址和交换机名称。
debezium.sink.type=rabbitmq debezium.sink.rabbitmq.host=localhost debezium.sink.rabbitmq.port=5672 debezium.sink.rabbitmq.username=guest debezium.sink.rabbitmq.password=guest debezium.sink.rabbitmq.exchange=cdc.exchange debezium.source.connector.class=io.debezium.connector.postgresql.PostgresConnector debezium.source.database.hostname=localhost debezium.source.database.port=5432 debezium.source.database.user=postgres debezium.source.database.password=postgres debezium.source.database.dbname=mydb debezium.source.plugin.name=pgoutput debezium.source.slot.name=debezium_slot debezium.source.publication.name=debezium_pub
如果选择Debezium Server,还需要在PostgreSQL中提前创建发布(Publication),并将需要同步的表加入发布。Debezium会自动创建复制槽,但也可以手动指定。这种方式的消息顺序由Debezium保证,但需要注意RabbitMQ的ack机制,避免消费者异常导致消息积压。
三、实现一个最小可用的Java逻辑复制消费者
如果你希望完全掌控同步过程,可以直接使用PostgreSQL JDBC驱动的复制API。下面的代码展示了如何连接数据库、读取已有复制槽的变更流,并将解析后的变更消息发布到RabbitMQ。为了简化示例,这里采用默认的pgoutput插件,消息解析只做示意,实际生产环境需要解析完整的消息类型。
首先需要在PostgreSQL中开启逻辑复制支持。修改postgresql.conf中的wal_level为logical,重启数据库,然后创建发布和复制槽。示例SQL如下:
ALTER SYSTEM SET wal_level = logical;
-- 重启数据库后执行
CREATE PUBLICATION cdc_pub FOR TABLE orders, customers;
SELECT * FROM pg_create_logical_replication_slot('cdc_slot', 'pgoutput');
接下来是Java代码。核心思路是获取PGConnection的复制API,打开已存在的复制槽,然后循环调用readPendingData获取ByteBuffer,解析pgoutput消息。为了简单,代码中略去了完整的事务和Relation处理,只演示如何将行级变更转换为字符串并发送到RabbitMQ。
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import org.postgresql.PGConnection;
import org.postgresql.replication.LogSequenceNumber;
import org.postgresql.replication.PGReplicationStream;
import java.nio.ByteBuffer;
import java.sql.DriverManager;
import java.util.Properties;
import java.util.concurrent.TimeUnit;
public class PostgresCDCToRabbitMQ {
public static void main(String[] args) throws Exception {
String jdbcUrl = "jdbc:postgresql://localhost:5432/mydb?replication=database";
Properties props = new Properties();
props.setProperty("user", "postgres");
props.setProperty("password", "postgres");
PGConnection replConn = DriverManager.getConnection(jdbcUrl, props)
.unwrap(PGConnection.class);
PGReplicationStream stream = replConn.getReplicationAPI()
.replicationStream()
.logical()
.withSlotName("cdc_slot")
.withStartPosition(LogSequenceNumber.valueOf("0/0"))
.start();
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
factory.setPort(5672);
factory.setUsername("guest");
factory.setPassword("guest");
try (Connection rmqConn = factory.newConnection();
Channel channel = rmqConn.createChannel()) {
channel.exchangeDeclare("cdc.exchange", "topic", true);
while (true) {
ByteBuffer buffer = stream.readPending();
if (buffer == null) {
TimeUnit.MILLISECONDS.sleep(100);
continue;
}
int offset = buffer.arrayOffset();
byte[] data = buffer.array();
String json = parseMessage(data, offset, buffer.remaining());
if (json != null) {
channel.basicPublish("cdc.exchange", "cdc.orders", null, json.getBytes());
}
stream.setAppliedLSN(stream.getLastReceiveLSN());
stream.setFlushedLSN(stream.getLastReceiveLSN());
}
}
}
private static String parseMessage(byte[] data, int offset, int length) {
// 解析pgoutput二进制消息,此处省略具体实现
return new String(data, offset, length);
}
}
上面的代码只是一个骨架,真正解析pgoutput消息需要处理消息头、事务开始、关系元数据、以及插入更新删除的具体格式。好在社区有一些开源解析库,或者可以直接基于Debezium的解析模块进行二次开发。发布到RabbitMQ时建议使用持久化消息和手动ack,并妥善处理网络异常重连。
四、生产环境的关键考量:顺序、重复与复制槽膨胀
将CDC事件从PostgreSQL推进RabbitMQ,最需要注意的是消息顺序。PostgreSQL逻辑复制流在单个复制槽内是有序的,但如果在消费者内部使用多线程并行处理,顺序就可能被打乱。对于强一致要求的场景,建议单线程消费复制流,将消息按顺序发送到RabbitMQ,让下游按顺序消费。如果允许一定程度乱序,可以在消息中带上LSN或事务ID,由消费者端做排序或去重。
重复消费是另一个常见问题。当消费者崩溃或网络中断后,从上次提交的LSN重新开始读取,会导致部分消息被重复发送。RabbitMQ本身只提供at-least-once投递,无法完全避免重复,因此下游消费者需要实现幂等。可以在消息体中包含事件唯一标识,例如LSN加上事务内序号,下游通过Redis或数据库去重。
复制槽膨胀也不容忽视。如果消费者长时间停止运行,WAL会不断累积,最终撑爆磁盘。设置合理的slot状态监控和告警很有必要。可以通过查询pg_replication_slots视图查看restart_lsn与pg_current_wal_lsn()之间的差距。另外,定期发送心跳消息或使用pg_recvlogical的status_interval参数,也能帮助数据库及时回收无用的WAL。
最后,关于消息格式,建议使用JSON并携带足够的元数据,包括表名、操作类型、时间戳、LSN以及行的前后值。这样RabbitMQ的消费者可以独立处理,不必再回查数据库。在实际部署中,还需要考虑RabbitMQ的交换机类型、队列绑定、死信队列以及连接重试策略,这些细节直接决定了整个数据链路的稳定性。
PostgreSQL CDCRabbitMQ逻辑复制修改时间:2026-09-20 00:54:18