StarRocks SQL连接器

更新时间:
复制 MD 格式

本文为您介绍如何在 SQL 作业中使用 StarRocks 连接器。

背景信息

StarRocks是新一代极速全场景MPP(Massively Parallel Processing)数据仓库,致力于构建极速和统一分析体验。StarRocks具有以下优势:

  • StarRocks兼容MySQL协议,可以使用MySQL客户端和常用BI工具对接StarRocks来分析数据。

  • StarRocks采用分布式架构:

    • 对数据表进行水平划分并以多副本存储。

    • 集群规模可以灵活伸缩,支持10 PB级别的数据分析。

    • 支持MPP框架,并行加速计算。

    • 支持多副本,具有弹性容错能力。

Flink连接器内部的结果表是通过缓存并批量由Stream Load导入实现,源表是通过批量读取数据实现。StarRocks连接器支持的信息如下。

类别

详情

支持类型

源表、维表和结果表、数据摄入目标端

运行模式

流模式和批模式

数据格式

CSV

特有监控指标

暂无

API种类

Datastream、SQL和数据摄入YAML

是否支持更新或删除结果表数据

前提条件

已创建StarRocks集群,包括EMRStarRocks或基于ECS的云上自建StarRocks。

使用限制

  • 仅实时计算引擎VVR 11.1及以上版本支持维表JOIN。

  • 为避免网络访问限制,必须将 StarRocks 集群的以下端口加入安全组或防火墙白名单:9030/8030/8040/9060/8060/9020。

  • 目标 StarRocks 表含隐藏生成列(Generated Column)时,需在 Flink 作业中仅声明实际写入的物理列,不包含生成列字段。StarRocks 表达式分区自动创建的隐藏生成列(如 __generated_partition_column_0)不接受外部写入,Connector 默认按完整 Schema 构造写入请求会导致作业失败。关于 StarRocks 生成列,请参见 Generated columns

特色功能

EMRStarRocks支持通过Flink CDC数据摄入作业实现单表的结构和数据同步,实现整库同步或者同一库中的多表结构和数据同步,详情请参见基于实时计算Flink使用CTAS&CDAS功能同步MySQL数据至StarRocks

语法结构

CREATE TABLE USER_RESULT(
 name VARCHAR,
 score BIGINT
 ) WITH (
 'connector' = 'starrocks',
 'jdbc-url'='jdbc:mysql://fe1_ip:query_port,fe2_ip:query_port,fe3_ip:query_port?xxxxx',
 'load-url'='fe1_ip:http_port;fe2_ip:http_port;fe3_ip:http_port',
 'database-name' = 'xxx',
 'table-name' = 'xxx',
 'username' = 'xxx',
 'password' = 'xxx'
 );

WITH参数

类型

参数

说明

数据类型

是否必填

默认值

备注

通用

connector

表类型。

String

固定值为starrocks。

jdbc-url

JDBC连接的URL。

String

指定FE(Front End)的IPJDBC端口,格式为jdbc:mysql://ip:port

database-name

StarRocks数据库名称。

String

无。

table-name

StarRocks表名称。

String

无。

username

StarRocks连接用户名。

String

无。

password

StarRocks连接密码。

String

无。

starrocks.create.table.properties

StarRocks表属性。

String

设置数据表初始属性,如引擎、副本数等。例如,'starrocks.create.table.properties' = 'buckets 8','starrocks.create.table.properties' = 'replication_num=1'。

源表独有

scan-url

数据扫描的url。

String

指定FE(Front End)的IPHTTP端口,格式为fe_ip:http_port;fe_ip:http_port

说明

填写多个IP和端口号时,请使用半角逗号(,)进行分隔。

scan.connect.timeout-ms

flink-connector-starrocks连接StarRocks的时间上限。

超过该时间上限,将报错。

String

1000

单位为毫秒。

scan.params.keep-alive-min

查询任务的保活时间。

String

10

无。

scan.params.query-timeout-s

查询任务的超时时间。

如果超过该时间,仍未返回查询结果,则停止查询任务。

String

600

单位为秒。

scan.params.mem-limit-byte

BE节点中单个查询的内存上限。

String

1073741824(1 GB)

单位为字节。

scan.max-retries

查询失败时的最大重试次数。

超过该数量上限,则将报错。

String

1

无。

结果表独有

load-url

数据导入的URL。

String

指定FE(Front End)的IPHTTP端口,格式为fe_ip:http_port;fe_ip:http_port

说明

填写多个IP和端口号时,请使用半角分号(;)进行分隔。

sink.semantic

数据写入语义。

String

at-least-once

取值如下:

  • at-least-once(默认值):至少一次。

  • exactly-once:恰好一次。

