Apache Flink
Continuously load data from Apache Flinkยฎ
StarRocks provides a self-developed connector named StarRocks Connector for Apache Flinkยฎ (Flink connector for short) to help you load data into a StarRocks table by using Flink. The basic principle is to accumulate the data and then load it all at a time into StarRocks through STREAM LOAD.
The Flink connector supports DataStream API, Table API & SQL, and Python API. It has a higher and more stable performance than flink-connector-jdbc provided by Apache Flinkยฎ.
NOTICE
Loading data into StarRocks tables with Flink connector needs SELECT and INSERT privileges on the target StarRocks table. If you do not have these privileges, follow the instructions provided in GRANT to grant these privileges to the user that you use to connect to your StarRocks cluster.
Version requirementsโ
| Connector | Flink | StarRocks | Java | Scala |
|---|---|---|---|---|
| 1.2.14 | 1.16,1.17,1.18,1.19,1.20 | 2.1 and later | 8 | 2.11,2.12 |
| 1.2.12 | 1.16,1.17,1.18,1.19,1.20 | 2.1 and later | 8 | 2.11,2.12 |
| 1.2.11 | 1.15,1.16,1.17,1.18,1.19,1.20 | 2.1 and later | 8 | 2.11,2.12 |
| 1.2.10 | 1.15,1.16,1.17,1.18,1.19 | 2.1 and later | 8 | 2.11,2.12 |
Obtain Flink connectorโ
You can obtain the Flink connector JAR file in the following ways:
- Directly download the compiled Flink connector JAR file.
- Add the Flink connector as a dependency in your Maven project and then download the JAR file.
- Compile the source code of the Flink connector into a JAR file by yourself.
The naming format of the Flink connector JAR file is as follows:
-
Since Flink 1.15, it's
flink-connector-starrocks-${connector_version}_flink-${flink_version}.jar. For example, if you install Flink 1.15 and you want to use Flink connector 1.2.7, you can useflink-connector-starrocks-1.2.7_flink-1.15.jar. -
Prior to Flink 1.15, it's
flink-connector-starrocks-${connector_version}_flink-${flink_version}_${scala_version}.jar. For example, if you install Flink 1.14 and Scala 2.12 in your environment, and you want to use Flink connector 1.2.7, you can useflink-connector-starrocks-1.2.7_flink-1.14_2.12.jar.
NOTICE
In general, the latest version of the Flink connector only maintains compatibility with the three most recent versions of Flink.
Download the compiled Jar fileโ
Directly download the corresponding version of the Flink connector Jar file from the Maven Central Repository.
Maven Dependencyโ
In your Maven project's pom.xml file, add the Flink connector as a dependency according to the following format. Replace flink_version, scala_version, and connector_version with the respective versions.
-
In Flink 1.15 and later
<dependency>
<groupId>com.starrocks</groupId>
<artifactId>flink-connector-starrocks</artifactId>
<version>${connector_version}_flink-${flink_version}</version>
</dependency> -
In versions earlier than Flink 1.15
<dependency>
<groupId>com.starrocks</groupId>
<artifactId>flink-connector-starrocks</artifactId>
<version>${connector_version}_flink-${flink_version}_${scala_version}</version>
</dependency>
Compile by yourselfโ
-
Download the Flink connector source code.
-
Execute the following command to compile the source code of Flink connector into a JAR file. Note that
flink_versionis replaced with the corresponding Flink version.sh build.sh <flink_version>For example, if the Flink version in your environment is 1.15, you need to execute the following command:
sh build.sh 1.15 -
Go to the
target/directory to find the Flink connector JAR file, such asflink-connector-starrocks-1.2.7_flink-1.15-SNAPSHOT.jar, generated upon compilation.
NOTE
The name of Flink connector which is not formally released contains the
SNAPSHOTsuffix.
Connect to StarRocksโ
Set the StarRocks connection options to addresses of the FE:
jdbc-url:jdbc:mysql://<fe_host>:<fe_query_port>. The query port defaults to9030.load-url:<fe_host>:<fe_http_port>. The HTTP port defaults to8030.
For <fe_host>, use an address that stays the same when an FE restarts or is replaced, such as a DNS name or a load balancer in front of the FEs, not the IP address of an FE. If the address stops leading to an FE, the Flink job can no longer load data.
The Flink TaskManagers must be able to reach these FE ports, and also the HTTP port (default 8040) of every BE or CN, because the FE redirects each load request to a BE or CN.
Optionsโ
General Optionsโ
connectorโ
- Required: Yes
- Default value: NONE
- Description: The connector that you want to use. The value must be "starrocks".
jdbc-urlโ
- Required: Yes
- Default value: NONE
- Description: The address that is used to connect to the MySQL server of the FE. You can specify multiple addresses, which must be separated by a comma (,). Format:
jdbc:mysql://<fe_host1>:<fe_query_port1>,<fe_host2>:<fe_query_port2>,<fe_host3>:<fe_query_port3>.
load-urlโ
- Required: Yes
- Default value: NONE
- Description: The address that is used to connect to the HTTP server of the FE. You can specify multiple addresses, which must be separated by a semicolon (;). Format:
<fe_host1>:<fe_http_port1>;<fe_host2>:<fe_http_port2>.
database-nameโ
- Required: Yes
- Default value: NONE
- Description: The name of the StarRocks database into which you want to load data.
table-nameโ
- Required: Yes
- Default value: NONE
- Description: The name of the table that you want to use to load data into StarRocks.
usernameโ
- Required: Yes
- Default value: NONE
- Description: The username of the account that you want to use to load data into StarRocks. The account needs SELECT and INSERT privileges on the target StarRocks table.
passwordโ
- Required: Yes
- Default value: NONE
- Description: The password of the preceding account.
sink.versionโ
- Required: No
- Default value: AUTO
- Description: The interface used to load data. This parameter is supported from Flink connector version 1.2.4 onwards. Valid Values:
V1: Use Stream Load interface to load data. Connectors before 1.2.4 only support this mode.V2: Use Stream Load transaction interface to load data. It requires StarRocks to be at least version 2.4. RecommendsV2because it optimizes the memory usage and provides a more stable exactly-once implementation.AUTO: If the version of StarRocks supports transaction Stream Load, will chooseV2automatically, otherwise chooseV1
sink.label-prefixโ
- Required: No
- Default value: NONE
- Description: The label prefix used by Stream Load. Recommend to configure it if you are using exactly-once with connector 1.2.8 and later. See exactly-once usage notes.
sink.semanticโ
- Required: No
- Default value: at-least-once
- Description: The semantic guaranteed by sink. Valid values: at-least-once and exactly-once.
sink.buffer-flush.max-bytesโ
- Required: No
- Default value: 94371840(90M)
- Description: The maximum size of data that can be accumulated in memory before being sent to StarRocks at a time. The maximum value ranges from 64 MB to 10 GB. Setting this parameter to a larger value can improve loading performance but may increase loading latency. This parameter only takes effect when
sink.semanticis set toat-least-once. Ifsink.semanticis set toexactly-once, the data in memory is flushed when a Flink checkpoint is triggered. In this circumstance, this parameter does not take effect.
sink.buffer-flush.max-rowsโ
- Required: No
- Default value: 500000
- Description: The maximum number of rows that can be accumulated in memory before being sent to StarRocks at a time. This parameter is available only when
sink.versionisV1andsink.semanticisat-least-once. Valid values: 64000 to 5000000.
sink.buffer-flush.interval-msโ
- Required: No
- Default value: 300000
- Description: The interval at which data is flushed. This parameter is available only when
sink.semanticisat-least-once. Unit: ms. Valid value range:- For versions earlier than v1.2.14: [1000, 3600000]
- For v1.2.14 and later: (0, 3600000].
sink.max-retriesโ
- Required: No
- Default value: 3
- Description: The number of times that the system retries to perform the Stream Load job. This parameter is available only when you set
sink.versiontoV1. Valid values: 0 to 10.
sink.connect.timeout-msโ
- Required: No
- Default value: 30000
- Description: The timeout for establishing HTTP connection. Valid values: 100 to 60000. Unit: ms. Before Flink connector v1.2.9, the default value is
1000.
sink.socket.timeout-msโ
- Required: No
- Default value: -1
- Description: Supported since 1.2.10. The time duration for which the HTTP client waits for data. Unit: ms. The default value
-1means there is no timeout.
sink.sanitize-error-logโ
- Required: No
- Default value: false
- Description: Supported since 1.2.12. Whether to sanitize sensitive data in the error log for production security. When this item is set to
true, sensitive row data and column values in Stream Load error logs are redacted in both the connector and SDK logs. The value defaults tofalsefor backward compatibility.
sink.wait-for-continue.timeout-msโ
- Required: No
- Default value: 10000
- Description: Supported since 1.2.7. The timeout for waiting response of HTTP 100-continue from the FE. Valid values:
3000to60000. Unit: ms
sink.ignore.update-beforeโ
- Required: No
- Default value: true
- Description: Supported since version 1.2.8. Whether to ignore
UPDATE_BEFORErecords from Flink when loading data to Primary Key tables. If this parameter is set to false, the record is treated as a delete operation to StarRocks table.
sink.parallelismโ
- Required: No
- Default value: NONE
- Description: The parallelism of loading. Only available for Flink SQL. If this parameter is not specified, Flink planner decides the parallelism. In the scenario of multi-parallelism, users need to guarantee data is written in the correct order.
sink.properties.*โ
- Required: No
- Default value: NONE
- Description: The parameters that are used to control Stream Load behavior. For example, the parameter
sink.properties.formatspecifies the format used for Stream Load, such as CSV or JSON. For a list of supported parameters and their descriptions, see STREAM LOAD.
sink.properties.formatโ
- Required: No
- Default value: csv
- Description: The format used for Stream Load. The Flink connector will transform each batch of data to the format before sending them to StarRocks. Valid values:
csvandjson.
sink.properties.column_separatorโ
- Required: No
- Default value: \t
- Description: The column separator for CSV-formatted data.
sink.properties.row_delimiterโ
- Required: No
- Default value: \n
- Description: The row delimiter for CSV-formatted data.
sink.properties.max_filter_ratioโ
- Required: No
- Default value: 0
- Description: The maximum error tolerance of the Stream Load. It's the maximum percentage of data records that can be filtered out due to inadequate data quality. Valid values:
0to1. Default value:0. See Stream Load for details.
sink.properties.partial_updateโ
- Required: NO
- Default value:
FALSE - Description: Whether to use partial updates. Valid values:
TRUEandFALSE. Default value:FALSE, indicating to disable this feature.
sink.properties.partial_update_modeโ
- Required: NO
- Default value:
row - Description: Specifies the mode for partial updates. Valid values:
rowandcolumn.- The value
row(default) means partial updates in row mode, which is more suitable for real-time updates with many columns and small batches. - The value
columnmeans partial updates in column mode, which is more suitable for batch updates with few columns and many rows. In such scenarios, enabling the column mode offers faster update speeds. For example, in a table with 100 columns, if only 10 columns (10% of the total) are updated for all rows, the update speed of the column mode is 10 times faster.
- The value
sink.properties.strict_modeโ
- Required: No
- Default value: false
- Description: Specifies whether to enable the strict mode for Stream Load. It affects the loading behavior when there are unqualified rows, such as inconsistent column values. Valid values:
trueandfalse. Default value:false. See Stream Load for details.
sink.properties.compressionโ
- Required: No
- Default value: NONE
- Description: The compression algorithm used for Stream Load. Valid values:
lz4_frame. Compression for the JSON format requires Flink connector 1.2.10+ and StarRocks v3.2.7+. Compression for the CSV format only requires Flink connector 1.2.11+.
sink.properties.prepared_timeoutโ
- Required: No
- Default value: NONE
- Description: Supported since 1.2.12 and only effective when
sink.versionis set toV2. Requires StarRocks 3.5.4 or later. Sets the timeout in seconds for the Transaction Stream Load phase fromPREPAREDtoCOMMITTED. Typically, only needed for exactly-once; at-least-once usually does not require setting this (the connector defaults to 300s). If not set in exactly-once, StarRocks FE configurationprepared_transaction_default_timeout_second(default 86400s) applies. See StarRocks Transaction timeout management.
sink.publish-timeout.msโ
- Required: No
- Default value: -1
- Description: Supported since 1.2.14 and only effective when
sink.versionis set toV2. Timeout in milliseconds for the Publish phase. If a transaction stays in COMMITTED status longer than this timeout, the system will consider it as successful. The default value-1means using StarRocks server-side default behavior. When Merge Commit is enabled, the default timeout is 10000 ms.
Merge Commit optionsโ
Supported from v1.2.14 onwards. Merge Commit allows the system to merge data from multiple subtasks into a single Stream Load transaction for better performance. You can enable this feature by setting sink.properties.enable_merge_commit to true. For more details about the merge commit feature in StarRocks, see Merge Commit parameters.
The following Stream Load properties are used to control the Merge Commit behavior: