全部產品
Search
文件中心

PolarDB:相容PolarDB PostgreSQL版(相容Oracle)的Debezium connector

更新時間:Aug 27, 2026

相容PolarDB PostgreSQL版(相容Oracle)的Debezium connector(簡稱Debezium PolarDBO connector),可用於捕獲PolarDB PostgreSQL版(相容Oracle)資料庫中的行層級更改,產生資料變更事件記錄,並將它們串流到Kafka Topic中。具體功能及用法請參考社區Debezium PostgreSQL connector。

由於PolarDB PostgreSQL版(相容Oracle)與社區PostgreSQL僅在少量資料類型和內建對象處理存在差異,本文為您介紹如何基於社區Debezium PostgreSQL connector,通過少量代碼適配打包出支援PolarDB PostgreSQL版(相容Oracle)的Debezium connector。

打包Debezium PolarDBO connector

重要

Debezium PolarDBO connector基於社區Debezium PostgreSQL connector適配開發,無論是您自行打包,還是使用本文中提供的JAR包,Debezium PolarDBO connector都不提供SLA保障。

操作前提

  • 配置Java環境

    目前Debezium各版本均要求Java11及以上版本,請在打包和正式運行時提前配置Java11環境。

  • 確定Debezium版本

    根據您使用的Kafka/Kafka Connect和PolarDB PostgreSQL版(相容Oracle)版本,確定Debezium版本。具體的版本相容資訊,請參考Debezium發布概覽。

    說明
    • Debezium代碼倉庫請參考Debezium。

    • 對於PolarDB PostgreSQL版(相容Oracle)匹配的社區版本如下:

      • Oracle文法相容 2.0對應社區PostgreSQL 14。

      • Oracle文法相容 1.0對應社區PostgreSQL 11。

  • 確定PgJDBC版本

    在對應版本的Debezium的pom.xml中通過尋找關鍵字version.postgresql.driver確定PgJDBC版本。

    說明

    PgJDBC代碼倉庫請參考PgJDBC。

操作步驟

社區Debezium 2.6.2.Final支援Kafka Connect 2.x、3.x版本,支援PostgreSQL 10、11、12、13、14、15、16版本。

接下來以社區Debezium 2.6.2.Final版本為例,為您介紹具體的打包步驟:

  1. 複製對應版本的Debezium和PgJDBC的代碼檔案。

    git clone -b v2.6.2.Final --depth=1 https://github.com/debezium/debezium.git
    git clone -b REL42.6.1 --depth=1 https://github.com/pgjdbc/pgjdbc.git
  2. 複製PgJDBC部分檔案到Debezium中。

    mkdir -p debezium/debezium-connector-postgres/src/main/java/org/postgresql/core/v3       
    mkdir -p debezium/debezium-connector-postgres/src/main/java/org/postgresql/jdbc
    cp pgjdbc/pgjdbc/src/main/java/org/postgresql/core/v3/ConnectionFactoryImpl.java debezium/debezium-connector-postgres/src/main/java/org/postgresql/core/v3 
    cp pgjdbc/pgjdbc/src/main/java/org/postgresql/core/v3/QueryExecutorImpl.java debezium/debezium-connector-postgres/src/main/java/org/postgresql/core/v3 
    cp pgjdbc/pgjdbc/src/main/java/org/postgresql/jdbc/PgDatabaseMetaData.java debezium/debezium-connector-postgres/src/main/java/org/postgresql/jdbc
  3. 應用適配PolarDB PostgreSQL版(相容Oracle)的patch檔案。

    git apply v2.6.2.Final-support-polardbo-v1.patch
    說明
    • 以上使用的Debezium PolarDBO connector相容patch檔案:v2.6.2.Final-support-polardbo-v1.patch。

    • 該patch文檔預設將依賴debezium-api、debezium-core、PgJDBC和protobuf-java打包至JAR包,如不需要可以從pom.xml中移除。

  4. 使用Maven打包Debezium PolarDBO connector。

    mvn clean package -pl :debezium-connector-postgres -DskipITs -Dquick
    # 打包完成後可以在debezium-connector-postgres/的target目錄中擷取到jar包

    按照以上流程基於JDK11打包出Debezium PolarDBO connector的JAR包:debezium-connector-postgres-polardbo-v1.0-2.6.2.Final.jar。

