PostgreSQL 实时同步到 Apache Doris:一站式 CDC 解决方案

PostgreSQL 实时同步到 Apache Doris:一站式 CDC 解决方案
作者吴迪SelectDB 内核研发工程师、Apache Doris PMC MemberPostgreSQL 已成为众多业务系统的核心数据库在 各类 TP 场景中被广泛采用。然而PostgreSQL 的设计目标是高并发在线事务处理并不适合承载分析型负载。当执行跨表聚合、大范围扫描等查询时容易占用大量系统资源从而影响线上业务的稳定性。因此在实时分析、实时大屏、数据仓库以及 AI 数据供给等场景中通常会将 PostgreSQL 中的数据实时同步至 Apache Doris 等分析型数据库由 Doris 承担分析查询任务实现 TP 与 AP 负载分离。相比分析系统本身真正具有挑战性的往往是实时同步链路的搭建与维护。如何持续、稳定、准确地将 PostgreSQL 的数据同步至 Doris是整个架构中的关键环节。PostgreSQL 至 Apache Doris 的实时同步先来看传统方案与 Apache Doris 一站式 CDC 同步方案的整体流程。传统同步方案的挑战由上图可知传统方案的同步链路一般由三个部分组成CDC 采集通过 Debezium、Flink CDC 等组件持续捕获 PostgreSQL WAL 日志中的数据变更消息传递根据业务需求将变更数据写入 Kafka或直接发送至下游系统数据处理依赖 Flink、Spark 等计算框架完成数据转换、清洗并最终写入 Apache Doris。这种方案具有较强的灵活性能够支持复杂的数据处理逻辑但同时也意味着需要维护一条由多个独立组件组成的数据链路。对于仅需完成 PostgreSQL 到 Doris 实时同步的场景这样的架构往往存在以下问题组件众多部署复杂除数据库外还需要维护 CDC、Kafka、Flink/Spark 等多个系统涉及部署、监控、升级以及扩缩容等工作。对于仅有数据同步需求的用户而言引入如此多的组件成本较高。同步链路较长故障排查困难当数据延迟、丢失或异常时需要在 CDC、消息队列、计算引擎等多个系统之间逐层排查定位问题耗时且复杂。Apache Doris 一站式 CDC 同步方案针对上述问题Apache Doris 提供了内置 PostgreSQL CDC 同步能力是专门面向 TP 数据库的实时同步场景。该方案将数据采集、任务调度、数据写入、断点续传等能力统一集成到 Doris Job 调度框架中无需部署额外的 CDC 工具、消息队列或计算引擎。因此整体同步链路被进一步简化为PostgreSQL → Apache Doris。对于用户而言需要维护的对象也从多个独立组件收敛为一个 Doris 集群降低了部署和运维复杂度。创建同步任务也非常简单只需执行一条 SQLCREATE JOB test_postgres_job ON STREAMING FROM POSTGRES ( jdbc_url jdbc:postgresql://127.0.0.1:5432/postgres, driver_url postgresql-42.5.0.jar, driver_class org.postgresql.Driver, user postgres, password postgres, database postgres, schema cdc_test, include_tables student1,student2, offset initial ) TO DATABASE target_db任务创建完成后Doris 将自动完成以下工作根据 PostgreSQL 表结构自动创建目标表如不存在导入历史存量数据持续消费 PostgreSQL WAL实时同步增量数据。一站式 CDC 同步方案的实现原理内置 PostgreSQL CDC 并不是一个独立运行的同步组件而是构建在 Apache Doris 的 Job 调度框架之上。通过FE调度、BE 执行以及 CDC Client 拉取变更数据三部分协同工作实现持续运行的实时同步任务。接下来我们了解下具体的执行流程。用户在 FE 中通过 SQL 创建 PostgreSQL Streaming Job并指定同步库、同步表及相关参数FE 根据 PostgreSQL 表结构自动创建 Doris 目标表若不存在Job Scheduler 生成同步 Task并调度至某个 BE 节点执行BE 接收到任务后将请求转发给本机 CDC Client若 CDC Client 尚未启动则自动拉起对应进程CDC Client 持续读取 PostgreSQL WAL完成全量及增量数据采集数据序列化后通过 Stream Load 写入 Doris当前批次写入完成后BE 将同步 Offset 上报给 FEFE 持久化 Offset并继续调度下一轮同步任务实现持续同步。整个过程中数据同步、任务调度以及状态管理均由 Doris 自身完成无需依赖 Kafka、Flink 或 Spark 等外部系统。这种设计一方面缩短了同步链路另一方面也充分复用了 Doris 已有的 Job 调度能力使整个同步过程更加稳定、易于维护。快速体验创建 PostgreSQL 同步任务仅需执行一条 SQL下面示例将 PostgreSQL 中的数据同步至目标数据库target_test_dbCREATE JOB test_postgres_job ON STREAMING FROM POSTGRES ( jdbc_url jdbc:postgresql://127.0.0.1:5432/postgres, driver_url postgresql-42.5.0.jar, driver_class org.postgresql.Driver, user postgres, password postgres, database postgres, schema cdc_test, include_tables student1,student2, offset initial ) TO DATABASE target_test_db创建完成后Job 会自动完成目标表创建、全量数据导入以及后续增量同步。同步过程中可以通过如下命令查看任务状态mysql select * from jobs(typeinsert) where nametest_postgres_job\G *************************** 1. row *************************** Id: 1774599451260 Name: test_postgres_job Definer: root ExecuteType: STREAMING RecurringStrategy: Status: RUNNING ExecuteSql: FROMPOSTGRES(schemacdc_test,include_tablesstudent1,student2,databasepostgres,driver_classorg.postgresql.Driver,offsetinitial,driver_urlpostgresql-42.5.0.jar,jdbc_urljdbc:postgresql://127.0.0.1:5432/postgres,userpostgres) TO DATABASE target_test_db (table.create.properties.replication_num1) CreateTime: 2026-03-27 17:02:48 SucceedTaskCount: 3 FailedTaskCount: 0 CanceledTaskCount: 0 Comment: Properties: CurrentOffset: {lsn:53547640,txId:1523,splitId:binlog-split} EndOffset: {lsn:53547640} LoadStatistic: {scannedRows:7,loadBytes:341,fileNumber:0,fileSize:0,filteredRows:0} ErrorMsg: JobRuntimeMsg: No data available for consumption at the moment, will retry after 1774602192496其中几个关键字段含义如下CurrentOffset记录当前同步进度、LoadStatistic反映已写入的行数与字节数。通过这些指标可以快速判断同步任务是否正常运行以及当前数据同步的实时进度。更多参数配置和高级用法可参考 PostgreSQL CDC 自动建表同步 - Apache Doris。总结对于 PostgreSQL 到 Apache Doris 的实时同步而言传统方案通常需要引入 CDC、消息队列和计算引擎等多个组件虽然具备较高的灵活性但也增加了部署和运维成本。Apache Doris 将 PostgreSQL CDC、Job 调度、自动建表、数据导入以及断点续传等能力统一集成到 Job 调度框架中使实时同步能够简化为一条 SQL 即可完成。后续该能力也将陆续集成至 SelectDB进一步降低生产环境下 PostgreSQL 实时同步的部署与运维成本。对于以数据同步为主要需求、无需复杂实时计算的场景这种内置方案能够有效缩短同步链路降低系统复杂度并减少后续的运维成本。