當您需要對PolarDB MySQL版中的業務資料進行全文檢索索引或複雜分析時,直接在資料庫上操作可能會影響核心業務的穩定性。PolarDB提供的AutoETL功能,能將資料從讀寫節點自動、持續地同步至叢集內的PolarSearch節點,為您提供一站式的資料服務。您可以通過搜尋視圖(Search View)或ETL預存程序快速建立資料同步鏈路,無需額外部署和維護ETL工具,即可實現資料同步,並將搜尋分析負載與線上交易處理負載隔離。
使用AutoETL建立鏈路,預設授權允許AutoETL引擎訪問PolarDB資料進行資料同步。
功能簡介
AutoETL是PolarDB MySQL版內建的資料同步能力,它允許資料在叢集內不同類型的節點間自動流轉。目前的版本僅支援從PolarDB MySQL版同步至同一叢集內的PolarSearch節點,以用於高效能的搜尋和分析。
AutoETL提供兩種建立資料同步鏈路的方式:
搜尋視圖(Search View):通過
CREATE SEARCH VIEW文法,以標準SQL的方式定義資料同步邏輯。適合大多數單表同步和多表匯聚情境,系統自動處理底層串連細節。ETL預存程序(
dbms_etl.sync_by_sql):使用相容Flink SQL的文法,通過預存程序定義複雜的資料清洗、轉換和彙總邏輯。
適用範圍
使用AutoETL功能前,需確保環境滿足以下條件:
叢集版本:
搜尋視圖(Search View):
MySQL 8.0.1,且修訂版本需為8.0.1.1.54或以上。
MySQL 8.0.2,且修訂版本需為8.0.2.2.34或以上。
ETL預存程序(sync_by_sql):
MySQL 8.0.1,且修訂版本需為8.0.1.1.52或以上。
MySQL 8.0.2,且修訂版本需為8.0.2.2.33或以上。
Binlog:叢集需要開啟Binlog。
同步方向:僅支援從PolarDB MySQL版同步至同一叢集內的PolarSearch節點。
DDL限制:對已建立搜尋視圖/ETL預存程序的源表進行DDL操作時,需遵循特定的規則以避免同步中斷。部分不相容的變更需要重建搜尋視圖。詳情請參見DDL變更規則與實踐。
資料類型:暫不支援
BIT類型以及GEOMETRY、POINT、LINESTRING、POLYGON、MULTIPOINT、MULTILINESTRING、MULTIPOLYGON、GEOMETRYCOLLECTION等空間資料類型的同步。搜尋視圖查詢限制:搜尋視圖目前僅支援定義同步語義,不支援資料查詢。資料查詢請直接連接PolarSearch節點執行。
計算資源:AutoETL 使用 CU 作為計算單元,預設叢集的 CU 數為 PolarSearch節點所有節點 CPU 之和的兩倍。您可以在控制台頁面的AutoETL頁簽中查看當前叢集所使用的 CU 數。
搜尋視圖
搜尋視圖(Search View)是AutoETL提供的一種聲明式資料同步機制。您可以使用標準SQL文法建立搜尋視圖,系統將自動建立從源表到PolarSearch節點的持續資料同步鏈路。
建立搜尋視圖
文法
CREATE SEARCH VIEW view_name [(column_list, PRIMARY KEY (pk_column_list))]
[WITH (option_list)]
AS select_statement;參數說明
參數 | 必填 | 說明 |
| 是 | 搜尋視圖名稱,同時也是PolarSearch節點中目標索引的名稱。 |
| 否 | 手動定義搜尋視圖的列,多列使用 說明 目前僅支援單表同步無需指定 |
| 否 | 指定搜尋視圖的主鍵列。搜尋視圖的結構對應於PolarSearch節點的Index mapping, 如果不指定,預設會自動使用 |
| 否 | 建立搜尋視圖時顯式指定的同步配置,例如同步的並行度、單個同步 worker 的計算資源用量、目標索引名等,多個配置之間使用 |
| 是 | 定義資料來源和同步邏輯的 |
使用限制與說明
源表需包含主鍵或唯一鍵。
需具有搜尋視圖中所有源表的
ALTER許可權,以及相關列(或整表)的SELECT許可權。建立搜尋視圖後,源表新增的列預設不會被自動同步。如需同步新增列,請參見變更搜尋視圖。
如果您希望使用自訂的目標索引配置,可以先在PolarSearch節點中手動建立索引並定義其配置,然後再建立搜尋視圖。如果建立時目標索引不存在,系統將自動建立。
如需為多表匯聚或複雜查詢配置進階同步參數(如JSON欄位轉換、路由欄位等),請參見AutoETL 參數配置和實踐案例。
資料準備
以下樣本使用的測試資料,您可以在PolarDB MySQL版中執行以下SQL語句建立。
CREATE DATABASE IF NOT EXISTS db1;
CREATE DATABASE IF NOT EXISTS db2;
USE db1;
CREATE TABLE IF NOT EXISTS t1 (
id INT PRIMARY KEY,
c1 VARCHAR(100),
c2 VARCHAR(100)
);
INSERT INTO t1(id, c1, c2) VALUES
(1, 'apple', 'red'),
(2, 'banana', 'yellow'),
(3, 'grape', 'purple');
USE db2;
CREATE TABLE IF NOT EXISTS t2 (id INT PRIMARY KEY, c2 INT);
INSERT INTO t2(id, c2) VALUES (1, 111), (2, 222), (4, 444);樣本
全表同步:將
db1.t1全部資料同步到PolarSearch。視圖名view_test即為PolarSearch中的目標索引名。CREATE SEARCH VIEW view_test AS SELECT * FROM db1.t1;指定列同步:僅同步
c1和c2列,並手動定義列類型和主鍵。CREATE SEARCH VIEW view_test1 AS SELECT c1, c2 FROM db1.t1;指定同步參數:建立搜尋視圖時通過
WITH子句指定同步配置。以下樣本將db1.t1的c1和c2列同步到PolarSearch,並將同步的並行度設定為 2。CREATE SEARCH VIEW view_test6 WITH ('parallelism' = '2') AS SELECT c1, c2 FROM db1.t1;條件過濾同步:僅同步滿足
WHERE條件的資料。CREATE SEARCH VIEW view_test2 AS SELECT id, c1, c2 FROM db1.t1 WHERE c1 > 10;多表JOIN:將
db1.t1和db2.t2通過id欄位進行JOIN,將結果同步到PolarSearch。CREATE SEARCH VIEW view_test3(id, c1, c2) AS SELECT t1.id, t1.c1, t2.c2 FROM db1.t1 AS t1 LEFT JOIN db2.t2 AS t2 ON t1.id = t2.id;多表UNION:將多個結構相同的表合并後同步。要求各
SELECT語句的列數和類型一致。CREATE SEARCH VIEW view_test4(id, c2) AS SELECT id, c2 FROM db1.t1 UNION ALL SELECT id, c2 FROM db2.t2;分組彙總:對資料進行分組彙總後同步。使用
GROUP BY時需要手動定義列和主鍵。CREATE SEARCH VIEW view_test5 (id, max_c) AS SELECT t1.id, MAX(t1.c1) AS max_c FROM db1.t1 GROUP BY t1.id;
驗證資料
查看搜尋視圖的同步狀態:
SHOW SEARCH VIEW STATUS;當狀態為active時,表示搜尋視圖資料同步正常。串連到PolarSearch節點,使用與Elasticsearch相容的REST API驗證資料:
# 將<user>:<password>替換為polarsearch節點的帳號,<polarsearch_endpoint>替換為PolarSearch節點的串連地址與連接埠
curl -u <user>:<password> -X GET "http://<polarsearch_endpoint>/view_test/_search"管理搜尋視圖
您可以使用以下命令查看已建立的搜尋視圖,或停止、重啟、重建搜尋視圖的資料同步。以下命令均需串連叢集後在資料庫用戶端中執行。
查看所有搜尋視圖的狀態
SHOW SEARCH VIEW STATUS;返回結果如下:當狀態為active時,表示搜尋視圖資料同步正常。
+------------+--------+----------+---------+---------------------+---------------------+
| View Name | Type | Status | Message | Created_at | Updated_at |
+------------+--------+----------+---------+---------------------+---------------------+
| view_test | search | active | | 2026-03-18 18:44:12 | 2026-03-18 18:51:37 |
+------------+--------+----------+---------+---------------------+---------------------+查看指定搜尋視圖的建立語句
SHOW CREATE SEARCH VIEW view_test;返回結果如下:
+-----------+------------------------------------------------------+
| View Name | Create Search View |
+-----------+------------------------------------------------------+
| view_test | CREATE SEARCH VIEW view_test AS SELECT * FROM db1.t1 |
+-----------+------------------------------------------------------+停止搜尋視圖
需要修改或重建PolarSearch節點上的目標索引時,為避免同步寫入報錯,您可以先停止搜尋視圖的資料同步,待PolarSearch節點上的索引修改完成後再重啟搜尋視圖。
ALTER SEARCH VIEW view_test STOP;重啟搜尋視圖
重啟已停止或正在啟動並執行搜尋視圖。
ALTER SEARCH VIEW view_test RESTART;重建搜尋視圖
重新全量讀取源表資料並寫入PolarSearch節點。
重建會重新掃描源表的全量資料,資料量較大時耗時較長。
重建不會清理PolarSearch節點上已有的索引資料,而是直接覆蓋。
ALTER SEARCH VIEW view_test REBUILD;刪除搜尋視圖
刪除搜尋視圖是高危操作,執行前請務必確認。此操作用於停止搜尋視圖的資料同步並清理相關資源,但不會刪除PolarSearch的索引資料。
DROP SEARCH VIEW view_name;對不同狀態的搜尋視圖執行刪除時,系統的處理邏輯存在差異:
active狀態的搜尋視圖:首先會變為dropping,待系統完成資源清理和目標索引資料的刪除後,狀態才會變為dropped。dropped狀態的搜尋視圖:系統將徹底清除該搜尋視圖的資訊。其他狀態的搜尋視圖:系統不支援刪除操作。
變更搜尋視圖
搜尋視圖建立後,您可以調整它的運行參數,也可以變更它的同步邏輯(如新增同步欄位、修改查詢條件等)。AutoETL 提供三種變更方式:
只調整運行參數(如並行度、計算資源用量)時,使用參數變更,無需改動同步 SQL。
變更同步邏輯時,優先使用原地 SQL 變更。變更前後的 SQL 定義相容時,同步從原有位點繼續,無需重新全量同步。您可以根據DDL變更規則與實踐判斷變更前後的 SQL 定義是否相容。
變更前後的 SQL 定義不相容、無法原地變更時,採用“新索引 + 新搜尋視圖”變更進行重建。
參數變更
對於運行中的搜尋視圖,您可以通過以下文法修改它的運行參數。AutoETL 會自動讀取新配置並重啟該搜尋視圖。
文法
ALTER SEARCH VIEW view_name UPDATE WITH (new_option_list);樣本
將view_test的並行度調整為 8,單個 worker 的 CPU 為 4,每個 worker 支援 8 個並發:
ALTER SEARCH VIEW view_test UPDATE WITH ('parallelism' = '8', 'link.tm.cpu' = '4', 'link.tm.slots' = '8');原地 SQL 變更
原地 SQL 變更直接修改已有搜尋視圖的同步 SQL,無需建立PolarSearch索引和搜尋視圖並切換業務查詢,可以減少全量同步的耗時。變更前後的 SQL 定義相容時,同步從原有位點繼續;不相容時,本次變更失敗,系統自動回退到變更前的同步定義並重啟搜尋視圖。
文法
ALTER SEARCH VIEW view_name UPDATE
[TO (column_list, PRIMARY KEY (pk_column_list))]
AS new_select_statement;樣本
同步源表新增的列:全表同步的搜尋視圖不會自動同步源表在視圖建立後新增的列。源表
db1.t1新增列後,執行以下語句將新列的資料同步到PolarSearch。ALTER SEARCH VIEW view_test UPDATE;調整同步的列:原來同步
c1和c2列,變更為僅同步c1列。ALTER SEARCH VIEW view_test UPDATE AS SELECT c1 FROM db1.t1;調整過濾條件:將
WHERE條件由c1 > 10變更為c1 > 20。ALTER SEARCH VIEW view_test UPDATE AS SELECT id, c1, c2 FROM db1.t1 WHERE c1 > 20;
“新索引 + 新搜尋視圖”變更
變更前後的 SQL 定義不相容、無法原地變更時,採用“新索引 + 新搜尋視圖”的方式進行重建,以確保業務查詢不受影響。
建立一個新的搜尋視圖,同步至新的PolarSearch索引。
通過
SHOW SEARCH VIEW STATUS查看新搜尋視圖的狀態,當等待新搜尋視圖的同步時延降至0~1秒時,將業務查詢邏輯從舊索引切換到新索引。刪除舊的搜尋視圖。
關於源表DDL變更對搜尋視圖的影響和詳細的變更實踐,請參見DDL變更規則與實踐。
ETL預存程序(sync_by_sql)
對於需要複雜轉換、彙總或計算的情境,您可以使用CALL dbms_etl.sync_by_sql預存程序,通過相容Flink SQL的文法定義資料同步邏輯。
建立同步鏈路
文法
調用dbms_etl.sync_by_sql前,您可以通過AutoETL參數配置和實踐案例中的會話變數(如esl_link_options、esl_sink_options)設定鏈路的同步配置,AutoETL引擎在建立鏈路時會自動讀取這些變數。
CALL dbms_etl.sync_by_sql("search", "<sync_sql>");樣本
源表和目標表的串連資訊(如主機地址、連接埠、帳號密碼等)由系統自動設定,您無需在WITH子句中手動指定。
CALL dbms_etl.sync_by_sql("search", "
-- 步驟1:定義 PolarDB 源表
CREATE TEMPORARY TABLE `db1`.`t1` (
`id` BIGINT,
`c1` STRING,
PRIMARY KEY (`id`) NOT ENFORCED
) WITH (
'connector' = 'mysql',
'database-name' = 'db1',
'table-name' = 't1'
);
-- 步驟2:定義 PolarSearch 目標表
CREATE TEMPORARY TABLE `dest` (
`id` BIGINT,
`max_c` STRING,
PRIMARY KEY (`id`) NOT ENFORCED
) WITH (
'connector' = 'opensearch',
'index' = 'dest'
);
-- 步驟3:定義計算和插入邏輯
INSERT INTO `dest`
SELECT
`t1`.`id`,
MAX(`t1`.`c1`)
FROM `db1`.`t1` AS `t1`
GROUP BY `t1`.`id`;
");驗證資料
串連到PolarSearch節點,使用與Elasticsearch相容的REST API進行查詢,確認資料已同步。
# 將<polarsearch_endpoint>替換為PolarSearch節點的串連地址
curl -u <user>:<password> -X GET "http://<polarsearch_endpoint>/dest/_search"管理同步鏈路
您可以使用以下命令查看已建立的同步鏈路,或停止、重啟、重建同步鏈路的資料同步。以下命令均需串連叢集後在資料庫用戶端中執行。
查看所有鏈路
CALL dbms_etl.show_sync_link();根據ID查看指定鏈路
將<sync_id>替換為建立鏈路時返回的ID。
CALL dbms_etl.show_sync_link_by_id('<sync_id>')\G返回結果說明:
*************************** 1. row ***************************
SYNC_ID: crb5rmv8rttsg
NAME: crb5rmv8rttsg
SYSTEM: search
SYNC_DEFINITION: db1.t1 -> dest
SOURCE_TABLES: db1.t1
SINK_TABLES: dest
STATUS: active -- 鏈路狀態,active表示正常運行
MESSAGE: -- 如果出錯,此處會顯示錯誤資訊
CREATED_AT: 2024-05-20 11:55:06
UPDATED_AT: 2024-05-20 17:28:04
OPTIONS: ...停止鏈路
需要修改或重建PolarSearch節點上的目標索引時,為避免同步寫入報錯,您可以先停止同步鏈路,待PolarSearch節點上的索引修改完成後再重啟鏈路。
CALL dbms_etl.stop_sync_link('<sync_id>');重啟鏈路
重啟已停止或正在啟動並執行同步鏈路。
CALL dbms_etl.restart_sync_link('<sync_id>');重建鏈路
重新全量讀取源表資料並寫入PolarSearch節點。重建會重新掃描源表的全量資料,資料量較大時耗時較長;重建不會清理PolarSearch節點上已有的索引資料,而是直接覆蓋。
CALL dbms_etl.rebuild_sync_link('<sync_id>');刪除同步鏈路
此操作用於停止資料同步並清理相關資源。
刪除同步鏈路是高危操作,執行前請務必確認。此操作用於停止同步鏈路的資料同步並清理相關資源,但不會刪除PolarSearch的索引資料。
CALL dbms_etl.drop_sync_link('<sync_id>');對不同狀態的鏈路執行drop_sync_link刪除時,系統的處理邏輯存在差異:
active狀態的鏈路:首先會變為dropping,待系統完成鏈路資源和目標索引資料的清理後,狀態才會變為dropped。dropped狀態的鏈路:系統將徹底清除該鏈路的資訊。其他狀態的鏈路:系統不支援刪除操作。
變更同步鏈路
同步鏈路建立後,您可以調整它的運行參數,也可以變更它的同步 SQL。與搜尋視圖一致,AutoETL 提供三種變更方式:
只調整運行參數(如並行度、計算資源用量)時,使用參數變更,無需改動同步 SQL。
變更同步 SQL 時,優先使用原地 SQL 變更。變更前後的 SQL 定義相容時,同步從原有位點繼續,無需重新全量同步。您可以根據DDL變更規則與實踐判斷變更前後的 SQL 定義是否相容。
變更前後的 SQL 定義不相容、無法原地變更時,採用“新索引 + 新鏈路”變更進行重建。
參數變更
對於運行中的同步鏈路,您可以先通過會話變數esl_link_options設定新的鏈路配置,再調用dbms_etl.update_sync_link使配置生效。AutoETL 會自動讀取新的esl_link_options配置並重啟該同步鏈路。
文法
SET esl_link_options = "<new_option_list>";
CALL dbms_etl.update_sync_link('<sync_id>', '');樣本
將鏈路8f4228x2uq12z的並行度調整為 8,單個 worker 的 CPU 為 4,每個 worker 支援 8 個並發:
SET esl_link_options = "'parallelism' = '8', 'link.tm.cpu' = '4', 'link.tm.slots' = '8'";
CALL dbms_etl.update_sync_link('8f4228x2uq12z', '');原地 SQL 變更
原地 SQL 變更直接修改已有鏈路的同步 SQL,無需建立PolarSearch索引和鏈路並切換業務查詢,可以減少全量同步的耗時。變更前後的 SQL 定義相容時,同步從原有位點繼續;不相容時,本次變更失敗,系統自動回退到變更前的同步定義並重啟鏈路。
文法
<new_sync_sql>需傳入完整的新同步 SQL,其結構與建立鏈路時dbms_etl.sync_by_sql的同步 SQL 一致。
CALL dbms_etl.update_sync_link('<sync_id>', '<new_sync_sql>');樣本
將鏈路8f4228x2uq12z的過濾條件變更為c1 > 20:
CALL dbms_etl.update_sync_link('8f4228x2uq12z', "
CREATE TEMPORARY TABLE `db1`.`t1` (
`id` BIGINT,
`c1` STRING,
PRIMARY KEY (`id`) NOT ENFORCED
) WITH (
'connector' = 'mysql',
'database-name' = 'db1',
'table-name' = 't1'
);
CREATE TEMPORARY TABLE `dest` (
`id` BIGINT,
`c1` STRING,
PRIMARY KEY (`id`) NOT ENFORCED
) WITH (
'connector' = 'opensearch',
'index' = 'dest'
);
INSERT INTO `dest` SELECT `id`, `c1` FROM `db1`.`t1` WHERE `c1` > 20;
");“新索引 + 新鏈路”變更
變更前後的 SQL 定義不相容、無法原地變更時,採用“新索引 + 新鏈路”的方式進行重建,以確保業務查詢不受影響。以變更鏈路 A 為例:
根據新的同步 SQL 建立鏈路 B,同步至新的PolarSearch索引。
通過
CALL dbms_etl.show_sync_link_by_id('<sync_id>')查看鏈路 B 的狀態,當等待鏈路 B 的同步時延降至0~1秒時,將業務查詢邏輯從舊索引切換到新索引。確認鏈路 B 運行穩定後,執行
CALL dbms_etl.drop_sync_link('<sync_id>')刪除舊鏈路 A。