All Products
Search
Document Center

E-MapReduce:Use the Flink service in a Dataflow cluster to synchronize data from MySQL to StarRocks by using the CTAS statement

Last Updated:Oct 10, 2026

Use the CREATE TABLE AS (CTAS) statement in a Flink SQL Client to stream MySQL changes into EMR Serverless StarRocks — including automatic table creation and real-time schema change propagation.

How it works

When you run a CTAS statement, Flink performs two operations in sequence:

  1. Table creation check: Flink checks whether the destination table exists in StarRocks.

    • If it doesn't exist, Flink creates it automatically with the same schema as the source MySQL table.

    • If it exists but the schemas don't match, the job returns an error.

  2. Data synchronization: Flink starts a continuous job that streams data changes and schema changes from MySQL to StarRocks.

Schema changes are detected by comparing the schema of consecutive data records — not by parsing DDL statements. This means schema changes only sync when data actually flows through. For example, adding a column to MySQL does not trigger a schema update in StarRocks until a new row arrives with that column populated.

Supported schema changes

ChangeBehavior in StarRocks
Add a nullable columnAdded to the end of the destination table; data synchronized
Delete a nullable columnColumn retained in destination table, filled with NULL
Rename a columnRenamed column added to end; original column filled with NULL

Unsupported schema changes

  • Data type changes (for example, VARCHAR to BIGINT, or NOT NULL to NULLABLE)

  • Constraint changes (primary key or index)

  • Addition or deletion of a non-nullable column

If any unsupported schema change occurs in MySQL, delete the destination table in StarRocks and restart the CTAS job. StarRocks will recreate the table and resync all historical data.

Because CTAS uses record-level schema comparison, deleting then re-adding a column with no intervening data changes is treated as no schema change. The schema sync triggers only when new data arrives.

For supported field type mappings between Flink and StarRocks, see Continuously load data from Apache Flink®.

Version requirements

ComponentVersion
Dataflow clusterEMR-3.42.0 or later, or EMR-5.8.0 or later
ApsaraDB RDS for MySQL5.7 or later
Flink (older versions)If using a version earlier than vvr-6.0.5-flink-1.15, add 'sink.use.new-apiapi' = 'false' to the WITH clause

Prerequisites

Before you begin, ensure that you have:

This tutorial uses MySQL 5.7 and a Dataflow cluster of EMR-3.42.0 as an example.

Limitations

  • The Dataflow cluster, StarRocks instance, and ApsaraDB RDS for MySQL instance must be in the same virtual private cloud (VPC).

  • The Dataflow cluster and StarRocks instance must be accessible over the Internet.

  • The ApsaraDB RDS for MySQL engine version must be 5.7 or later.

  • The Dataflow cluster must be EMR-3.42.0 or later, or EMR-5.8.0 or later.

Step 1: Prepare test data

  1. Create a test database and a test account in your ApsaraDB RDS for MySQL instance, then grant read and write permissions to the account. For more information, see Create an ApsaraDB RDS for MySQL instance and configure a database.

    This tutorial uses test_cdc as the database name and emr_test as the account name.
  2. Connect to the ApsaraDB RDS for MySQL instance with the test account. For more information, see Connect to an ApsaraDB RDS for MySQL instance.

  3. Create a test table and insert a row:

    USE test_cdc;
    
    CREATE TABLE IF NOT EXISTS `runoob_tbl` (
      `runoob_id`      INT UNSIGNED AUTO_INCREMENT,
      `runoob_title`   VARCHAR(100) NOT NULL,
      `runoob_author`  VARCHAR(40)  NOT NULL,
      `submission_date` DATE,
      `add_col`        INT DEFAULT NULL,
      PRIMARY KEY (`runoob_id`)
    ) ENGINE=InnoDB DEFAULT CHARSET=utf8;
    
    INSERT INTO test_cdc.`runoob_tbl`
      (`runoob_id`, `runoob_title`, `runoob_author`, `submission_date`, `add_col`)
    VALUES (18, 'first', 'tom', '2022-06-22 17:13:44', 3);
  4. Connect to your EMR Serverless StarRocks instance. For more information, see Connect to a StarRocks instance by using a MySQL client.

  5. Create the test_cdc database and a StarRocks user, then grant permissions. You can create a super administrator user named test (with the example password 1qaz!QAZ), or create a regular user named test and grant permissions on the database to the user. For more information, see Manage users.

    CREATE DATABASE test_cdc;
    CREATE USER 'test' IDENTIFIED BY '1qaz!QAZ';
    GRANT ALL ON test_cdc TO test;

Step 2: Upload custom connectors

SSH into the Dataflow cluster and upload the Flink connectors for StarRocks and MySQL. For more information, see Log on to a cluster.