sink.buffer-flush.max-bytes

Buffer可容纳的最大数据量。

String

94371840(90 MB)

取值范围为64 MB~10 GB。

sink.buffer-flush.max-rows

Buffer可容纳的最大数据行数。

String

500000

取值范围为1,000~5000,000。

sink.buffer-flush.interval-ms

Buffer刷新时间间隔。

String

300000

取值范围为1000毫秒~3600000毫秒。

sink.max-retries

最大重试次数。

String

3

取值范围为0~1000。

sink.connect.timeout-ms

连接到starrocks的超时时间。

String

1000

取值范围为100~60000。单位为毫秒。

sink.properties.*

结果表属性。

String

Stream Load的参数控制Stream Load导入行为。例如,参数 sink.properties.format表示Stream Load所导入的数据格式,如CSV。更多参数和解释,请参见Stream Load

维表独有

lookup.cache.enabled

是否启用维表缓存机制。

Boolean

true

取值如下:

  • true:启用。首次读取表数据后缓存至内存,后续请求在缓存有效期内直接使用内存数据,减少IO开销。

  • false:关闭。每次查询均直接访问数据源。

重要
  • 仅实时计算引擎VVR 11.1及以上版本支持。

  • 建议关闭场景:

    • 维表数据更新频繁,需保证实时性;

    • 单表数据量过大,避免内存溢出风险。

类型映射

StarRocks字段类型

Flink字段类型

NULL

NULL

BOOLEAN

BOOLEAN

TINYINT

TINYINT

SMALLINT

SMALLINT

INT

INT

BIGINT

BIGINT

BIGINT UNSIGNED

说明

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

DECIMAL(20,0)

LARGEINT

DECIMAL(20,0)

FLOAT

FLOAT

DOUBLE

DOUBLE

DATE

DATE

DATETIME

TIMESTAMP

DECIMAL

DECIMAL

DECIMALV2

DECIMAL

DECIMAL32

DECIMAL

DECIMAL64

DECIMAL

DECIMAL128

DECIMAL

CHAR(m)

说明
  • 仅实时计算引擎VVR 8.0.10版本,CHAR类型长度自动扩展至三倍(m=n*3,n<=85),以适配MySQLStarRocks之间的编码差异。

  • 仅实时计算引擎VVR 8.0.11及以上版本,CHAR类型长度自动扩展至四倍(m=n*4,n<=63),以适配MySQLStarRocks之间的编码差异。

  • StarRocks CHAR类型长度最长不可超过255,因此只有Flink CHAR类型长度自动扩容后不超过255才会被映射到StarRocks CHAR类型。

CHAR(n)

VARCHAR(m)

说明
  • 仅实时计算引擎VVR 8.0.10版本,VARCHAR类型长度自动扩展至三倍(m=n*3,n>85),以适配MySQLStarRocks之间的编码差异。

  • 仅实时计算引擎VVR 8.0.11及以上版本,VARCHAR类型长度自动扩展至四倍(m=n*4,n>63),以适配MySQLStarRocks之间的编码差异。

  • StarRocks CHAR类型长度最长不可超过255,因此Flink CHAR类型长度自动扩容后超过255会被映射到StarRocks VARCHAR类型。

CHAR(n)

VARCHAR

STRING

VARBINARY

说明

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

VARBINARY

使用示例

以下示例中的源表、结果表和维表可组合使用:源表示例声明的runoob_tbl_source同时作为结果表示例和维表示例的数据来源。使用前请在StarRocks中准备字段类型匹配的表,并替换示例中的地址、库表名与密钥变量。

源表示例

CREATE TEMPORARY TABLE runoob_tbl_source (
    runoob_id BIGINT NOT NULL,
    runoob_title STRING NOT NULL,
    runoob_author STRING NOT NULL,
    submission_date DATE,
    proc_time AS PROCTIME()                                  --处理时间列,维表示例的Temporal Join使用
) WITH (
    'connector' = 'starrocks',
    'jdbc-url' = 'jdbc:mysql://<fe_host>:9030',              --FE的MySQL协议地址,多地址用逗号分隔
    'scan-url' = '<fe_host>:8030',                           --FE的HTTP地址,多地址用逗号分隔
    'database-name' = '<database_name>',
    'table-name' = '<source_table_name>',
    'username' = '${secret_values.starrocks_username}',      --推荐使用变量管理防止密钥泄露
    'password' = '${secret_values.starrocks_password}',
    'scan.params.query-timeout-s' = '600'                    --单次查询超时时间,单位秒
);
说明

源表为批量读取,作业启动后扫描目标表当前数据,读取完成后源表结束输出,不持续产出增量变更。

结果表示例

说明

