StarRocks DataStream连接器支持通过DataStream API读取和写入StarRocks数据。本文介绍连接器依赖及Source、Lookup Source和Sink的开发示例。
配置连接器依赖
通过DataStream的方式读写数据时,需要使用对应的DataStream连接器连接实时计算Flink版,DataStream连接器设置方法请参见DataStream连接器设置方法。StarRocks DataStream连接器已发布到Maven中央仓库。本地运行和调试包含连接器的作业介绍了本地调试所需的Uber JAR。
在Maven项目的pom.xml文件中添加以下依赖:
<dependency>
<groupId>com.alibaba.ververica</groupId>
<artifactId>ververica-connector-starrocks</artifactId>
<version>${vvr-version}</version>
</dependency>
StarRocks Source
StarRocks Source通过批量扫描读取数据,是有界读取,不代表持续订阅StarRocks的增量变更。如需读取变更流,请使用StarRocks数据摄入连接器从上游数据库同步。
以下示例构建StarRocks Source并读取数据,该写法在VVR 8.x与11.x上一致:
import com.starrocks.connector.flink.StarRocksSource;
import com.starrocks.connector.flink.table.source.StarRocksSourceOptions;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.api.DataTypes;
import org.apache.flink.table.api.TableSchema;
public class Sample {
public static void main(String[] args) throws Exception {
// StarRocks连接参数。
StarRocksSourceOptions options = StarRocksSourceOptions.builder()
// FE的HTTP地址,多地址用逗号分隔
.withProperty("scan-url", "<fe_host>:8030")
// FE的MySQL协议地址
.withProperty("jdbc-url", "jdbc:mysql://<fe_host>:9030")
.withProperty("username", "<yourUsername>")
.withProperty("password", "<yourPassword>")
.withProperty("database-name", "<yourDatabaseName>")
.withProperty("table-name", "<yourTableName>")
// 可选:指定读取列
.withProperty("scan.columns", "user_id,user_name")
// 可选:过滤条件
.withProperty("scan.filter", "user_id > 100")
.build();
// 初始化读取表的Schema,需与StarRocks表字段匹配,可以只定义部分字段。
TableSchema tableSchema = TableSchema.builder()
.field("user_id", DataTypes.BIGINT())
.field("user_name", DataTypes.STRING())
.build();
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.addSource(StarRocksSource.source(tableSchema, options))
.setParallelism(5)
.print();
env.execute();
}
}
-
读取支持谓词下推,过滤条件会被转换为由StarRocks执行的语句,无需额外配置。
-
StarRocks读取侧不支持Checkpoint机制,作业重启后会重新扫描。
-
常用可选参数:scan.connect.timeout-ms(默认1000)、scan.max-retries(默认1)、scan.params.query-timeout-s(默认600)、scan.params.keep-alive-min(默认10)、scan.params.mem-limit-byte(默认1073741824)、scan.params.batch-rows(默认1000)。数据类型映射请参见StarRocks SQL连接器。
StarRocks Lookup Source
StarRocks连接器没有提供可在DataStream算子中直接调用的Lookup Function,其维表实现仅供Table与SQL的维表JOIN使用。DataStream作业需要维表关联时,请通过Table API声明StarRocks维表,使用Temporal Join完成关联后再转回DataStream:
import org.apache.flink.api.common.typeinfo.Types;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.api.Schema;
import org.apache.flink.table.api.Table;
import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
import org.apache.flink.types.Row;
public class Sample {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tEnv = StreamTableEnvironment.create(env);
// 声明StarRocks维表,维表WITH参数请参见StarRocks SQL连接器。
tEnv.executeSql(
"CREATE TEMPORARY TABLE sr_dim (\n"
+ " user_id BIGINT NOT NULL,\n"
+ " user_name STRING,\n"
+ " PRIMARY KEY (user_id) NOT ENFORCED\n"
+ ") WITH (\n"
+ " 'connector' = 'starrocks',\n"
+ " 'jdbc-url' = 'jdbc:mysql://<fe_host>:9030',\n"
+ " 'scan-url' = '<fe_host>:8030',\n"
+ " 'database-name' = '<yourDatabaseName>',\n"
+ " 'table-name' = '<yourDimTableName>',\n"
+ " 'username' = '<yourUsername>',\n"
+ " 'password' = '<yourPassword>',\n"
+ " 'lookup.cache.ttl-ms' = '5000'\n"
+ ")");
// 将主流注册为视图,并声明处理时间属性。
tEnv.createTemporaryView(
"orders",
env.fromElements(Row.of(1L, 1001L), Row.of(2L, 1002L))
.returns(Types.ROW_NAMED(
new String[] {"order_id", "user_id"},
Types.LONG,
Types.LONG)),
Schema.newBuilder().columnByExpression("proc_time", "PROCTIME()").build());
Table joined = tEnv.sqlQuery(
"SELECT o.order_id, o.user_id, d.user_name\n"
+ "FROM orders AS o\n"
+ "JOIN sr_dim FOR SYSTEM_TIME AS OF o.proc_time AS d\n"
+ "ON o.user_id = d.user_id");
tEnv.toDataStream(joined).print();
env.execute();
}
}
-
维表能力的支持版本、缓存参数(lookup.cache.max-rows、lookup.cache.ttl-ms、lookup.max-retries)与使用限制请参见StarRocks SQL连接器。
StarRocks Sink
StarRocks Sink支持写入自定义Java对象(推荐,可标记UPSERT或DELETE),也可以直接写入CSV或JSON格式的String记录(通过sink.properties.format指定格式)。
以下示例写入自定义对象,演示主键表的UPSERT写入:
import com.starrocks.connector.flink.StarRocksSink;
import com.starrocks.connector.flink.row.sink.StarRocksSinkOP;
import com.starrocks.connector.flink.row.sink.StarRocksSinkRowBuilder;
import com.starrocks.connector.flink.table.sink.StarRocksSinkOptions;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.sink.SinkFunction;
import org.apache.flink.table.api.DataTypes;
import org.apache.flink.table.api.TableSchema;
public class Sample {
/** 自定义数据容器 */
public static class ScoreRecord {
public int userId;
public String userName;
public int score;
public ScoreRecord() {}
public ScoreRecord(int userId, String userName, int score) {
this.userId = userId;
this.userName = userName;
this.score = score;
}
}
// 按Schema顺序把输入对象填充到Object[]中。
private static class ScoreRecordTransformer implements StarRocksSinkRowBuilder<ScoreRecord> {
@Override
public void accept(Object[] internalRow, ScoreRecord record) {
internalRow[0] = record.userId;
internalRow[1] = record.userName;
internalRow[2] = record.score;
// 目标表为主键表时,需设置最后一个元素表示本次写入是UPSERT还是DELETE。
internalRow[internalRow.length - 1] = StarRocksSinkOP.UPSERT.ordinal();
}
}
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.enableCheckpointing(30000);
DataStream<ScoreRecord> stream = env.fromElements(
new ScoreRecord(1, "starrocks-object", 100),
new ScoreRecord(2, "flink-object", 100));
StarRocksSinkOptions options = StarRocksSinkOptions.builder()
// FE的MySQL协议地址,多地址用逗号分隔
.withProperty("jdbc-url", "jdbc:mysql://<fe_host>:9030")
// FE的HTTP地址,多地址用分号分隔
.withProperty("load-url", "<fe_host>:8030")
.withProperty("database-name", "<yourDatabaseName>")
.withProperty("table-name", "<yourTableName>")
.withProperty("username", "<yourUsername>")
.withProperty("password", "<yourPassword>")
.build();
// Schema需与StarRocks表结构匹配,主键列必须声明notNull()。
TableSchema schema = TableSchema.builder()
.field("user_id", DataTypes.INT().notNull())
.field("user_name", DataTypes.STRING())
.field("score", DataTypes.INT())
.primaryKey("user_id")
.build();
SinkFunction<ScoreRecord> sink =
StarRocksSink.sink(schema, options, new ScoreRecordTransformer());
stream.addSink(sink);
env.execute();
}
}
写入的记录会在连接器内部批量缓存并重试,因此作业不能开启对象重用(ExecutionConfig#enableObjectReuse()),否则可能导致写入数据异常。
-
目标表为主键表时,
Object[]的长度为Schema列数加1,最后一个元素用于标记UPSERT或DELETE;非主键表则不需要该元素。 -
主键列在
TableSchema中必须使用notNull()声明,例如DataTypes.INT().notNull()。 -
sink.semantic默认值为at-least-once;使用exactly-once时需要开启Checkpoint。 -
sink.version默认值为V1,支持V1、V2、AUTO,VVR 8.x与11.x均可用。V2使用事务Stream Load接口,要求StarRocks为2.4及以上版本,推荐优先使用。 -
VVR 8.x与11.x的差异:
sink.ignore.update-before、sink.socket.timeout-ms等参数仅VVR 11.x支持;sink.connect.timeout-ms默认值VVR 8.x为1000、VVR 11.x为30000。其余参数差异请参见StarRocks SQL连接器。