使用說明

Debezium PolarDBO connector是通過PolarDB PostgreSQL版(相容Oracle)資料庫的邏輯複製讀取增量變更,使用時需要滿足以下條件:

  • wal_level參數的值需設定為logical,即在預寫式日誌WAL(Write-ahead logging)中增加支援邏輯複製所需的資訊。

    說明

    您可以通過控制台設定wal_level參數,詳細操作請參考設定叢集參數。修改該參數後叢集將會重啟,請在修改參數前做好業務安排,謹慎操作。

  • 執行ALTER TABLE schema.table REPLICA IDENTITY FULL;命令設定訂閱表的REPLICA IDENTITY為FULL(發出的插入和更新操作事件包含表中所有列的舊值),以保障該表資料同步的一致性。

    說明
    • REPLICA IDENTITY是PostgreSQL特有的表級設定,決定了邏輯解碼外掛程式在發生(INSERT)和更新(UPDATE)事件時,是否包含涉及的表列的舊值。REPLICA IDENTITY取值含義詳情,請參見REPLICA IDENTITY。

    • 設定訂閱表的REPLICA IDENTITY為FULL時可能需要鎖表,進而影響業務,請在修改參數前做好業務安排。您可以通過以下命令查看當前配置是否為FULL:

      SELECT relreplident = 'f' FROM pg_class WHERE relname = 'tablename';
  • 需要確保max_wal_senders和max_replication_slots的參數值均大於當前資料庫複寫槽已使用數和Kafka作業所需要的slot數量。

  • 確保使用的是高許可權帳號或者同時擁有LOGIN和REPLICATION許可權的普通帳號,並且具有訂閱表的SELECT許可權用於全量資料查詢。

  • 只能串連PolarDB叢集的主地址,叢集地址不支援邏輯複製。

  • connector.class參數指定為io.debezium.connector.postgresql.PolarDBOConnector。

  • 建議將plugin.name參數設定為pgoutput,否則非UTF-8編碼的資料庫可能會發生增量解析亂碼,詳細介紹請參考社區文檔。

樣本

以下樣本用於說明,如何通過Debezium PolarDBO connector,將PolarDB PostgreSQL版Oracle文法相容 2.0叢集中dbz_db庫的t1和t2表,同步到Kafka訊息佇列中。

前提準備

  1. Kafka準備

    1. 部署Kafka執行個體,並確保在Kafka Connect的主機上能夠成功訪問。您也可以直接使用雲訊息佇列 Kafka 版,詳情請參考快速入門。

    2. 在Kafka執行個體中建立一個名為pg_dbz_event的Topic,用於接收訊息。

      說明

      測試情境為便於查看可以建立單分區Topic,對於生產環境請建立多分區Topic。

  2. 在本地以distributed模式啟動Kafka Connect,連接埠為8083。

    • 將上文打包的Debezium PolarDBO connectorJAR包拷貝到Kafka Connect的plugin.path目錄中。

      # ${plugin.path} 請替換為具體的路徑
      mkdir ${plugin.path}/debezium-connector-polardbo
      cp debezium-connector-postgres-polardbo-v1.0-2.6.2.Final.jar ${plugin.path}/debezium-connector-polardbo
  3. PolarDB PostgreSQL版(相容Oracle)準備

    1. 在PolarDB叢集購買頁面,購買PolarDB PostgreSQL版(相容Oracle) 2.0叢集。

    2. 按照使用說明,完成PolarDB叢集配置,滿足Debezium PolarDBO connector使用前提。

    3. 建立高許可權賬戶,詳細操作請參考建立帳號。

    4. 擷取叢集主地址,詳細操作請參考查看串連地址。如果PolarDB叢集和Kafka Connect在同一可用性區域,可直接使用私網地址,否則需要申請公網地址。將Kafka Connect執行個體地址添加到PolarDB叢集白名單中,請參考設定叢集白名單。

    5. 在控制台建立資料庫dbz_db,詳細步驟請參考建立資料庫。

    6. 執行如下語句,在資料庫dbz_db中建立表t1、t2,並寫入資料。

      CREATE TABLE public.t1 (a int PRIMARY KEY, b text, c TIMESTAMP);
      ALTER TABLE public.t1 REPLICA IDENTITY FULL;
      INSERT INTO public.t1(a, b, c) VALUES(1, 'a', now());
      CREATE TABLE public.t2 (a int PRIMARY KEY, b text, c DATE);
      ALTER TABLE public.t2 REPLICA IDENTITY FULL;
      INSERT INTO public.t2(a, b, c) VALUES(1, 'a', now());