在 StarRocks 中,表定义允许主键列为NULLABLE,但 Flink 不支持主键包含可空列。要求主键必须具有唯一且非空的语义,这是其数据一致性模型的基础,否则将抛出错误:Invalid primary key. Column 'xxx' is nullable。详情请参见报错:“Invalid primary key. Column 'xxx' is nullable.”

VVR 11+

CREATE TEMPORARY TABLE runoob_tbl_sink (
    runoob_id BIGINT NOT NULL,                               --主键列必须声明NOT NULL
    runoob_title STRING NOT NULL,
    runoob_author STRING NOT NULL,
    submission_date DATE,
    PRIMARY KEY (runoob_id) NOT ENFORCED
) WITH (
    'connector' = 'starrocks',
    'jdbc-url' = 'jdbc:mysql://<fe_host>:9030',              --FE的MySQL协议地址,多地址用逗号分隔
    'load-url' = '<fe_host>:8030',                           --FE的HTTP地址,多地址用分号分隔
    'database-name' = '<database_name>',
    'table-name' = '<sink_table_name>',
    'username' = '${secret_values.starrocks_username}',      --推荐使用变量管理防止密钥泄露
    'password' = '${secret_values.starrocks_password}',
    'sink.version' = 'V2',                                   --事务Stream Load,需StarRocks 2.4及以上版本
    'sink.semantic' = 'at-least-once',                       --可选exactly-once,需开启Checkpoint
    'sink.buffer-flush.interval-ms' = '5000'                 --攒批写入间隔,单位毫秒
);

INSERT INTO runoob_tbl_sink
SELECT runoob_id, runoob_title, runoob_author, submission_date
FROM runoob_tbl_source;

VVR 8+

VVR 8.xsink.version默认值为V1(支持V1/V2/AUTO),但不支持sink.ignore.update-before等仅VVR 11.x注册的参数:

CREATE TEMPORARY TABLE runoob_tbl_sink (
    runoob_id BIGINT NOT NULL,                               --主键列必须声明NOT NULL
    runoob_title STRING NOT NULL,
    runoob_author STRING NOT NULL,
    submission_date DATE,
    PRIMARY KEY (runoob_id) NOT ENFORCED
) WITH (
    'connector' = 'starrocks',
    'jdbc-url' = 'jdbc:mysql://<fe_host>:9030',              --FE的MySQL协议地址,多地址用逗号分隔
    'load-url' = '<fe_host>:8030',                           --FE的HTTP地址,多地址用分号分隔
    'database-name' = '<database_name>',
    'table-name' = '<sink_table_name>',
    'username' = '${secret_values.starrocks_username}',      --推荐使用变量管理防止密钥泄露
    'password' = '${secret_values.starrocks_password}',
    'sink.semantic' = 'at-least-once',                       --可选exactly-once,需开启Checkpoint
    'sink.buffer-flush.interval-ms' = '5000',                --攒批写入间隔,单位毫秒
    'sink.max-retries' = '3'                                 --Stream Load失败后的重试次数
);

INSERT INTO runoob_tbl_sink
SELECT runoob_id, runoob_title, runoob_author, submission_date
FROM runoob_tbl_source;

维表示例

维表JOIN仅实时计算引擎VVR 11.1及以上版本支持。维表需声明主键作为关联键,使用处理时间进行Temporal Join。以下示例承接源表示例中的runoob_tbl_source(含proc_time处理时间列):

CREATE TEMPORARY TABLE sr_dim (
    runoob_id BIGINT NOT NULL,
    runoob_author STRING,
    PRIMARY KEY (runoob_id) NOT ENFORCED                     --维表需声明主键作为关联键
) WITH (
    'connector' = 'starrocks',
    'jdbc-url' = 'jdbc:mysql://<fe_host>:9030',              --FE的MySQL协议地址,多地址用逗号分隔
    'scan-url' = '<fe_host>:8030',                           --FE的HTTP地址,多地址用逗号分隔
    'database-name' = '<database_name>',
    'table-name' = '<dim_table_name>',
    'username' = '${secret_values.starrocks_username}',      --推荐使用变量管理防止密钥泄露
    'password' = '${secret_values.starrocks_password}',
    'lookup.cache.ttl-ms' = '5000',                          --维表缓存存活时间,单位毫秒
    'lookup.cache.enabled' = 'true'                          --默认true;false为直连查询不走缓存
);

SELECT o.runoob_id, o.runoob_title, d.runoob_author
FROM runoob_tbl_source AS o
JOIN sr_dim FOR SYSTEM_TIME AS OF o.proc_time AS d           --Temporal Join,使用源表的处理时间列
ON o.runoob_id = d.runoob_id;