本文为您介绍如何在数据摄入 YAML 作业中,使用 StarRocks 连接器进行数据同步。
背景信息
StarRocks是面向实时分析的MPP(Massively Parallel Processing)数据仓库,兼容MySQL协议,采用分布式架构,支持集群弹性伸缩和并行计算。StarRocks YAML连接器支持将上游数据记录和表结构变更写入StarRocks,同时支持社区版StarRocks和阿里云E-MapReduce Serverless StarRocks全托管版本。StarRocks YAML连接器支持的信息如下。
类别 | 详情 |
支持类型 | 数据摄入目标端(Sink) |
运行模式 | 流模式和批模式 |
数据格式 | JSON |
特有监控指标 | 暂无 |
API种类 | YAML |
是否支持更新或删除结果表数据 | 是 |
YAML连接器当前仅支持至少一次语义。即使显式设置sink.semantic: exactly-once,也会被覆盖为at-least-once,不会报错。如需恰好一次语义,请参见StarRocks SQL连接器。
特色功能
自动建库建表。
如果来自上游的数据库及数据表不存在于下游StarRocks实例中,则对应的数据库及数据表会被自动创建。您可以通过
table.create.properties.*参数设定自动创建表时的选项。表结构变更同步。
目前,StarRocks连接器支持自动将建表事件(CreateTableEvent)、增加列事件(AddColumnEvent)和删除列(DropColumnEvent)事件自动应用到下游数据库中。
实时计算引擎VVR 11.1及以上版本支持兼容的列类型变更,详情请参见ALTER TABLE | StarRocks。
注意事项
目前,同步的表必须包含主键。不含主键的表必须通过
transform语句块指定主键方可正常写入下游。例如:transform: - source-table: ... primary-keys: id, ...自动创建的表分桶键与主键相同,且不可有分区键。
进行表结构变更同步时,新增列只能追加到已有列的尾部。在默认的表结构演化模式Lenient下,会自动将其他位置的插入转换到尾部。
如果您使用的StarRocks版本低于2.5.7,则必须显式地通过
table.create.num-buckets参数指定分桶数量。更高版本的StarRocks可以自动设定合适的分桶数。如果您使用的是StarRocks 3.2或更高版本,建议开启
table.create.properties.fast_schema_evolution选项来加快表结构变更的速度。如果您使用的是实时计算引擎 VVR 11.9 及更高版本、StarRocks 3.3.2及更高版本,则可以将上游的重命名列事件同步到下游。
使用CDC YAML数据摄入写入 EMR Serverless StarRocks 时可能出现串流问题,您可以采用以下选项之一来规避:
使用 Flink SQL StarRocks 连接器,并使用
sink.version=V1参数;开启 FE emr_internal_redirect 参数;
使用 StarRocks Private Zone 域名而不是 SLB。
语法结构
source:
...
sink:
type: starrocks
name: StarRocks Sink
jdbc-url: jdbc:mysql://127.0.0.1:9030
load-url: 127.0.0.1:8030
username: root
password: pass
sink.buffer-flush.interval-ms: 5000 # 设定 Flush 数据的间隔配置项
参数名称 | 描述 | 类型 | 是否必填 | 默认值 | 备注 |
| 连接器的名称。 | String | 是 | 无 | 固定值为 |
| Sink的显示名称。 | String | 否 | 无 | 无。 |
| JDBC连接的URL。 | String | 是 | 无 | 支持传入多个地址,使用英文逗号 ( |
| 连接到FE节点的HTTP服务URL。 | String | 是 | 无 | 支持传入多个地址,使用英文分号 ( |
| 连接到 StarRocks时使用的用户名。 | String | 是 | 无 | 该用户至少需要具备对目标表的SELECT和INSERT权限。您可以使用StarRocks的GRANT命令赋予相应的权限。 |
| 连接到 StarRocks时使用的密码。 | String | 是 | 无 | 无。 |
| 数据写入语义。 | String | 否 | at-least-once | 当前仅支持 at-least-once。显式配置为 exactly-once 不会报错,但会被自动重置为 at-least-once。如需恰好一次语义,请使用 Flink SQL StarRocks 连接器,详情请参见本文 SQL 章节。 |
| 在进行Stream Load导入时使用的标签前缀。 | String | 否 | 无 | 取值仅支持英文字母、数字、短横线(-)和下划线(_)。使用其他字符可能导致导入失败。 |
| 建立HTTP连接时的超时时间。 | Integer | 否 | 30000 | 单位为毫秒,取值需要介于100 ~ 60000。 |
| 从服务器得到100 Continue请求前的超时时间)。 | Integer | 否 | 30000 | 单位为毫秒。取值需要介于3000 ~ 600000。 |
| 在将数据写入StarRocks前,最多可以在内存中缓存多少字节的数据。 | Long | 否 | 94371840 | 单位为字节,取值需要介于64 MB ~ 10 GB。 说明
|
| 在将数据写入StarRocks前,最多可以在内存中缓存多少行数据。 | Long | 否 | 500000 | 取值范围需要介于1,000 和 5,000,000 之间。 |
| 每张表连续两次Flush之间的间隔时间。 | Long | 否 | 300000 | 单位为毫秒。 说明 对于同步的数据量不多的作业,需要将此参数适当降低,以免数据长时间无法落盘。 |
| 最大重试次数。 | Long | 否 | 3 | 取值范围需要介于 0 和 1000 之间。 |
| 连续两次检查是否应该进行Flush之间的间隔时间。 | Long | 否 | 50 | 单位为毫秒。 |
| 在进行 Stream Load导入时的线程数量。 | Integer | 否 | 2 | 无。 |
| 是否使用Stream Load事务接口进行导入。 | Boolean | 否 | true | 仅在数据库支持的情况下生效。 |
| 是否忽略更新操作中的 update-before 记录。 | Boolean | 否 | true | 当通过 Transform 模块变更了主键(例如使用 primary-keys 指定了与上游不同的主键),必须将 sink.ignore.update-before 设置为 false,否则旧主键对应的行不会被删除,导致数据残留。 仅实时计算引擎 VVR 11.8 及以上版本支持。 |
| 是否忽略删除记录。 | Boolean | 否 | false | 设置为 true 时,Delete 记录会被过滤,不写入 StarRocks。适用于希望在下游保留历史数据、仅同步插入和更新操作的场景。 仅实时计算引擎 VVR 11.8 及以上版本支持。 |
| 提供给Sink的额外参数。 | String | 否 | 无 | 可以在STREAM LOAD查看支持的参数。 |
| 自动建表时的Bucket数量。 | Integer | 否 | 无 |
|
| 在自动建表时需要传递的额外参数。 | String | 否 | 无 | 例如,可以传递 |
| 执行表结构变更的超时时间。 | Duration | 否 | 30 min | 必须设定为整数秒。 说明 如果某个表结构变更操作耗时超过此限制,作业将运行失败。 |
| 为每个Unicode字符分配多少个字节。 | Integer | 否 | 3 | CDC中的VARCHAR类型长度为字符数,而StarRocks对应的VARCHAR长度则是字节数。 大部分情况下Unicode字符经UTF-8编码后的字节数不会超过3,但部分生僻字及Emoji符号可能占用超过4字节。 |
| 向StarRocks flush时的HTTP客户端超时时间 | Long | 否 | -1 | 向StarRocks中flush发送stream load 请求时候的http 客户端超时时间,单位毫秒。-1 表示使用系统默认值,表示不超时,无限等。 仅实时计算引擎 VVR 11.8 及以上版本支持。 |
| close的超时时间 | Long | 否 | 60000 | 作业关闭时等待 Flush 队列timeout结束的超时时间,单位毫秒。仅实时计算引擎 VVR 11.8 及以上版本支持。 |
复用已有 Catalog
自VVR 11.5版本起,您可以在Flink CDC数据摄入作业中直接引用“数据管理”页面中创建的内置StarRocks Catalog,减少手写连接属性工作量。
sink:
type: starrocks
using.built-in-catalog: starrocks_catalog目前,数据摄入作业支持自动复用以下 StarRocks Catalog 参数:
jdbc-url
http-url
username
password
table.num-buckets
如果希望覆盖以上自动复用的参数,可显式写出相应的 YAML 参数,其具备更高的优先级。
类型映射
StarRocks并不支持所有的CDC YAML类型,尝试将不支持的类型写入下游会导致作业失败。您可以使用Transform CAST内置函数对不支持的数据进行转换,或是使用Projection语句将其从结果表中移除。详情请参考数据摄入作业开发参考。
CDC类型 | StarRocks类型 | 附注 |
TINYINT | TINYINT | 无 |
SMALLINT | SMALLINT | |
INT | INT | |
BIGINT | BIGINT | |
FLOAT | FLOAT | |
DOUBLE | DOUBLE | |
BOOLEAN | BOOLEAN | |
DATE | DATE | |
TIMESTAMP | DATETIME | |
TIMESTAMP_LTZ | DATETIME | |
DECIMAL(p, s) | DECIMAL(p, s) | StarRocks不支持DECIMAL作为主键。因此当上游数据表的字段类型为DECIMAL且该字段作为主键时,同步至StarRocks的表结构会自动将主键字段类型从DECIMAL变更为VARCHAR。 |
CHAR(n) (n <= 85 时) | CHAR(n × 3) | CDC中的CHAR类型长度表示字符数,而StarRocks中的CHAR类型长度表示UTF-8编码后的字节数。通常情况下,一个中文字符经过UTF-8编码后不会超过3字节,因此映射到的StarRocks CHAR类型长度为原来的3倍。 说明 StarRocks的CHAR类型长度最长不可超过255,因此只有长度不超过85的CDC CHAR类型才会被映射到StarRocks CHAR类型。 说明 设置 |
CHAR(n) (n > 85 时) | VARCHAR(n × 3) | CDC中的CHAR类型长度表示字符数,而StarRocks中的CHAR类型长度表示UTF-8编码后的字节数。通常情况下,一个中文字符经过UTF-8编码后不会超过3字节,因此映射到的 StarRocks VARCHAR类型长度为原来的3倍。 说明 StarRocks的CHAR类型长度最长不可超过255,因此长度大于85的CDC CHAR类型会被映射到StarRocks VARCHAR类型。 说明 设置 |
VARCHAR(n) | VARCHAR(n × 3) | CDC中的VARCHAR类型长度表示字符数,而StarRocks中的VARCHAR类型长度表示UTF-8编码后的字节数。通常情况下,一个中文字符经过UTF-8编码后不会超过3字节,因此映射到的StarRocks VARCHAR类型长度为原来的3倍。 说明 设置 |
BINARY(n) | BINARY(n+2) | 增加长度为2的padding,防止数据问题。 |
VARBINARY(n) | VARBINARY(n+1) | 增加长度为1的padding,防止数据问题。 |
表结构变更
CDC YAML Pipeline作业在处理表结构变更时有不同的策略,通过pipeline级别的配置项schema.change.behavior来设置,取值有IGNORE、LENIENT、TRY_EVOLVE、EVOLVE和EXCEPTION,默认值为LENIENT。事件是否执行还受Sink的include.schema.changes、exclude.schema.changes配置影响。
其中LENIENT和EVOLVE涉及到表结构变更,接下来会说明如何处理不同表结构变更事件(Schema Change Event)。
以下说明针对普通一对一同步。多对一路由会先推导合并后的Schema,再规范化事件,不能将上述一对一行为直接理解为任一分表的删除操作都会直接删除合并目标。
支持的事件
CREATE TABLE EVENT
说明在下游 StarRocks 表已经存在时,不会尝试重复建表。您需要保证下游表结构与上游兼容。
ADD COLUMN EVENT
说明StarRocks 要求主键列总是位于最前。新插入的列也需要保持此限制。
ALTER COLUMN TYPE EVENT
说明支持的表结构变更路径,请参考 StarRocks 官方文档。
RENAME COLUMN EVENT
说明需要实时计算引擎版本高于 VVR 11.9、且 StarRocks 集群版本高于 3.3.2。
DROP COLUMN EVENT
TRUNCATE TABLE EVENT
DROP TABLE EVENT
LENIENT(默认)
LENIENT模式下支持的Schema变更策略详情如下:
添加可空列:会自动在结果表Schema末尾添加对应的列,并自动同步新增列的数据。
添加非空列:会自动在结果表Schema末尾添加对应的列,新增的列会设置为可空列,对于添加列发生之前的数据自动设置为NULL值。
删除列:不会直接在结果表中删除该列。旧列原为非空时,将其变更为可空列;已可空时不产生变更。
重命名列:被看作为添加新列并保留旧列。在结果表末尾添加改名后的可空列;旧列原为非空时,将其变更为可空列。
列顺序调整:忽略,不同步到下游。
列类型变更:不因LENIENT模式自动忽略,仍交给StarRocks执行,仅支持兼容的变更路径。允许的路径请参见ALTER TABLE,VVR 11.1及以上版本的支持范围沿用原帮助文档。
删除表、清空表:默认被排除,不同步到下游。
默认排除删除表和清空表,是YAML解析器在未配置exclude.schema.changes时加入的默认项,并不是无条件的保护。显式配置排除列表(包括空列表)会取代这两个默认排除项,请自行保留需要禁止的事件。LENIENT也不保证忽略DDL执行失败。
EVOLVE
EVOLVE模式下支持的Schema变更策略详情如下:
添加列:支持。StarRocks执行器忽略上游新增位置,追加到已有列尾部。
删除列:支持。会在结果表中实际删除该列。
重命名列:支持(需要 StarRocks 3.3.2 及更高版本)。会修改下游列的名字。
列类型变更:支持。仅支持兼容的变更路径。
列顺序调整:支持。要求主键列仍在最前且保持原主键顺序。
删除表、清空表:未被排除时实际执行。
在EVOLVE模式下,如果在未删除结果表的情况下无状态重启,有可能出现上游数据与结果表的结构不一致的情况导致作业失败,需要用户手动调整下游表结构。
代码示例
下面展示几个典型的用户使用场景下的配置示例。
单表同步
将 MySQL 中的一张表同步到 StarRocks,下游库表不存在时会自动创建为主键表。
pipeline:
name: MySQL to StarRocks Pipeline
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: <yourUsername>
password: ${secret_values.mysql_password}
tables: test_db.test_source_table
server-id: 5401-5499
# (可选)增量阶段实时同步新创建表的数据,无需重启作业
scan.binlog.newly-added-table.enabled: true
# (可选)向下游同步表注释和字段注释
include-comments.enabled: true
# (可选)只解析被捕获表的 binlog,加速读取
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: starrocks
name: StarRocks Sink
jdbc-url: jdbc:mysql://<yourFeHostname>:9030
load-url: <yourFeHostname>:8030
username: <yourUsername>
password: ${secret_values.starrocks_password}
# (可选)数据量不大的作业建议调低 flush 间隔,避免数据长时间不落盘(默认 300000,即 5 分钟)
sink.buffer-flush.interval-ms: 5000
# (可选)上游为 utf8mb4 字符集时建议设为 4,避免文本截断(默认 3)
unicode-char.max-bytes: 4
# (可选)自动建表的分桶数;StarRocks 2.5.7 以下必须显式配置,更高版本可自动推断
table.create.num-buckets: 8
# (可选)自动建表的副本数,按集群情况配置
table.create.properties.replication_num: 3
# (可选)StarRocks 3.2 及以上建议开启,加速表结构变更
table.create.properties.fast_schema_evolution: true
# 注意:通过 transform 变更主键时,必须同时设置 sink.ignore.update-before: false,
# 否则旧主键对应的行会残留在下游
pipeline:
name: MySQL to StarRocks Pipeline整库同步
将 MySQL 一个库中的所有表一次性同步到 StarRocks,下游自动建库、建主键表,无需提前逐表建表。
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: <yourUsername>
password: ${secret_values.mysql_password}
# 正则匹配整个库的表;多个库可用英文逗号分隔写多个模式
tables: test_db.\.*
server-id: 5401-5499
# (可选)增量阶段实时同步新创建表的数据,无需重启作业
scan.binlog.newly-added-table.enabled: true
# (可选)向下游同步表注释和字段注释
include-comments.enabled: true
sink:
type: starrocks
name: StarRocks Sink
jdbc-url: jdbc:mysql://<yourFeHostname>:9030
load-url: <yourFeHostname>:8030
username: <yourUsername>
password: ${secret_values.starrocks_password}
# (可选)数据量不大的作业建议调低 flush 间隔,避免数据长时间不落盘(默认 300000,即 5 分钟)
sink.buffer-flush.interval-ms: 5000
# (可选)上游为 utf8mb4 字符集时建议设为 4,避免文本截断(默认 3)
unicode-char.max-bytes: 4
# (可选)自动建表的分桶数;StarRocks 2.5.7 以下必须显式配置,更高版本可自动推断
table.create.num-buckets: 8
# (可选)自动建表的副本数,按集群情况配置
table.create.properties.replication_num: 3
# (可选)StarRocks 3.2 及以上建议开启,加速表结构变更
table.create.properties.fast_schema_evolution: true
pipeline:
name: MySQL to StarRocks Pipeline整库同步时排除部分表
整库同步时用正则跳过不需要同步到下游的表,例如临时表、敏感表。
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: <yourUsername>
password: ${secret_values.mysql_password}
tables: test_db.\.*
# 此正则命中的表都不会被同步
tables.exclude: test_db.tmp_.\*
server-id: 5401-5499
sink:
type: starrocks
name: StarRocks Sink
jdbc-url: jdbc:mysql://<yourFeHostname>:9030
load-url: <yourFeHostname>:8030
username: <yourUsername>
password: ${secret_values.starrocks_password}
# (可选)导入接口版本:V2 需 StarRocks 2.4 及以上;EMR Serverless 如遇串流可改为 V1
sink.version: V2
# (可选)数据量不大的作业建议调低,避免数据长时间不落盘(默认 300000,即 5 分钟)
sink.buffer-flush.interval-ms: 5000
# (可选)自动建表的分桶数;StarRocks 2.5.7 以下必须显式配置
table.create.num-buckets: 8
# (可选)StarRocks 3.2 及以上建议开启,加速表结构变更
table.create.properties.fast_schema_evolution: true
pipeline:
name: MySQL to StarRocks Pipeline同步到指定库表
下游 StarRocks 的库名、表名需要与上游不一致时(例如写入 ODS 层库),用 route 统一重命名。
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: <yourUsername>
password: ${secret_values.mysql_password}
tables: test_db.\.*
server-id: 5401-5499
sink:
type: starrocks
name: StarRocks Sink
jdbc-url: jdbc:mysql://<yourFeHostname>:9030
load-url: <yourFeHostname>:8030
username: <yourUsername>
password: ${secret_values.starrocks_password}
# (可选)数据量不大的作业建议调低,避免数据长时间不落盘(默认 300000,即 5 分钟)
sink.buffer-flush.interval-ms: 5000
# (可选)自动建表的分桶数;StarRocks 2.5.7 以下必须显式配置
table.create.num-buckets: 8
route:
# 将 MySQL test_db 中的所有表同步到 StarRocks test_db2 库,表名保持不变;
# <> 为占位符,会被匹配到的源表名替换
- source-table: test_db.\.*
sink-table: test_db2.<>
replace-symbol: <>
pipeline:
name: MySQL to StarRocks Pipeline分库分表合并
将多张结构相同的分表合并写入一张 StarRocks 表,统一查询分析;要求各分表 schema 一致。
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: <yourUsername>
password: ${secret_values.mysql_password}
# 匹配所有分表,如 user_0、user_1……
tables: test_db.user\.*
server-id: 5401-5499
sink:
type: starrocks
name: StarRocks Sink
jdbc-url: jdbc:mysql://<yourFeHostname>:9030
load-url: <yourFeHostname>:8030
username: <yourUsername>
password: ${secret_values.starrocks_password}
# (可选)数据量不大的作业建议调低,避免数据长时间不落盘(默认 300000,即 5 分钟)
sink.buffer-flush.interval-ms: 5000
# (可选)合并表建议显式指定分桶数,按合并后的数据量评估
table.create.num-buckets: 8
route:
# 所有分表合并到一张 StarRocks test_db.user 表
- source-table: test_db.user\.*
sink-table: test_db.user
pipeline:
name: MySQL to StarRocks Pipeline开启 EVOLVE 模式
默认(LENIENT 模式)下,删列、删表、清空表等结构变更不会同步到下游;如确需严格同步表结构,可开启 EVOLVE 模式。该模式限制明显,使用前务必确认以下注意事项。
限制与注意事项
不支持列改名:上游出现列改名事件时作业会直接失败。
删列、删表、清空表会在下游真实执行:上游的误操作会直接影响下游表;默认 LENIENT 模式下删表、清空表事件不会同步,更安全。
在未删除结果表的情况下无状态重启,可能出现上游与结果表结构不一致导致作业失败,需要手动调整下游表结构。
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: <yourUsername>
password: ${secret_values.mysql_password}
tables: test_db.test_source_table
server-id: 5401-5499
sink:
type: starrocks
name: StarRocks Sink
jdbc-url: jdbc:mysql://<yourFeHostname>:9030
load-url: <yourFeHostname>:8030
username: <yourUsername>
password: ${secret_values.starrocks_password}
pipeline:
name: MySQL to StarRocks Pipeline
# 开启 EVOLVE:严格同步表结构变更,遇到不支持的变更(如列改名)作业失败
schema.change.behavior: evolve