StarRocks DataStream连接器

更新时间:
复制 MD 格式

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连接器。

相关文档