Download and upload the following JAR files to /opt/apps/FLINK/flink-current/lib on the Dataflow cluster:

Step 3: Execute the CTAS statement

Start a YARN session and open the SQL Client

  1. SSH into the Dataflow cluster. For more information, see Log on to a cluster.

  2. Go to the Flink directory and start a YARN session in detached mode:

    cd /opt/apps/FLINK/flink-current
    ./bin/yarn-session.sh --detached

    A successful start prints application_XXXX_YY in the output. This is the session ID.

    sessionid

  3. Open the SQL Client with the session ID:

    ./bin/sql-client.sh -s <application_XXXX_YY>

    Replace <application_XXXX_YY> with the session ID from the previous step.

Create catalogs for MySQL and StarRocks

Run the following statements in the SQL Client to register both data sources as catalogs:

CREATE CATALOG sr WITH (
  'type'     = 'starrocks',
  'endpoint' = 'fe-c-9b354c83e891****-internal.starrocks.aliyuncs.com:9030',
  'username' = 'test',
  'password' = '1qaz!QAZ',
  'dbname'   = 'test_cdc'
);

CREATE CATALOG mysql WITH (
  'type'             = 'mysql',
  'hostname'         = 'rm-2zepd6e20u3od****.mysql.rds.aliyuncs.com',
  'port'             = '3306',
  'username'         = 'emr_test',
  'password'         = '123456',
  'default-database' = 'test_cdc'
);

StarRocks catalog parameters

ParameterDescription
typeSet to starrocks
endpointInternal endpoint and query port of the FE node, in the format <FE-internal-endpoint>:9030. To get the internal endpoint, see View the instance list and details.
usernameUsername for the StarRocks database (created in Step 1; test in this example)
passwordPassword for the StarRocks database (1qaz!QAZ in this example)
dbnameStarRocks database name (test_cdc in this example)

MySQL catalog parameters

ParameterDescription
typeSet to mysql
hostnameInternal endpoint of the ApsaraDB RDS for MySQL instance. Copy it from the Database Connection page in the ApsaraDB RDS console.
portMySQL port number. Default: 3306
usernameMySQL account name (created in Step 1; emr_test in this example)
passwordMySQL account password (123456 in this example)
default-databaseDefault MySQL database name (test_cdc in this example)

Choose a delivery semantics and run the CTAS statement

Choose one of the following three modes based on your requirements, then run the corresponding CTAS statement in the SQL Client.

ModeDelivery semanticsKey trade-off
At-least-once (default)At-least-onceLow write latency, lower memory usage; duplicate records possible on failure
Exactly-onceExactly-onceNo data loss or duplication; data visibility depends on checkpoint interval
SimpleAt-least-onceEasiest to configure — no need to specify table properties manually; no partition support

At-least-once mode

Use sink.buffer-flush.interval-ms to control how often data is flushed to StarRocks.

USE CATALOG sr;

CREATE TABLE IF NOT EXISTS runoob_tbl1 WITH (
  'starrocks.create.table.properties' = 'engine = olap primary key(runoob_id) distributed by hash(runoob_id) buckets 8',
  'database-name'                     = 'test_cdc',
  'jdbc-url'                          = 'jdbc:mysql://fe-c-9b354c83e891****-internal.starrocks.aliyuncs.com:9030',
  'load-url'                          = 'fe-c-9b354c83e891****-internal.starrocks.aliyuncs.com:8030',
  'table-name'                        = 'runoob_tbl_sr',
  'username'                          = 'test',
  'password'                          = '1qaz!QAZ',
  'sink.buffer-flush.interval-ms'     = '5000',
  'sink.properties.row_delimiter'     = '\x02',
  'sink.properties.column_separator'  = '\x01'
) AS TABLE mysql.test_cdc.runoob_tbl /*+ OPTIONS (
  'connector'     = 'mysql-cdc',
  'hostname'      = 'rm-2zepd6e20u3od****.mysql.rds.aliyuncs.com',
  'port'          = '3306',
  'username'      = 'test',
  'password'      = '123456',
  'database-name' = 'test_cdc',
  'table-name'    = 'runoob_tbl'
) */;

Exactly-once mode

Set the checkpoint interval before running the CTAS statement. Data becomes visible only after each checkpoint completes. For more information, see Checkpointing.

SET 'execution.checkpointing.interval' = '1 min';
SET 'execution.checkpointing.mode'     = 'EXACTLY_ONCE';
SET 'execution.checkpointing.timeout'  = '10 min';

USE CATALOG sr;

