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:
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.
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
| Change | Behavior in StarRocks |
|---|---|
| Add a nullable column | Added to the end of the destination table; data synchronized |
| Delete a nullable column | Column retained in destination table, filled with NULL |
| Rename a column | Renamed column added to end; original column filled with NULL |
Unsupported schema changes
Data type changes (for example,
VARCHARtoBIGINT, orNOT NULLtoNULLABLE)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
| Component | Version |
|---|---|
| Dataflow cluster | EMR-3.42.0 or later, or EMR-5.8.0 or later |
| ApsaraDB RDS for MySQL | 5.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:
A Dataflow cluster created in the new console with the Flink service selected. For more information, see Create a cluster
An EMR Serverless StarRocks instance. For more information, see Create an instance
An ApsaraDB RDS for MySQL instance. For more information, see Create an ApsaraDB RDS for MySQL instance and configure a database
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
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_cdcas the database name andemr_testas the account name.Connect to the ApsaraDB RDS for MySQL instance with the test account. For more information, see Connect to an ApsaraDB RDS for MySQL instance.
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);Connect to your EMR Serverless StarRocks instance. For more information, see Connect to a StarRocks instance by using a MySQL client.
Create the
test_cdcdatabase and a StarRocks user, then grant permissions. You can create a super administrator user namedtest(with the example password1qaz!QAZ), or create a regular user namedtestand 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
SSH into the Dataflow cluster. For more information, see Log on to a cluster.
Go to the Flink directory and start a YARN session in detached mode:
cd /opt/apps/FLINK/flink-current ./bin/yarn-session.sh --detachedA successful start prints
application_XXXX_YYin the output. This is the session ID.
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
| Parameter | Description |
|---|---|
type | Set to starrocks |
endpoint | Internal 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. |
username | Username for the StarRocks database (created in Step 1; test in this example) |
password | Password for the StarRocks database (1qaz!QAZ in this example) |
dbname | StarRocks database name (test_cdc in this example) |
MySQL catalog parameters
| Parameter | Description |
|---|---|
type | Set to mysql |
hostname | Internal endpoint of the ApsaraDB RDS for MySQL instance. Copy it from the Database Connection page in the ApsaraDB RDS console. |
port | MySQL port number. Default: 3306 |
username | MySQL account name (created in Step 1; emr_test in this example) |
password | MySQL account password (123456 in this example) |
default-database | Default 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.
| Mode | Delivery semantics | Key trade-off |
|---|---|---|
| At-least-once (default) | At-least-once | Low write latency, lower memory usage; duplicate records possible on failure |
| Exactly-once | Exactly-once | No data loss or duplication; data visibility depends on checkpoint interval |
| Simple | At-least-once | Easiest 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
| Parameter | Required | Description |
|---|---|---|
starrocks.create.table.properties | Yes | Table definition properties appended after field definitions when creating the StarRocks table — for example, engine, key type, and buckets. |
database-name | Yes | StarRocks database name (test_cdc in this example) |
jdbc-url | Yes | JDBC URL for StarRocks queries. Format: jdbc:mysql://<FE-internal-endpoint>:9030. To get the internal endpoint, see View the instance list and details. |
load-url | Yes | HTTP load endpoint of the FE node. Format: <FE-internal-endpoint>:8030. To get the internal endpoint, see View the instance list and details. |
sink.semantic | No | Delivery semantics. Set to exactly-once for strict consistency. Default: at-least-once. |
starrocks.create.table.mode | No | normal (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_delimiter | No | Custom row delimiter (for example, \x02) |
sink.properties.column_separator | No | Custom 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.
| Parameter | Description |
|---|---|
connector | Set to mysql-cdc |
hostname | Internal endpoint of the ApsaraDB RDS for MySQL instance. Copy it from the Database Connection page in the ApsaraDB RDS console. |
port | MySQL port number. Default: 3306 |
username | MySQL account name (emr_test in this example) |
password | MySQL account password |
table-name | Source MySQL table name (runoob_tbl in this example) |
database-name | Source 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.