You can use a CREATE TABLE AS SELECT (CTAS) statement in Realtime Compute for Apache Flink to synchronize data from an ApsaraDB RDS for MySQL instance to an EMR Serverless StarRocks cluster in real time, including schema changes.
Background information
You can use CTAS or CDAS statements to synchronize MySQL data to EMR Serverless StarRocks. CTAS synchronizes the schema and data of a single table. CDAS synchronizes an entire database, or the schemas and data of multiple tables in the same database. This topic uses a CTAS statement. CDAS statements are used in a similar way to CTAS statements. For more information, see Synchronize MySQL to StarRocks using CTAS in Flink.
You can use a CTAS (CREATE TABLE AS) statement to automatically create a table in StarRocks that has the same schema as the table in MySQL, and to synchronize the data. Schema changes of the upstream table are also synchronized to the downstream table in real time. This improves the efficiency of creating tables in the destination storage and maintaining schema changes of the source table.
When a CTAS statement is executed, Flink performs the following steps:
-
Checks whether the destination table exists in the destination storage.
-
If the table does not exist, Flink uses the destination catalog to create the destination table in the destination storage. The destination table has the same schema as the data source.
-
If the table exists, table creation is skipped. If the schema of the existing destination table is inconsistent with that of the source table, an error is reported.
-
-
Submits and starts the data synchronization job to synchronize the data of the data source and its schema changes to the destination table.
Schema change synchronization policy: A CTAS statement synchronizes data in real time and also synchronizes schema changes of the source table to the destination table.
Schema changes include the creation of the initial table and subsequent table changes.
-
Schema changes that are currently supported:
-
Add a nullable column: The corresponding column is automatically added to the end of the destination table schema, and the data of the new column is automatically synchronized.
-
Delete a nullable column: The column is not directly deleted from the destination table. Instead, the data of the column is automatically filled with NULL values.
-
Rename a column: The renamed column is added to the end of the destination table, and the data of the column before renaming is automatically filled with NULL values.
For example, if col_a is renamed col_b, col_b is added to the end of the destination table, and the data of col_a is automatically filled with NULL values.
-
-
Schema changes that are not supported:
-
Data type changes.
For example, changing VARCHAR to BIGINT, or changing the NOT NULL attribute to NULLABLE.
-
Changes to constraints such as primary keys or indexes.
-
Adding or deleting non-nullable columns.
-
- If the schema of the source table has one of the preceding changes, you must delete the destination table and restart the job that executes the CTAS statement. This way, the destination table is created again and the historical data is resynchronized to the destination table.
- The CTAS statement does not identify the types of DDL statements, but compares the schema differences between the two data records before and after the schema is changed. Therefore, if you delete a column and then add the column again, and no data changes between the two DDL statements that are used to delete and add the column, the CTAS statement considers that no schema change occurs. Similarly, the CTAS statement does not trigger schema change synchronization even if you add a column to the source table. The statement identifies the schema change only when the data changes in the source table. In this case, the statement synchronizes the schema change to the destination table.
- For more information about the field types supported by the CTAS statement, see Continuously load data from Apache Flink®.
Prerequisites
-
You have activated fully managed Realtime Compute for Apache Flink and created a Flink cluster. For more information, see Activate fully managed Flink and Quick start with Flink SQL jobs.
-
You have created an EMR Serverless StarRocks instance. For more information, see Create an instance.
-
You have created an ApsaraDB RDS for MySQL instance. For more information, see Create an ApsaraDB RDS for MySQL instance.
This topic uses MySQL 5.7 and Flink vvr-8.0.11-flink-1.17 as examples.
Limitations
-
The Flink cluster, EMR Serverless StarRocks instance, and ApsaraDB RDS for MySQL instance must be in the same Virtual Private Cloud (VPC).
-
The ApsaraDB RDS for MySQL instance must be version 5.7 or later.
Step 1: Prepare test data
-
Create a test database and account. For more information, see Create a database and an account.
After creating the database and account, grant read and write permissions to the test account.
NoteIn this example, the database is named
test_cdcand the account is namedtest. -
Use the test account to connect to the MySQL instance. For more information, see Log on to an ApsaraDB RDS for MySQL instance by using DMS.
-
Run the following commands in MySQL to create a data table.
use test_cdc; -- Create a table 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 data INSERT INTO test_cdc.`runoob_tbl` (`runoob_id`,`runoob_title`,`runoob_author`,`submission_date`,`add_col`) values (18,'first','tom','2025-06-22 17:13:44',3) -
Log on to and connect to the EMR Serverless StarRocks instance. For more information, see Connect to a StarRocks instance through a MySQL client.
-
Run the following commands to create the
test_cdcdatabase, create a user namedtestwith the password1qaz!QAZ, and grant the user permissions on the database and its tables. For more information, see Manage users.-- Create a database CREATE DATABASE test_cdc; -- Create a user CREATE USER 'test' IDENTIFIED by '1qaz!QAZ'; -- Grant database permissions to the user GRANT ALL on test_cdc to test; -- Grant table permissions to the user GRANT ALL ON ALL TABLES IN DATABASE test_cdc to test;
Step 2: Create catalogs
On the Data Management page in the Realtime Compute for Apache Flink console, create MySQL and StarRocks catalogs. For more information, see Data Management.
The following parameter configurations are examples. Adjust the values based on your environment.
-
MySQL catalog
-
Sample code
CREATE CATALOG mysql WITH ( 'type' = 'mysql', 'hostname' = 'rm-2zepd6e20u3od****.mysql.rds.aliyuncs.com', 'port' = '3306', 'username' = 'emr-test', 'password' = '123456', 'default-database' = 'test_cdc' ); -
Parameters
Parameter
Description
type
The type of the catalog. Set this to
mysql.hostname
The internal endpoint of the ApsaraDB RDS for MySQL instance, which you can copy from the Database Connection page in the ApsaraDB RDS for MySQL console. Example:
rm-2zepd6e20u3od****.mysql.rds.aliyuncs.com.port
The port number of the MySQL database service. The default value is 3306.
username
The username for the MySQL database service.
Use the username for the account that you created in Step 1: Prepare test data. In this example, the username is
test.password
The password for the MySQL database service.
Use the password for the account that you created in Step 1: Prepare test data.
default-database
The name of the default MySQL database.
Use the name of the database that you created in Step 1: Prepare test data. In this example, the name is
test_cdc.
-
-
StarRocks catalog
-
Sample code
CREATE CATALOG sr WITH ( 'type' = 'starrocks', 'endpoint' = 'fe-c-9b354c83e891****-internal.starrocks.aliyuncs.com:9030', 'username' = 'test', 'password' = '1qaz!QAZ', 'dbname' = 'test_cdc' ); -
Parameters
Parameter
Description
type
The type of the catalog. Set this to
starrocks.endpoint
The internal endpoint and query port of an FE node, specified in the format
Internal endpoint of an FE node of the EMR Serverless StarRocks instance:9030.Example:
fe-c-9b354c83e891****-internal.starrocks.aliyuncs.com:9030.NoteTo obtain the internal endpoint of the FE node, see View instance list and details.
username
The username for the StarRocks cluster.
Use the username for the account that you created in Step 1: Prepare test data. In this example, the username is
test.password
The password for the StarRocks database service.
Use the password for the account that you created in Step 1: Prepare test data.
dbname
The name of the StarRocks database.
Use the name of the database that you created in Step 1: Prepare test data. In this example, the name is
test_cdc.
-
Step 3: Create and publish a job
-
On the Data Development > ETL page in the Realtime Compute for Apache Flink console, write a CTAS statement.
The following examples show three different modes:
-
At-least-once semantics: Uses the sink.buffer-flush.interval-ms parameter to control the write interval to StarRocks, offering short write intervals and low memory usage.
/* At-least-once semantics */ use CATALOG sr; CREATE TABLE IF NOT EXISTS runoob_tbl 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.buffer-flush.interval-ms' = '5000', 'sink.properties.row_delimiter' = '\x02', 'sink.properties.column_separator' = '\x01' ) as table mysql.test_cdc.runoob_tbl; -
Exactly-once semantics: Requires a checkpoint interval. Guarantees no data loss or duplication, even during exceptions. Data visibility depends on the checkpoint interval. For more information, see Checkpointing.
/* Exactly-once semantics. */ 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_tbl 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; -
Simple mode: Automatically mirrors the source MySQL table's schema, eliminating the need to specify fields. Does not support partitions. Use the normal mode for tables that require partitions.
/* The two examples above use the normal mode. This example demonstrates the simple mode. */ use CATALOG sr; CREATE TABLE IF NOT EXISTS runoob_tbl with ( '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', '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;
Table 1.
WITHparametersParameter
Required
Description
starrocks.create.table.properties
Yes
Additional clauses for the StarRocks
CREATE TABLEstatement, excluding column definitions. Examples includeengine,key, andbuckets.database-name
Yes
The name of the StarRocks database.
In this example, the name is
test_cdc.jdbc-url
Yes
The JDBC URL used to run query operations in StarRocks.
For example, jdbc:mysql://fe-c-9b354c83e891****-internal.starrocks.aliyuncs.com:9030. In this string,
fe-c-9b354c83e891****-internal.starrocks.aliyuncs.comis the internal endpoint of an FE node of the EMR Serverless StarRocks instance.NoteTo obtain the internal endpoint of the FE node, see View instance list and details.
load-url
Yes
Specify the internal endpoint and query port of the FE in the format
Internal endpoint of the FE node of the EMR Serverless StarRocks instance:8030.Example:
fe-c-9b354c83e891****-internal.starrocks.aliyuncs.com:8030.NoteTo obtain the internal endpoint of the FE node, see View instance list and details.
sink.semantic
No
The delivery guarantee. Set to
exactly-onceto ensure data consistency. Default:at-least-once.starrocks.create.table.mode
No
The mode for creating the table. Valid values:
-
normal mode (default): You must specify the complete properties, such as engine, key, and buckets, in the starrocks.create.table.properties configuration.
-
simplemode: The default engine isolapand the key type isprimary key. The primary key must be identical to the primary key of the MySQL table. By default, the table is distributed by hash on all primary key columns and is not partitioned. In the starrocks.create.table.properties configuration, you must specify thebucketsparameter. Other parameters, such asproperties, are optional.
sink.properties.row_delimiter
No
A custom row delimiter.
sink.properties.column_separator
No
A custom column delimiter.
Note-
The
sink.use.new-apiparameter is removed in Flink versions 1.15-vvr-6.0.5 and later. If you use a version earlier than 1.15-vvr-6.0.5, you must add'sink.use.new-api'='false',to the WITH parameters. -
For information about other configurations, see Continuously load data from Apache Flink.
Table 2. Source connector options
Parameter
Description
connector
The type of connector. Set this to
mysql-cdc.hostname
The internal endpoint of the ApsaraDB RDS for MySQL instance. You can copy it from the Database Connection page in the ApsaraDB RDS for MySQL console. Example:
rm-bp1nu0c46fn9k****.mysql.rds.aliyuncs.com.port
The port number of the MySQL database service. The default value is 3306.
username
The username for the MySQL database service.
Use the username for the account that you created in Step 1: Prepare test data. In this example, the username is
test.password
The password for the MySQL database service.
Use the password for the account that you created in Step 1: Prepare test data.
table-name
The name of the table in StarRocks.
Use the name of the table that you created in Step 1: Prepare test data. In this example, the name is
runoob_tbl.database-name
The name of the default MySQL database.
Use the name of the database that you created in Step 1: Prepare test data. In this example, the name is
test_cdc. -
-
Click Publish.
-
On the Publish new version page, select a deployment target and click OK.
-
On the Job O&M page, find the target job and click START in the Actions column.
The Realtime Compute for Apache Flink console does not support debugging CTAS statements.
Step 4: Verify data synchronization
Data query
-
Log on to and connect to the EMR Serverless StarRocks instance. For more information, see Connect to a StarRocks instance through a MySQL client.
-
In the StarRocks connection window, run the following commands to view the table data.
use test_cdc; select * from runoob_tbl;The following output indicates that the data from MySQL has been synchronized to StarRocks.
+-----------+--------------+---------------+-----------------+---------+ | runoob_id | runoob_title | runoob_author | submission_date | add_col | +-----------+--------------+---------------+-----------------+---------+ | 18 | first | tom | 2025-06-22 | 3 | +-----------+--------------+---------------+-----------------+---------+
Data insertion
-
In the ApsaraDB RDS for MySQL window, run the following command to insert data.
INSERT INTO runoob_tbl(`runoob_id`,`runoob_title`,`runoob_author`,`submission_date`,`add_col`) values(1,'second','tom2','2022-06-23',1) -
In the StarRocks connection window, run the following command to view the table data.
select * from runoob_tbl;The following output indicates that the data has been successfully inserted.
+-----------+--------------+---------------+-----------------+---------+ | runoob_id | runoob_title | runoob_author | submission_date | add_col | +-----------+--------------+---------------+-----------------+---------+ | 1 | second | tom2 | 2022-06-23 | 1 | | 18 | first | tom | 2025-06-22 | 3 | +-----------+--------------+---------------+-----------------+---------+
Data update
-
In the ApsaraDB RDS for MySQL window, run the following command to update a record.
update runoob_tbl set runoob_title= 'new' where runoob_id = 18 -
In the StarRocks connection window, run the following command to view the table data.
select * from runoob_tbl;The following output indicates that the data update has been synchronized.
+-----------+--------------+---------------+-----------------+---------+ | runoob_id | runoob_title | runoob_author | submission_date | add_col | +-----------+--------------+---------------+-----------------+---------+ | 1 | second | tom2 | 2022-06-23 | 1 | | 18 | new | tom | 2025-06-22 | 3 | +-----------+--------------+---------------+-----------------+---------+
Data deletion
-
In the ApsaraDB RDS for MySQL window, run the following command to delete a record.
DELETE FROM runoob_tbl WHERE runoob_id = 1 -
In the StarRocks connection window, run the following command to view the table data.
select * from runoob_tbl;The following output indicates that the data deletion has been synchronized.
+-----------+--------------+---------------+-----------------+---------+ | runoob_id | runoob_title | runoob_author | submission_date | add_col | +-----------+--------------+---------------+-----------------+---------+ | 18 | new | tom | 2025-06-22 | 3 | +-----------+--------------+---------------+-----------------+---------+
Nullable column addition
-
In the ApsaraDB RDS for MySQL window, run the following command to add a nullable column.
alter table `runoob_tbl` add COLUMN `add_col2` INT; -
Run the following command to insert data.
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) -
In the StarRocks connection window, run the following command to view the table data.
select * from runoob_tbl;The following output indicates that the schema change has been successfully synchronized.
+-----------+--------------+---------------+-----------------+---------+----------+ | runoob_id | runoob_title | runoob_author | submission_date | add_col | add_col2 | +-----------+--------------+---------------+-----------------+---------+----------+ | 18 | new | tom | 2025-06-22 | 3 | NULL | +-----------+--------------+---------------+-----------------+---------+----------+ | 1 | second | tom2 | 2025-06-23 | 1 | 2 | +-----------+--------------+---------------+-----------------+---------+----------+
CDAS
A CDAS statement synchronizes an entire MySQL database to multiple corresponding tables in StarRocks. You can use the including table clause to select specific tables for synchronization.
As with CTAS, you must create the MySQL and StarRocks catalogs before running a 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',
'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';
References
You can also synchronize data to StarRocks using a data ingestion YAML file. For more information, see Data ingestion.