CREATE TABLE IF NOT EXISTS runoob_tbl1 WITH (
  'starrocks.create.table.properties' = 'engine = olap primary key(runoob_id) distributed by hash(runoob_id) buckets 8',
  'database-name'                     = 'test_cdc',
  'jdbc-url'                          = 'jdbc:mysql://fe-c-9b354c83e891****-internal.starrocks.aliyuncs.com:9030',
  'load-url'                          = 'fe-c-9b354c83e891****-internal.starrocks.aliyuncs.com:8030',
  'table-name'                        = 'runoob_tbl',
  'username'                          = 'test',
  'password'                          = '1qaz!QAZ',
  'sink.semantic'                     = 'exactly-once',
  'sink.properties.row_delimiter'     = '\x02',
  'sink.properties.column_separator'  = '\x01'
) AS TABLE mysql.test_cdc.runoob_tbl /*+ OPTIONS (
  'connector'     = 'mysql-cdc',
  'hostname'      = 'rm-2zepd6e20u3od****.mysql.rds.aliyuncs.com',
  'port'          = '3306',
  'username'      = 'test',
  'password'      = '123456',
  'database-name' = 'test_cdc',
  'table-name'    = 'runoob_tbl'
) */;

Simple mode

StarRocks automatically infers table properties from the MySQL schema: engine = olap, primary key matching MySQL's primary key, and distributed by hash on all primary key columns. No partitions are created. Use normal mode if you need partitions.

USE CATALOG sr;

CREATE TABLE IF NOT EXISTS runoob_tbl1 WITH (
  'starrocks.create.table.properties' = 'buckets 8',
  'starrocks.create.table.mode'       = 'simple',
  'database-name'                     = 'test_cdc',
  'jdbc-url'                          = 'jdbc:mysql://fe-c-9b354c83e891****-internal.starrocks.aliyuncs.com:9030',
  'load-url'                          = 'fe-c-9b354c83e891****-internal.starrocks.aliyuncs.com:8030',
  'table-name'                        = 'runoob_tbl_sr',
  'username'                          = 'test',
  'password'                          = '1qaz!QAZ',
  'sink.buffer-flush.interval-ms'     = '5000',
  'sink.properties.row_delimiter'     = '\x02',
  'sink.properties.column_separator'  = '\x01'
) AS TABLE mysql.test_cdc.runoob_tbl /*+ OPTIONS (
  'connector'     = 'mysql-cdc',
  'hostname'      = 'rm-2zepd6e20u3od****.mysql.rds.aliyuncs.com',
  'port'          = '3306',
  'username'      = 'emr_test',
  'password'      = '123456',
  'database-name' = 'test_cdc',
  'table-name'    = 'runoob_tbl'
) */;

WITH clause parameters

ParameterRequiredDescription
starrocks.create.table.propertiesYesTable definition properties appended after field definitions when creating the StarRocks table — for example, engine, key type, and buckets.
database-nameYesStarRocks database name (test_cdc in this example)
jdbc-urlYesJDBC URL for StarRocks queries. Format: jdbc:mysql://<FE-internal-endpoint>:9030. To get the internal endpoint, see View the instance list and details.
load-urlYesHTTP load endpoint of the FE node. Format: <FE-internal-endpoint>:8030. To get the internal endpoint, see View the instance list and details.
sink.semanticNoDelivery semantics. Set to exactly-once for strict consistency. Default: at-least-once.
starrocks.create.table.modeNonormal (default): specify engine, key type, and buckets in starrocks.create.table.properties. simple: StarRocks infers engine, primary key, and distribution from the MySQL schema; only buckets is required.
sink.properties.row_delimiterNoCustom row delimiter (for example, \x02)
sink.properties.column_separatorNoCustom column separator (for example, \x01)
If you use a Flink version earlier than vvr-6.0.5-flink-1.15, add 'sink.use.new-apiapi' = 'false' to the WITH clause. For other configuration options, see Continuously load data from Apache Flink.

OPTIONS clause parameters

All examples use the OPTIONS hint to pass MySQL CDC connector settings inline with the CTAS statement.

ParameterDescription
connectorSet to mysql-cdc
hostnameInternal endpoint of the ApsaraDB RDS for MySQL instance. Copy it from the Database Connection page in the ApsaraDB RDS console.
portMySQL port number. Default: 3306
usernameMySQL account name (emr_test in this example)
passwordMySQL account password
table-nameSource MySQL table name (runoob_tbl in this example)
database-nameSource MySQL database name (test_cdc in this example)

Step 4: Verify data synchronization

If checkpointing is enabled, wait up to one checkpoint interval before querying StarRocks.

Verify initial data

Connect to StarRocks and query the destination table:

USE test_cdc;
SELECT * FROM runoob_tbl1;

Expected output:

+-----------+--------------+---------------+-----------------+---------+
| runoob_id | runoob_title | runoob_author | submission_date | add_col |
+-----------+--------------+---------------+-----------------+---------+
|        18 | first        | tom           | 2022-06-22      |       3 |
+-----------+--------------+---------------+-----------------+---------+

