AutoETL支持汇聚多个PolarDB MySQL版集群的业务数据到同一个PolarSearch节点中,便于您对分布在不同集群中的数据进行统一检索和分析。本文介绍如何通过AutoETL创建多集群数据汇聚链路。
当前功能处于灰度阶段。如您有相关需求,请提交工单与我们联系,以便为您开启该功能。
背景信息
在实际业务场景中,数据可能分布在多个PolarDB MySQL版集群中。当您需要对这些分散的数据进行统一的全文检索或分析时,可以通过AutoETL的多集群数据汇聚功能,将多个源集群的数据同步到一个汇聚集群的PolarSearch节点中。
整个数据流转过程如下:
在源集群中,通过
CREATE ETL GRANT授权汇聚集群访问本集群的数据。在汇聚集群中,通过ETL存储过程(
dbms_etl.sync_by_sql)创建汇聚链路,在源表的WITH子句中通过polardb-mysql-instance参数指定数据来源的PolarDB MySQL版集群。链路启动后,AutoETL链路会检查每个源集群是否已授权汇聚集群的ETL访问。确认授权后,引擎从各源集群读取数据并同步到汇聚集群的PolarSearch节点中。
使用限制
所有参与汇聚的PolarDB MySQL版集群必须属于同一个阿里云账号(UID)。
目前仅支持通过ETL存储过程(
dbms_etl.sync_by_sql)创建多集群汇聚链路,不支持搜索视图方式。删除ETL授权不影响已创建且正在运行的同步链路。
步骤一:创建ETL授权
为了允许汇聚集群同步源集群的数据,您需要在每个源集群中创建ETL授权,授权汇聚集群通过AutoETL访问本集群的数据。
创建ETL授权
在源集群中执行以下SQL语句,通过<allow_instance_id>指定允许ETL同步数据的汇聚集群ID。
CREATE ETL GRANT `<allow_instance_id>`;删除ETL授权
DROP ETL GRANT `<allow_instance_id>`;查看ETL授权
SHOW ETL GRANTS;步骤二:创建汇聚链路
您可以通过ETL存储过程(dbms_etl.sync_by_sql)创建多集群的汇聚链路。在源表的WITH子句中,通过polardb-mysql-instance参数指定该表所在的源集群ID。
准备数据
假设有两个PolarDB MySQL版集群(集群A和集群B),需要将集群A的数据汇聚到集群B的PolarSearch节点中。
在集群A中创建测试数据,并授权集群B进行ETL同步:
-- 在集群A中执行 CREATE DATABASE IF NOT EXISTS db1; 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, '1', '1'), (2, '1', '1'), (3, '1', '1'); -- 授权集群B访问本集群数据 -- 将pc-xxx替换为集群B的实际集群ID CREATE ETL GRANT `pc-xxx`;在集群B中创建测试数据:
-- 在集群B中执行 CREATE DATABASE IF NOT EXISTS db2; USE db2; CREATE TABLE IF NOT EXISTS t2 ( id INT PRIMARY KEY, c1 VARCHAR(100), c2 VARCHAR(100) ); INSERT INTO t2(id, c1, c2) VALUES (1, '2', '2'), (2, '2', '2'), (3, '2', '2');
定义汇聚链路
在集群B中执行以下SQL语句,创建汇聚链路将两个集群的数据同步到集群B的PolarSearch节点中。
源表和目标表的连接信息(如连接地址、端口、账号密码等)由系统自动配置,您无需在WITH子句中手动指定。仅需通过polardb-mysql-instance参数指定源表所在的集群ID。
-- 在集群B中执行
CALL dbms_etl.sync_by_sql("search", "
-- 步骤1:定义集群A的源表
CREATE TEMPORARY TABLE `db1`.`t1` (
`id` BIGINT,
`c1` STRING,
`c2` STRING,
PRIMARY KEY (`id`) NOT ENFORCED
) WITH (
'connector' = 'mysql',
'database-name' = 'db1',
'table-name' = 't1',
'polardb-mysql-instance' = 'pc-xxx' -- 替换为集群A的实际集群ID
);
-- 步骤2:定义集群B的源表
CREATE TEMPORARY TABLE `db2`.`t2` (
`id` BIGINT,
`c1` STRING,
`c2` STRING,
PRIMARY KEY (`id`) NOT ENFORCED
) WITH (
'connector' = 'mysql',
'database-name' = 'db2',
'table-name' = 't2',
'polardb-mysql-instance' = 'pc-xxx' -- 替换为集群B的实际集群ID
);
-- 步骤3:定义PolarSearch目标表
CREATE TEMPORARY TABLE `dest` (
`id` BIGINT,
`p1_c1` STRING,
`p2_c1` STRING,
PRIMARY KEY (`id`) NOT ENFORCED
) WITH (
'connector' = 'opensearch',
'index' = 'dest'
);
-- 步骤4:定义计算和插入逻辑
INSERT INTO `dest`
SELECT
`t1`.`id`,
`t1`.`c1`,
`t2`.`c1`
FROM `db1`.`t1` AS `t1`
LEFT JOIN `db2`.`t2` AS `t2`
ON `t1`.`id` = `t2`.`id`;
");步骤三:验证数据
连接到汇聚集群的PolarSearch节点,使用与Elasticsearch兼容的REST API进行查询,确认数据已同步。
# 将<user>:<password>替换为PolarSearch节点的账号密码
# 将<polarsearch_endpoint>替换为PolarSearch节点的连接地址与端口
curl -u <user>:<password> -X GET "http://<polarsearch_endpoint>/dest/_search"