StarRocks数据摄入连接器

更新时间:
复制 MD 格式

本文为您介绍如何在数据摄入 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选项来加快表结构变更的速度。

  • 使用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 数据的间隔

配置项

参数名称

描述

类型

是否必填

默认值

备注

type

连接器的名称。

String

固定值为starrocks

name

Sink的显示名称。

String

无。

jdbc-url

JDBC连接的URL。

String

支持传入多个地址,使用英文逗号 (,) 分隔。例如 jdbc:mysql://fe_host1:fe_query_port1,fe_host2:fe_query_port2,fe_host3:fe_query_port3

load-url

连接到FE节点的HTTP服务URL。

String

支持传入多个地址,使用英文分号 (;) 分隔。例如 fe_host1:fe_http_port1;fe_host2:fe_http_port2

username

连接到 StarRocks时使用的用户名。

String

该用户至少需要具备对目标表的SELECTINSERT权限。您可以使用StarRocksGRANT命令赋予相应的权限。

password

连接到 StarRocks时使用的密码。

String

无。

sink.semantic

数据写入语义。

String

at-least-once

当前仅支持 at-least-once。显式配置为 exactly-once 不会报错,但会被自动重置为 at-least-once。如需恰好一次语义,请使用 Flink SQL StarRocks 连接器,详情请参见本文 SQL 章节。

sink.label-prefix

在进行Stream Load导入时使用的标签前缀。

String

取值仅支持英文字母、数字、短横线(-)和下划线(_)。使用其他字符可能导致导入失败。

sink.connect.timeout-ms

建立HTTP连接时的超时时间。

Integer

30000

单位为毫秒,取值需要介于100 ~ 60000。

sink.wait-for-continue.timeout-ms

从服务器得到100 Continue请求前的超时时间)。

Integer

30000

单位为毫秒。取值需要介于3000 ~ 600000。

sink.buffer-flush.max-bytes

在将数据写入StarRocks前,最多可以在内存中缓存多少字节的数据。

Long

94371840

单位为字节,取值需要介于64 MB ~ 10 GB。

说明
  • 该缓存大小被所有表共用。当缓冲区已满时,连接器将选择若干张表进行Flush。

  • 将此参数设置为较大的值可以提高吞吐量,但可能会增加导入时的延迟。

sink.buffer-flush.max-rows

在将数据写入StarRocks前,最多可以在内存中缓存多少行数据。

Long

500000

取值范围需要介于1,000 和 5,000,000 之间。

sink.buffer-flush.interval-ms

每张表连续两次Flush之间的间隔时间。

Long

300000

单位为毫秒。

说明

对于同步的数据量不多的作业,需要将此参数适当降低,以免数据长时间无法落盘。

sink.max-retries

最大重试次数。

Long

3

取值范围需要介于 0 和 1000 之间。

sink.scan-frequency.ms

连续两次检查是否应该进行Flush之间的间隔时间。

Long

50

单位为毫秒。

sink.io.thread-count

在进行 Stream Load导入时的线程数量。

Integer

2

无。

sink.at-least-once.use-transaction-stream-load

是否使用Stream Load事务接口进行导入。

Boolean

true

仅在数据库支持的情况下生效。

sink.ignore.update-before

是否忽略更新操作中的 update-before 记录。

Boolean

true

当通过 Transform 模块变更了主键(例如使用 primary-keys 指定了与上游不同的主键),必须将 sink.ignore.update-before 设置为 false,否则旧主键对应的行不会被删除,导致数据残留。

仅实时计算引擎 VVR 11.8 及以上版本支持。

sink.ignore.delete

是否忽略删除记录。

Boolean

false

设置为 true 时,Delete 记录会被过滤,不写入 StarRocks。适用于希望在下游保留历史数据、仅同步插入和更新操作的场景。

仅实时计算引擎 VVR 11.8 及以上版本支持。

sink.properties.*

提供给Sink的额外参数。

String

可以在STREAM LOAD查看支持的参数。

table.create.num-buckets

自动建表时的Bucket数量。

