全部產品
Search
文件中心

Realtime Compute for Apache Flink:StarRocks SQL連接器

更新時間:Sep 17, 2026

本文為您介紹如何在 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連接埠,格式為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)的IP和HTTP連接埠,格式為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)的IP和HTTP連接埠,格式為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:關閉。每次查詢均直接存取資料來源。

重要
  • 僅Realtime Compute引擎VVR 11.1及以上版本支援。

  • 建議關閉情境:

    • 維表資料更新頻繁,需保證即時性;

    • 單表資料量過大,避免記憶體溢出風險。

類型映射

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)

說明
  • 僅Realtime Compute引擎VVR 8.0.10版本,CHAR類型長度自動擴充至三倍(m=n*3,n<=85),以適配MySQL和StarRocks之間的編碼差異。

  • 僅Realtime Compute引擎VVR 8.0.11及以上版本,CHAR類型長度自動擴充至四倍(m=n*4,n<=63),以適配MySQL和StarRocks之間的編碼差異。

  • StarRocks CHAR類型長度最長不可超過255,因此只有Flink CHAR類型長度自動擴容後不超過255才會被映射到StarRocks CHAR類型。

CHAR(n)

VARCHAR(m)

說明
  • 僅Realtime Compute引擎VVR 8.0.10版本,VARCHAR類型長度自動擴充至三倍(m=n*3,n>85),以適配MySQL和StarRocks之間的編碼差異。

  • 僅Realtime Compute引擎VVR 8.0.11及以上版本,VARCHAR類型長度自動擴充至四倍(m=n*4,n>63),以適配MySQL和StarRocks之間的編碼差異。

  • StarRocks CHAR類型長度最長不可超過255,因此Flink CHAR類型長度自動擴容後超過255會被映射到StarRocks VARCHAR類型。

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;