全部产品
Search
文档中心

云原生数据库 PolarDB:AutoETL多集群数据汇聚

更新时间:Jun 16, 2026

AutoETL支持汇聚多个PolarDB MySQL版集群的业务数据到同一个PolarSearch节点中,便于您对分布在不同集群中的数据进行统一检索和分析。本文介绍如何通过AutoETL创建多集群数据汇聚链路。

说明

当前功能处于灰度阶段。如您有相关需求,请提交工单与我们联系,以便为您开启该功能。

背景信息

在实际业务场景中,数据可能分布在多个PolarDB MySQL版集群中。当您需要对这些分散的数据进行统一的全文检索或分析时,可以通过AutoETL的多集群数据汇聚功能,将多个源集群的数据同步到一个汇聚集群的PolarSearch节点中。

整个数据流转过程如下:

  1. 在源集群中,通过CREATE ETL GRANT授权汇聚集群访问本集群的数据。

  2. 在汇聚集群中,通过ETL存储过程(dbms_etl.sync_by_sql)创建汇聚链路,在源表的WITH子句中通过polardb-mysql-instance参数指定数据来源的PolarDB MySQL版集群。

  3. 链路启动后,AutoETL链路会检查每个源集群是否已授权汇聚集群的ETL访问。确认授权后,引擎从各源集群读取数据并同步到汇聚集群的PolarSearch节点中。

image

使用限制

  • 所有参与汇聚的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"

相关文档