Verify inserts

Insert a row in MySQL:

INSERT INTO runoob_tbl (`runoob_id`, `runoob_title`, `runoob_author`, `submission_date`, `add_col`)
VALUES (1, 'second', 'tom2', '2022-06-23', 1);

Query StarRocks to confirm the row was replicated:

SELECT * FROM runoob_tbl1;

Expected output:

+-----------+--------------+---------------+-----------------+---------+
| runoob_id | runoob_title | runoob_author | submission_date | add_col |
+-----------+--------------+---------------+-----------------+---------+
|         1 | second       | tom2          | 2022-06-23      |       1 |
|        18 | first        | tom           | 2022-06-22      |       3 |
+-----------+--------------+---------------+-----------------+---------+

Verify updates

Update a row in MySQL:

UPDATE runoob_tbl SET runoob_title = 'new' WHERE runoob_id = 18;

Query StarRocks to confirm the update was replicated:

SELECT * FROM runoob_tbl1;

Expected output:

+-----------+--------------+---------------+-----------------+---------+
| runoob_id | runoob_title | runoob_author | submission_date | add_col |
+-----------+--------------+---------------+-----------------+---------+
|         1 | second       | tom2          | 2022-06-23      |       1 |
|        18 | new          | tom           | 2022-06-22      |       3 |
+-----------+--------------+---------------+-----------------+---------+

Verify deletes

Delete a row in MySQL:

DELETE FROM runoob_tbl WHERE runoob_id = 1;

Query StarRocks to confirm the deletion was replicated:

SELECT * FROM runoob_tbl1;

Expected output:

+-----------+--------------+---------------+-----------------+---------+
| runoob_id | runoob_title | runoob_author | submission_date | add_col |
+-----------+--------------+---------------+-----------------+---------+
|        18 | new          | tom           | 2022-06-22      |       3 |
+-----------+--------------+---------------+-----------------+---------+

Verify schema changes

Add a nullable column in MySQL:

ALTER TABLE `runoob_tbl` ADD COLUMN `add_col2` INT;

Insert a row that populates the new column:

INSERT INTO runoob_tbl (`runoob_id`, `runoob_title`, `runoob_author`, `submission_date`, `add_col`, `add_col2`)
VALUES (1, 'second', 'tom2', '2022-06-23', 1, 2);

Query StarRocks to confirm the schema change was propagated:

SELECT * FROM runoob_tbl1;

Expected output:

+-----------+--------------+---------------+-----------------+---------+---------+
| runoob_id | runoob_title | runoob_author | submission_date | add_col | add_co2 |
+-----------+--------------+---------------+-----------------+---------+---------+
|         1 | second       | tom2          | 2022-06-23      |       1 |       2 |
|        18 | new          | tom           | 2022-06-22      |       3 |    NULL |
+-----------+--------------+---------------+-----------------+---------+---------+

The new add_col2 column appears in StarRocks. The existing row (runoob_id=18) shows NULL because the column was not present when that row was last written.

CDAS introduction

The CREATE DATABASE AS (CDAS) statement is syntactic sugar for CTAS. A single CDAS statement generates one Flink job that streams an entire MySQL database — or a selected subset of tables — into StarRocks. Use the INCLUDING TABLE clause to select specific tables.

Create the MySQL and StarRocks catalogs first (same as for CTAS), then run the CDAS statement:

CREATE DATABASE IF NOT EXISTS sr_db WITH (
  'starrocks.create.table.properties' = 'buckets 8',
  'starrocks.create.table.mode'       = 'simple',
  'jdbc-url'                          = 'jdbc:mysql://fe-c-9b354c83e891****-internal.starrocks.aliyuncs.com:9030',
  'load-url'                          = 'fe-c-9b354c83e891****-internal.starrocks.aliyuncs.com:8030',
  'username'                          = 'test',
  'password'                          = '1qaz!QAZ',
  'sink.buffer-flush.interval-ms'     = '5000',
  'sink.properties.row_delimiter'     = '\x02',
  'sink.properties.column_separator'  = '\x01'
) AS DATABASE mysql.test_cdc INCLUDING TABLE
  'tabl1', 'tbl2', 'tbl3' /*+ OPTIONS (
  'connector'     = 'mysql-cdc',
  'hostname'      = 'rm-2zepd6e20u3od****.mysql.rds.aliyuncs.com',
  'port'          = '3306',
  'username'      = 'test',
  'password'      = '123456',
  'database-name' = 'test_cdc'
) */;

What's next

  • To sync schema changes across an entire database with fewer configuration steps, consider using CDAS. The statement syntax and parameters follow the same pattern as CTAS.

  • To learn how CTAS handles additional field type mappings, see Continuously load data from Apache Flink.