Integer

  • StarRocks 2.5.7及更高版本:此参数可选,Bucket数会被自动推断

  • StarRocks 2.5.6及之前版本:此参数必填。

table.create.properties.*

在自动建表时需要传递的额外参数。

String

例如,可以传递'table.create.properties.fast_schema_evolution' = 'true'来启用快速表结构变更功能。参数详情请参见StarRocks文档

table.schema-change.timeout

执行表结构变更的超时时间。

Duration

30 min

必须设定为整数秒。

说明

如果某个表结构变更操作耗时超过此限制,作业将运行失败。

unicode-char.max-bytes

为每个Unicode字符分配多少个字节。

Integer

3

CDC中的VARCHAR类型长度为字符数,而StarRocks对应的VARCHAR长度则是字节数

大部分情况下Unicode字符经UTF-8编码后的字节数不会超过3,但部分生僻字及Emoji符号可能占用超过4字节。

sink.socket.time

StarRocks flush时的HTTP客户端超时时间

Long

-1

StarRocksflush发送stream load 请求时候的http 客户端超时时间,单位毫秒。-1 表示使用系统默认值,表示不超时,无限等。

仅实时计算引擎 VVR 11.8 及以上版本支持。

sink.close.eof-timeout-ms

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倍。

说明

StarRocksCHAR类型长度最长不可超过255,因此只有长度不超过85CDC CHAR类型才会被映射到StarRocks CHAR类型。

说明

设置unicode-char.max-bytes参数可以为每个Unicode字符分配更多字节的空间。

CHAR(n)

(n > 85 时)

VARCHAR(n × 3)

CDC中的CHAR类型长度表示字符数,而StarRocks中的CHAR类型长度表示UTF-8编码后的字节数。通常情况下,一个中文字符经过UTF-8编码后不会超过3字节,因此映射到的 StarRocks VARCHAR类型长度为原来的3倍。

说明

StarRocksCHAR类型长度最长不可超过255,因此长度大于85CDC CHAR类型会被映射到StarRocks VARCHAR类型。

说明

设置unicode-char.max-bytes参数可以为每个Unicode字符分配更多字节的空间。

VARCHAR(n)

VARCHAR(n × 3)

CDC中的VARCHAR类型长度表示字符数,而StarRocks中的VARCHAR类型长度表示UTF-8编码后的字节数。通常情况下,一个中文字符经过UTF-8编码后不会超过3字节,因此映射到的StarRocks VARCHAR类型长度为原来的3倍。

说明

设置unicode-char.max-bytes参数可以为每个Unicode字符分配更多字节的空间。

BINARY(n)

BINARY(n+2)

增加长度为2padding,防止数据问题。

VARBINARY(n)

VARBINARY(n+1)

增加长度为1padding,防止数据问题。

表结构变更

CDC YAML Pipeline作业在处理表结构变更时有不同的策略,通过pipeline级别的配置项schema.change.behavior来设置,取值有IGNORE、LENIENT、TRY_EVOLVE、EVOLVEEXCEPTION,默认值为LENIENT。事件是否执行还受Sinkinclude.schema.changesexclude.schema.changes配置影响。

其中LENIENTEVOLVE涉及到表结构变更,接下来会说明如何处理不同表结构变更事件(Schema Change Event)。

说明

以下说明针对普通一对一同步。多对一路由会先推导合并后的Schema,再规范化事件,不能将上述一对一行为直接理解为任一分表的删除操作都会直接删除合并目标。

支持的事件

  • CREATE TABLE EVENT

    说明

    在下游 StarRocks 表已经存在时,不会尝试重复建表。您需要保证下游表结构与上游兼容。

  • ADD COLUMN EVENT

    说明

    StarRocks 要求主键列总是位于最前。新插入的列也需要保持此限制。

  • 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执行器时,作业失败。

  • 列类型变更:支持。仅支持兼容的变更路径。

  • 列顺序调整:支持。要求主键列仍在最前且保持原主键顺序。

  • 删除表、清空表:未被排除时实际执行。

警告

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