本文为您介绍如何在 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集群,包括EMR的StarRocks或基于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。
特色功能
EMR的StarRocks支持通过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)的IP和JDBC端口,格式为 | |
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)的IP和HTTP端口,格式为 说明 填写多个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)的IP和HTTP端口,格式为 说明 填写多个IP和端口号时,请使用半角分号(;)进行分隔。 |
sink.semantic | 数据写入语义。 | String | 否 | at-least-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 | 取值如下:
重要
|
类型映射
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) 说明
| CHAR(n) |
VARCHAR(m) 说明
| 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.x的sink.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;