StarRocks YAML connector

Updated at:

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

Note

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 transform block 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-buckets parameter. 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_evolution option 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=V1 parameter.

    • 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

type

Specifies the sink connector type.

String

Yes

—

Set to starrocks.

name

The display name of the sink.

String

No

—

—

jdbc-url

The JDBC URL for the database connection.

String

Yes

—

Supports multiple addresses separated by commas (,). Example: jdbc:mysql://fe_host1:fe_query_port1,fe_host2:fe_query_port2,fe_host3:fe_query_port3.

load-url

The HTTP URL of an FE node for Stream Load.

String

Yes

—

Supports multiple addresses separated by semicolons (;). Example: fe_host1:fe_http_port1;fe_host2:fe_http_port2.

username

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.

password

The password for the StarRocks connection.

String

Yes

—

—

sink.semantic

The delivery semantics for data writes.

String

No

at-least-once

Only at-least-once is supported. Explicitly setting this parameter to exactly-once does not cause an error, but the value is automatically reset to at-least-once. To use exactly-once semantics, use the Flink SQL StarRocks connector. For more information, see the SQL section of this topic.

sink.label-prefix

The label prefix for Stream Load jobs.

String

No

—

The value can contain only English letters, digits, hyphens (-), and underscores (_). Other characters may cause the load to fail.

sink.connect.timeout-ms

The timeout for establishing an HTTP connection.

Integer

No

30000

Unit: milliseconds. The value must be between 100 and 60000.

sink.wait-for-continue.timeout-ms

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.

sink.buffer-flush.max-bytes

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
  • This cache size is shared by all tables. When the buffer is full, the connector selects several tables to flush.

  • Setting a larger value can improve throughput but may increase ingestion latency.

sink.buffer-flush.max-rows

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.

sink.buffer-flush.interval-ms

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.

sink.max-retries

The maximum number of retries.

Long

No

3

The value must be between 0 and 1000.

sink.scan-frequency.ms

The frequency at which the connector checks whether to flush the buffer.

Long

No

50

Unit: milliseconds.

sink.io.thread-count

The number of threads used for Stream Load.

Integer

No

2

—

sink.at-least-once.use-transaction-stream-load

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.

sink.ignore.update-before

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 primary-keys specifies a primary key different from the upstream), set sink.ignore.update-before to false. Otherwise, rows corresponding to the old primary key are not deleted, resulting in stale data.

Only Ververica Runtime (VVR) 11.8 or later supports this parameter.

sink.ignore.delete

Specifies whether to ignore delete records.

Boolean

No

false

If you set this parameter to true, delete records are filtered out and are not written to StarRocks. Use this setting to retain historical data in the sink and synchronize only insert and update operations.

Only Ververica Runtime (VVR) 11.8 or later supports this parameter.

sink.properties.*

Additional properties for the sink.

String

No

—

For supported properties, see STREAM LOAD.

table.create.num-buckets

The number of buckets for automatically created tables.

Integer

No

—

table.create.properties.*

Additional properties for automatic table creation.

String

No

—

For example, you can pass 'table.create.properties.fast_schema_evolution' = 'true' to enable fast schema change. For details, see the StarRocks documentation.

table.schema-change.timeout

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.

unicode-char.max-bytes

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.

sink.socket.time

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.

sink.close.eof-timeout-ms

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

Note

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 unicode-char.max-bytes parameter to allocate more bytes for each Unicode character.

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 unicode-char.max-bytes parameter to allocate more bytes for each Unicode character.

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 unicode-char.max-bytes parameter to allocate more bytes for each Unicode character.

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.

Note

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

    Note

    If 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

    Note

    StarRocks requires primary key columns to appear first in a table. Any new columns must be added after them.

  • RENAME COLUMN EVENT

    Note

    Requires 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.

Important

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.

Warning

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