相容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版本為例,為您介紹具體的打包步驟:
-
複製對應版本的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 -
複製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 -
應用適配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中移除。
-
-
使用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訊息佇列中。
前提準備
-
Kafka準備
-
部署Kafka執行個體,並確保在Kafka Connect的主機上能夠成功訪問。您也可以直接使用雲訊息佇列 Kafka 版,詳情請參考快速入門。
-
在Kafka執行個體中建立一個名為pg_dbz_event的Topic,用於接收訊息。
說明測試情境為便於查看可以建立單分區Topic,對於生產環境請建立多分區Topic。
-
-
在本地以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
-
-
PolarDB PostgreSQL版(相容Oracle)準備
-
在PolarDB叢集購買頁面,購買PolarDB PostgreSQL版(相容Oracle) 2.0叢集。
-
按照使用說明,完成PolarDB叢集配置,滿足Debezium PolarDBO connector使用前提。
-
建立高許可權賬戶,詳細操作請參考建立帳號。
-
擷取叢集主地址,詳細操作請參考查看串連地址。如果PolarDB叢集和Kafka Connect在同一可用性區域,可直接使用私網地址,否則需要申請公網地址。將Kafka Connect執行個體地址添加到PolarDB叢集白名單中,請參考設定叢集白名單。
-
在控制台建立資料庫dbz_db,詳細步驟請參考建立資料庫。
-
執行如下語句,在資料庫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());
-
測試
-
建立設定檔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做了彙總。
-
添加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。 -
在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。