StarRocks DataStream connector
The StarRocks DataStream connector supports reading data from and writing data to StarRocks by using the DataStream API. This topic describes the connector dependencies and provides development examples for Source, Lookup Source, and Sink.
Configure connector dependencies
To read and write data by using DataStream, you must use the corresponding DataStream connector to connect to Realtime Compute for Apache Flink. For information about how to set up a DataStream connector, see Integrate and use connectors in DataStream programs. The StarRocks DataStream connector is published to Maven Central. For the Uber JAR that is required for local debugging, see Run and debug connectors locally.
Add the following dependency to the pom.xml file of your Maven project:
<dependency>
<groupId>com.alibaba.ververica</groupId>
<artifactId>ververica-connector-starrocks</artifactId>
<version>${vvr-version}</version>
</dependency>
StarRocks Source
StarRocks Source reads data by batch scanning. It is a bounded read and does not continuously subscribe to incremental changes in StarRocks. To read a change stream, use the StarRocks YAML connector to synchronize data from the upstream database.
The following example builds a StarRocks Source and reads data. The code is the same in VVR 8.x and VVR 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 connection parameters.
StarRocksSourceOptions options = StarRocksSourceOptions.builder()
// HTTP endpoint of the FE. Separate multiple endpoints with commas (,).
.withProperty("scan-url", "<fe_host>:8030")
// MySQL protocol endpoint of the FE.
.withProperty("jdbc-url", "jdbc:mysql://<fe_host>:9030")
.withProperty("username", "<yourUsername>")
.withProperty("password", "<yourPassword>")
.withProperty("database-name", "<yourDatabaseName>")
.withProperty("table-name", "<yourTableName>")
// Optional. The columns to read.
.withProperty("scan.columns", "user_id,user_name")
// Optional. The filter condition.
.withProperty("scan.filter", "user_id > 100")
.build();
// Initialize the schema of the table to read. The schema must match the fields of the
// StarRocks table. You can define only some of the fields.
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();
}
}
-
Reads support predicate pushdown. Filter conditions are converted into statements that StarRocks executes, and no additional configuration is required.
-
The StarRocks read side does not support checkpoints. The scan starts over after the job restarts.
-
Common optional parameters: scan.connect.timeout-ms (default: 1000), scan.max-retries (default: 1), scan.params.query-timeout-s (default: 600), scan.params.keep-alive-min (default: 10), scan.params.mem-limit-byte (default: 1073741824), and scan.params.batch-rows (default: 1000). For data type mappings, see StarRocks SQL connector.
StarRocks Lookup Source
The StarRocks connector does not provide a lookup function that you can call directly in a DataStream operator. Its dimension table implementation is used only for dimension table joins in the Table API and SQL. To perform a dimension table join in a DataStream job, declare the StarRocks dimension table by using the Table API, complete the join with a temporal join, and then convert the result back to a 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);
// Declare the StarRocks dimension table. For the WITH parameters of a dimension table,
// see the StarRocks SQL connector topic.
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"
+ ")");
// Register the main stream as a view and declare the processing time attribute.
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();
}
}
-
For the supported versions of the dimension table feature, the cache parameters (lookup.cache.max-rows, lookup.cache.ttl-ms, and lookup.max-retries), and the usage limits, see StarRocks SQL connector.
StarRocks Sink
StarRocks Sink supports writing custom Java objects, which is recommended because each record can be marked as UPSERT or DELETE. You can also write String records in CSV or JSON format directly, and specify the format with sink.properties.format.
The following example writes custom objects and demonstrates UPSERT writes to a primary key table:
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 {
/** Custom data container. */
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;
}
}
// Fill the input object into Object[] in the order of the schema fields.
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;
// When the destination table is a primary key table, the last element must be set to
// indicate whether this write is an UPSERT or a 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()
// MySQL protocol endpoint of the FE. Separate multiple endpoints with commas (,).
.withProperty("jdbc-url", "jdbc:mysql://<fe_host>:9030")
// HTTP endpoint of the FE. Separate multiple endpoints with semicolons (;).
.withProperty("load-url", "<fe_host>:8030")
.withProperty("database-name", "<yourDatabaseName>")
.withProperty("table-name", "<yourTableName>")
.withProperty("username", "<yourUsername>")
.withProperty("password", "<yourPassword>")
.build();
// The schema must match the StarRocks table structure. Primary key columns must be
// declared with 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();
}
}
Records are buffered in batches and retried inside the connector. Therefore, object reuse (ExecutionConfig#enableObjectReuse()) must be disabled for the job. Otherwise, the written data may be incorrect.
-
When the destination table is a primary key table, the length of
Object[]is the number of schema columns plus one, and the last element marks the record asUPSERTorDELETE. This element is not required for non-primary key tables. -
Primary key columns must be declared with
notNull()inTableSchema, for example,DataTypes.INT().notNull(). -
The default value of
sink.semanticisat-least-once. To useexactly-once, enable checkpointing. -
The default value of
sink.versionisV1.V1,V2, andAUTOare supported in both VVR 8.x and VVR 11.x.V2uses the transactional Stream Load interface and requires StarRocks 2.4 or later. V2 is recommended. -
Differences between VVR 8.x and VVR 11.x: parameters such as
sink.ignore.update-beforeandsink.socket.timeout-msare supported only in VVR 11.x. The default value ofsink.connect.timeout-msis 1000 in VVR 8.x and 30000 in VVR 11.x. For other parameter differences, see StarRocks SQL connector.