測試

  1. 建立設定檔config/postgresql-connector.json,配置說明請參考社區文檔。

    {
      "name": "dbz-polardb",
      "config": {
        "connector.class": "io.debezium.connector.postgresql.PolarDBOConnector",
        "database.hostname": "<yourHostname>", 
        "database.port": "<yourPort>", 
        "database.user": "<yourUserName>", 
        "database.password": "<yourPassWord>", 
        "database.dbname" : "dbz_db",
        "plugin.name": "pgoutput",
        "slot.name": "dbz_polardb",
        "table.include.list": "public.t1,public.t2",
        "topic.prefix": "polardb"
        "transforms": "Combine",
        "transforms.Combine.type": "io.debezium.transforms.ByLogicalTableRouter",
        "transforms.Combine.topic.regex": "(.*)",
        "transforms.Combine.topic.replacement": "pg_dbz_event"
      }
    }
    說明

    預設需要為每個表建立一個Topic,以上配置對Topic做了彙總。

  2. 添加connector。

    curl -i -X POST -H "Accept:application/json" -H "Content-Type:application/json" 'http://localhost:8083/connectors' -d @config/postgresql-connector.json

    成功添加後,能夠在Kafka的Topic中查詢到全量資料。

    在 Kafka 執行個體的訊息查詢頁簽,選取查詢方式為按位點查詢,分區選擇 0,起始位點設為 0,單擊查詢。結果顯示兩條 Debezium CDC 訊息(位點 0 和 1),Key 中 __dbz__physicalTableIdentifier 分別為 polardb.public.t1 和 polardb.public.t2,確認全量資料已同步至 Kafka。

  3. 在PolarDB叢集的dbz_db庫中執行以下DML語句:

    INSERT INTO public.t1(a, b, c) VALUES(2, 'b', now());
    UPDATE public.t1 SET b = 'c' WHERE a = 1;
    DELETE FROM public.t1 WHERE a = 2;
    INSERT INTO public.t1(a, b, c) VALUES(4, 'd', now());
    INSERT INTO public.t2(a, b, c) VALUES(2, 'b', now());
    UPDATE public.t2 SET b = 'c' WHERE a = 1;
    DELETE FROM public.t2 WHERE a = 2;
    INSERT INTO public.t2(a, b, c) VALUES(4, 'd', now());

    能夠在Kafka的Topic中查詢到增量資料。

    在訊息佇列的訊息查詢頁面,將查詢方式設定為按位點查詢,選擇目標資料分割和起始位點後單擊查詢。查詢結果顯示 t1 和 t2 表的 DML 操作已被 Debezium 採集為 CDC 訊息,每條訊息的 Key 包含表標識符(如 polardb.public.t1)及主索引值,Value 包含 Debezium 格式的變更事件(含 before/after 欄位),其中 DELETE 操作對應的 Value 為 0 Bytes。