本文為您介紹如何在 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。
使用限制
僅Realtime Compute引擎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資料攝入作業實現單表的結構和資料同步,實現整庫同步或者同一庫中的多表結構和資料同步,詳情請參見基於Realtime ComputeFlink使用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 說明 僅Realtime Compute引擎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 說明 僅Realtime Compute引擎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僅Realtime Compute引擎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;