StarRocks YAML connector
This topic describes how to use the StarRocks connector to synchronize data in YAML-based data ingestion jobs.
Background
StarRocks is an MPP (Massively Parallel Processing) data warehouse designed for real-time analytics. It is compatible with the MySQL protocol, uses a distributed architecture, and supports elastic cluster scaling and parallel computing. The StarRocks YAML connector writes upstream data records and schema changes to StarRocks, and supports both the community edition of StarRocks and the fully managed EMR Serverless StarRocks from Alibaba Cloud. The following table describes what the StarRocks YAML connector supports.
|
Category |
Description |
|
Supported types |
Data ingestion sink |
|
Execution mode |
streaming mode and batch mode |
|
Data format |
JSON |
|
Connector-specific metrics |
None |
|
API types |
YAML |
|
Support for updates/deletions in sink tables |
Yes |
The YAML connector currently supports only at-least-once semantics. Even if you explicitly set sink.semantic: exactly-once, it is overridden to at-least-once without an error. For exactly-once semantics, see StarRocks SQL connector.
Features
-
Automatic database and table creation.
If an upstream database or table does not exist in the downstream StarRocks instance, the connector creates it automatically. You can use the
table.create.properties.*parameter to configure options for automatic table creation. -
Schema change synchronization.
The StarRocks connector automatically applies CreateTableEvent, AddColumnEvent, and DropColumnEvent events to the downstream database.
-
VVR 11.1 and later supports compatible column type changes. For more information, see ALTER TABLE | StarRocks.
Usage notes
-
Each synchronized table must have a primary key. For tables without a primary key, you must specify one in the
transformblock to write data downstream. For example:transform: - source-table: ... primary-keys: id, ... -
For automatically created tables, the bucket key is the same as the primary key, and the table cannot have a partition key.
-
When synchronizing schema changes, new columns can only be appended to the end of existing columns. In the default Lenient schema change mode, insertions at other positions are automatically moved to the end.
-
If you use a StarRocks version earlier than 2.5.7, you must explicitly specify the number of buckets with the
table.create.num-bucketsparameter. StarRocks 2.5.7 and later can automatically determine an appropriate number of buckets. -
If you use StarRocks 3.2 or later, we recommend that you enable the
table.create.properties.fast_schema_evolutionoption to accelerate schema changes. -
If you use Ververica Runtime (VVR) 11.9 or later and StarRocks 3.3.2 or later, you can synchronize upstream column rename events to the downstream table.
-
Streaming issues may occur when you use CDC YAML for data ingestion into EMR Serverless StarRocks. You can use one of the following workarounds:
-
Use the Flink SQL StarRocks connector and set the
sink.version=V1parameter. -
Enable the FE parameter
emr_internal_redirect. -
Use a StarRocks Private Zone domain name instead of an SLB.
-
Syntax
source:
...
sink:
type: starrocks
name: StarRocks Sink
jdbc-url: jdbc:mysql://127.0.0.1:9030
load-url: 127.0.0.1:8030
username: root
password: pass
sink.buffer-flush.interval-ms: 5000 # Set the data flush interval.
Configuration
|
Parameter |
Description |
Type |
Required |
Default |
Remarks |
|
|
Specifies the sink connector type. |
String |
Yes |
— |
Set to |
|
|
The display name of the sink. |
String |
No |
— |
— |
|
|
The JDBC URL for the database connection. |
String |
Yes |
— |
Supports multiple addresses separated by commas ( |
|
|
The HTTP URL of an FE node for Stream Load. |
String |
Yes |
— |
Supports multiple addresses separated by semicolons ( |
|
|
The username for the StarRocks connection. |
String |
Yes |
— |
This user must have at least SELECT and INSERT permissions on the target table. You can grant the required permissions with the StarRocks GRANT command. |
|
|
The password for the StarRocks connection. |
String |
Yes |
— |
— |
|
|
The delivery semantics for data writes. |
String |
No |
at-least-once |
Only |
|
|
The label prefix for Stream Load jobs. |
String |
No |
— |
The value can contain only English letters, digits, hyphens ( |
|
|
The timeout for establishing an HTTP connection. |
Integer |
No |
30000 |
Unit: milliseconds. The value must be between 100 and 60000. |
|
|
The timeout for waiting for a 100 Continue response from the server. |
Integer |
No |
30000 |
Unit: milliseconds. The value must be between 3000 and 600000. |
|
|
The maximum size of the in-memory cache, in bytes, before a flush is triggered. |
Long |
No |
94371840 |
Unit: bytes. The value must be between 64 MB and 10 GB. Note
|
|
|
The maximum number of rows in the in-memory cache before a flush is triggered. |
Long |
No |
500000 |
The value must be between 1,000 and 5,000,000. |
|
|
The time interval between flushes for each table's buffer. |
Long |
No |
300000 |
Unit: milliseconds. Note
For jobs that synchronize small amounts of data, reduce this value to avoid long delays before data is persisted. |
|
|
The maximum number of retries. |
Long |
No |
3 |
The value must be between 0 and 1000. |
|
|
The frequency at which the connector checks whether to flush the buffer. |
Long |
No |
50 |
Unit: milliseconds. |
|
|
The number of threads used for Stream Load. |
Integer |
No |
2 |
— |
|
|
Specifies whether to use the Stream Load transaction interface for data ingestion. |
Boolean |
No |
true |
This option takes effect only if the database supports it. |
|
|
Specifies whether to ignore update-before records in update operations. |
Boolean |
No |
true |
When the primary key is changed via the Transform module (for example, when Only Ververica Runtime (VVR) 11.8 or later supports this parameter. |
|
|
Specifies whether to ignore delete records. |
Boolean |
No |
false |
If you set this parameter to Only Ververica Runtime (VVR) 11.8 or later supports this parameter. |
|
|
Additional properties for the sink. |
String |
No |
— |
For supported properties, see STREAM LOAD. |
|
|
The number of buckets for automatically created tables. |
Integer |
No |
— |
|
|
|
Additional properties for automatic table creation. |
String |
No |
— |
For example, you can pass |
|
|
The timeout for schema change operations. |
Duration |
No |
30 min |
Must be an integer number of seconds. Note
If a schema change operation exceeds this limit, the job fails. |
|
|
The number of bytes to allocate for each Unicode character. |
Integer |
No |
3 |
In CDC, the length of a VARCHAR type is measured in characters, whereas in StarRocks, the length of a VARCHAR type is measured in bytes. In most cases, a Unicode character does not exceed 3 bytes after UTF-8 encoding. However, some rare characters and emoji symbols may occupy 4 or more bytes. |
|
|
The HTTP client timeout for flushing data to StarRocks. |
Long |
No |
-1 |
The HTTP client timeout, in milliseconds, for sending Stream Load requests when data is flushed to StarRocks. A value of -1 uses the system default, which means no timeout. Only Ververica Runtime (VVR) 11.8 or later supports this parameter. |
|
|
The timeout for closing the sink. |
Long |
No |
60000 |
The timeout, in milliseconds, for waiting for the flush queue to finish when the job closes. Only Ververica Runtime (VVR) 11.8 or later supports this parameter. |
Reuse a built-in catalog
VVR 11.5 and later lets you reference a built-in StarRocks catalog created on the Data Management page directly in a Flink CDC data ingestion job. This simplifies configuration by reducing the number of properties you need to set manually.
sink:
type: starrocks
using.built-in-catalog: starrocks_catalog
Data ingestion jobs can automatically reuse the following StarRocks catalog options:
-
jdbc-url
-
http-url
-
username
-
password
-
table.num-buckets
To override these values, you can explicitly set the corresponding YAML options, which take precedence.
Type mapping
StarRocks does not support all CDC YAML types. Writing an unsupported type to the sink causes the job to fail. You can use the CAST built-in function in a transform to convert unsupported data, or use a projection statement to remove it from the result table. For more information, see Develop a Flink CDC data ingestion job.
|
CDC type |
StarRocks type |
Remarks |
|
TINYINT |
TINYINT |
— |
|
SMALLINT |
SMALLINT |
|
|
INT |
INT |
|
|
BIGINT |
BIGINT |
|
|
FLOAT |
FLOAT |
|
|
DOUBLE |
DOUBLE |
|
|
BOOLEAN |
BOOLEAN |
|
|
DATE |
DATE |
|
|
TIMESTAMP |
DATETIME |
|
|
TIMESTAMP_LTZ |
DATETIME |
|
|
DECIMAL(p, s) |
DECIMAL(p, s) |
Because StarRocks does not support DECIMAL for a primary key, the connector automatically converts an upstream DECIMAL primary key column to VARCHAR in the synchronized StarRocks schema. |
|
CHAR(n) (n <= 85) |
CHAR(n × 3) |
CDC measures length in characters, while StarRocks uses bytes. The connector multiplies the length by 3 to account for multi-byte UTF-8 characters. Note
The maximum length of the StarRocks CHAR type is 255. Therefore, only CDC CHAR types with a length up to 85 are mapped to the StarRocks CHAR type. Note
You can set the |
|
CHAR(n) (n > 85) |
VARCHAR(n × 3) |
CDC measures length in characters, while StarRocks uses bytes. The connector multiplies the length by 3 to account for multi-byte UTF-8 characters. Note
CDC measures length in characters, while StarRocks uses bytes. The connector multiplies the length by 3. Since the result exceeds the 255-byte limit for the StarRocks CHAR type, it is mapped to VARCHAR. Note
You can set the |
|
VARCHAR(n) |
VARCHAR(n × 3) |
CDC measures length in characters, while StarRocks uses bytes. The connector multiplies the length by 3 to account for multi-byte UTF-8 characters. Note
You can set the |
|
BINARY(n) |
BINARY(n+2) |
Two bytes of padding are added to ensure data integrity. |
|
VARBINARY(n) |
VARBINARY(n+1) |
One byte of padding is added to ensure data integrity. |
Schema change
A CDC YAML pipeline job uses different strategies to handle schema changes, which are configurable with the pipeline-level schema.change.behavior parameter. Valid values are IGNORE, LENIENT, TRY_EVOLVE, EVOLVE, and EXCEPTION. Default value: LENIENT. Whether an event is applied is also affected by the include.schema.changes and exclude.schema.changes parameters of the sink.
The LENIENT and EVOLVE strategies involve schema changes. The following sections describe how different schema change events are handled.
The following descriptions apply to regular one-to-one synchronization. For many-to-one routing, the merged schema is derived first and the events are then normalized. Do not assume that a drop operation on any shard table directly drops the merged destination table.
Supported events
-
CREATE TABLE EVENT
NoteIf the downstream StarRocks table already exists, the connector does not attempt to create it again. Ensure that the downstream table schema is compatible with the upstream schema.
-
ADD COLUMN EVENT
NoteStarRocks requires primary key columns to appear first in a table. Any new columns must be added after them.
-
ALTER COLUMN TYPE EVENT
NoteFor supported schema change paths, refer to the official StarRocks documentation.
-
RENAME COLUMN EVENT
NoteRequires Ververica Runtime (VVR) 11.9 or later, and StarRocks 3.3.2 or later.
-
DROP COLUMN EVENT
-
TRUNCATE TABLE EVENT
-
DROP TABLE EVENT
LENIENT (default)
In LENIENT mode, schema changes are handled as follows:
-
Add a nullable column: The corresponding column is automatically added to the end of the sink table schema, and the data of the new column is synchronized automatically.
-
Add a non-nullable column: A corresponding column is automatically added to the end of the sink table schema and is set to nullable. For data that existed before the column was added, the value is automatically set to NULL.
-
Drop a column: The column is not dropped from the sink table. If the original column was non-nullable, it is changed to nullable. If it was already nullable, no change is made.
-
Rename a column: This operation is treated as adding a new column while keeping the old one. The renamed nullable column is added to the end of the sink table. If the original column was non-nullable, it is changed to nullable.
-
Reorder columns: Ignored. The change is not synchronized to the downstream table.
-
Change a column type: Not automatically ignored by LENIENT mode. The change is still passed to StarRocks, and only compatible change paths are supported. For the allowed paths, see ALTER TABLE. For VVR 11.1 and later, the supported scope remains the same as that described in the existing documentation.
-
Drop a table or truncate a table: Excluded by default. The change is not synchronized to the downstream.
Excluding drop-table and truncate-table events by default is a default entry that the YAML parser adds when exclude.schema.changes is not configured. It is not unconditional protection. Explicitly configuring the exclusion list, including an empty list, replaces these two default exclusions, so keep the events that you want to prohibit. LENIENT also does not guarantee that failed DDL executions are ignored.
EVOLVE
In EVOLVE mode, schema changes are handled as follows:
-
Add a column: Supported. The StarRocks executor ignores the upstream insertion position and appends the column to the end of the existing columns.
-
Drop a column: Supported. The column is actually dropped from the sink table.
-
Rename a column: Supported (requires StarRocks 3.3.2 or later). The name of the downstream column is changed.
-
Change a column type: Supported. Only compatible change paths are supported.
-
Reorder columns: Supported. The primary key columns must still come first and keep their original order.
-
Drop a table or truncate a table: Actually executed when the events are not excluded.
In EVOLVE mode, if you perform a stateless restart without dropping the sink table, the upstream data and the sink table schema may become inconsistent, which causes the job to fail. You must then adjust the downstream table schema manually.
Code examples
The following examples show configurations for common scenarios.
Synchronize a single table
Synchronize a single MySQL table to StarRocks. If the destination database and table do not exist, the connector automatically creates a primary key table.
pipeline:
name: MySQL to StarRocks Pipeline
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: <yourUsername>
password: ${secret_values.mysql_password}
tables: test_db.test_source_table
server-id: 5401-5499
# (Optional) Synchronize data from newly added tables during the incremental phase without restarting the job.
scan.binlog.newly-added-table.enabled: true
# (Optional) Synchronize table and column comments to the destination.
include-comments.enabled: true
# (Optional) Deserialize binlogs only for captured tables to improve read performance.
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: starrocks
name: StarRocks Sink
jdbc-url: jdbc:mysql://<yourFeHostname>:9030
load-url: <yourFeHostname>:8030
username: <yourUsername>
password: ${secret_values.starrocks_password}
# (Optional) For jobs with a small amount of data, reduce the flush interval to prevent long persistence delays. Default: 300000, or 5 minutes.
sink.buffer-flush.interval-ms: 5000
# (Optional) If the upstream character set is utf8mb4, set this parameter to 4 to prevent text truncation. Default: 3.
unicode-char.max-bytes: 4
# (Optional) The number of buckets for automatically created tables. This parameter is required for StarRocks versions earlier than 2.5.7. Later versions can infer the value automatically.
table.create.num-buckets: 8
# (Optional) The number of replicas for automatically created tables. Configure this value based on your cluster.
table.create.properties.replication_num: 3
# (Optional) For StarRocks 3.2 and later, enable this option to accelerate schema changes.
table.create.properties.fast_schema_evolution: true
# Note: If you use a transform to change the primary key, you must also set sink.ignore.update-before: false.
# Otherwise, rows associated with the old primary key remain in the destination.
pipeline:
name: MySQL to StarRocks Pipeline
Synchronize an entire database
Synchronize all tables in a MySQL database to StarRocks at the same time. The connector automatically creates the destination database and primary key tables, so you do not need to create each table in advance.
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: <yourUsername>
password: ${secret_values.mysql_password}
# Use a regular expression to match all tables in the database. To match multiple databases, separate patterns with commas.
tables: test_db.\.*
server-id: 5401-5499
# (Optional) Synchronize data from newly added tables during the incremental phase without restarting the job.
scan.binlog.newly-added-table.enabled: true
# (Optional) Synchronize table and column comments to the destination.
include-comments.enabled: true
sink:
type: starrocks
name: StarRocks Sink
jdbc-url: jdbc:mysql://<yourFeHostname>:9030
load-url: <yourFeHostname>:8030
username: <yourUsername>
password: ${secret_values.starrocks_password}
# (Optional) For jobs with a small amount of data, reduce the flush interval to prevent long persistence delays. Default: 300000, or 5 minutes.
sink.buffer-flush.interval-ms: 5000
# (Optional) If the upstream character set is utf8mb4, set this parameter to 4 to prevent text truncation. Default: 3.
unicode-char.max-bytes: 4
# (Optional) The number of buckets for automatically created tables. This parameter is required for StarRocks versions earlier than 2.5.7. Later versions can infer the value automatically.
table.create.num-buckets: 8
# (Optional) The number of replicas for automatically created tables. Configure this value based on your cluster.
table.create.properties.replication_num: 3
# (Optional) For StarRocks 3.2 and later, enable this option to accelerate schema changes.
table.create.properties.fast_schema_evolution: true
pipeline:
name: MySQL to StarRocks Pipeline
Exclude specific tables during full database synchronization
When synchronizing an entire database, use a regular expression to skip tables that you do not want to synchronize to the destination, such as temporary or sensitive tables.
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: <yourUsername>
password: ${secret_values.mysql_password}
tables: test_db.\.*
# Tables that match this regular expression are not synchronized.
tables.exclude: test_db.tmp_.\*
server-id: 5401-5499
sink:
type: starrocks
name: StarRocks Sink
jdbc-url: jdbc:mysql://<yourFeHostname>:9030
load-url: <yourFeHostname>:8030
username: <yourUsername>
password: ${secret_values.starrocks_password}
# (Optional) The load interface version. V2 requires StarRocks 2.4 or later. For EMR Serverless StarRocks, use V1 if streaming issues occur.
sink.version: V2
# (Optional) For jobs with a small amount of data, reduce the flush interval to prevent long persistence delays. Default: 300000, or 5 minutes.
sink.buffer-flush.interval-ms: 5000
# (Optional) The number of buckets for automatically created tables. This parameter is required for StarRocks versions earlier than 2.5.7.
table.create.num-buckets: 8
# (Optional) For StarRocks 3.2 and later, enable this option to accelerate schema changes.
table.create.properties.fast_schema_evolution: true
pipeline:
name: MySQL to StarRocks Pipeline
Synchronize to specified databases and tables
If the destination StarRocks database or table name must differ from the upstream name, such as when writing to an ODS-layer database, use a route to rename it.
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: <yourUsername>
password: ${secret_values.mysql_password}
tables: test_db.\.*
server-id: 5401-5499
sink:
type: starrocks
name: StarRocks Sink
jdbc-url: jdbc:mysql://<yourFeHostname>:9030
load-url: <yourFeHostname>:8030
username: <yourUsername>
password: ${secret_values.starrocks_password}
# (Optional) For jobs with a small amount of data, reduce the flush interval to prevent long persistence delays. Default: 300000, or 5 minutes.
sink.buffer-flush.interval-ms: 5000
# (Optional) The number of buckets for automatically created tables. This parameter is required for StarRocks versions earlier than 2.5.7.
table.create.num-buckets: 8
route:
# Synchronize all tables in the MySQL test_db database to the StarRocks test_db2 database without changing table names.
# <> is a placeholder that is replaced by the matched source table name.
- source-table: test_db.\.*
sink-table: test_db2.<>
replace-symbol: <>
pipeline:
name: MySQL to StarRocks Pipeline
Merge sharded tables
Merge multiple sharded tables with identical schemas into a single StarRocks table for unified queries and analytics.
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: <yourUsername>
password: ${secret_values.mysql_password}
# Match all sharded tables, such as user_0 and user_1.
tables: test_db.user\.*
server-id: 5401-5499
sink:
type: starrocks
name: StarRocks Sink
jdbc-url: jdbc:mysql://<yourFeHostname>:9030
load-url: <yourFeHostname>:8030
username: <yourUsername>
password: ${secret_values.starrocks_password}
# (Optional) For jobs with a small amount of data, reduce the flush interval to prevent long persistence delays. Default: 300000, or 5 minutes.
sink.buffer-flush.interval-ms: 5000
# (Optional) For merged tables, explicitly specify the number of buckets based on the total data volume.
table.create.num-buckets: 8
route:
# Merge all sharded tables into the StarRocks test_db.user table.
- source-table: test_db.user\.*
sink-table: test_db.user
pipeline:
name: MySQL to StarRocks Pipeline
Enable EVOLVE mode
By default, LENIENT mode does not synchronize schema changes such as dropping a column, dropping a table, or truncating a table to the destination. If you require strict schema synchronization, enable EVOLVE mode. This mode has significant limitations, so review the following notes before you use it.
Limitations and usage notes
-
Column renaming is not supported. The job fails if an upstream column rename event occurs.
-
Dropping a column or table, or truncating a table, is applied to the destination. An accidental upstream operation directly affects the destination table. The default LENIENT mode is safer because it does not synchronize drop-table or truncate-table events.
-
If you restart the job without state and do not delete the sink table, an upstream and sink schema mismatch may cause the job to fail. You must manually adjust the downstream table schema.
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: <yourUsername>
password: ${secret_values.mysql_password}
tables: test_db.test_source_table
server-id: 5401-5499
sink:
type: starrocks
name: StarRocks Sink
jdbc-url: jdbc:mysql://<yourFeHostname>:9030
load-url: <yourFeHostname>:8030
username: <yourUsername>
password: ${secret_values.starrocks_password}
pipeline:
name: MySQL to StarRocks Pipeline
# Enable EVOLVE mode to strictly synchronize schema changes. The job fails on unsupported changes, such as column renaming.
schema.change.behavior: evolve