All Products
Search
Document Center

E-MapReduce:Synchronize MySQL data to StarRocks using a CTAS statement in Realtime Compute for Apache Flink

Last Updated:Sep 20, 2026

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:

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

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

Note
  • 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

Note

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

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

    Note

    In this example, the database is named test_cdc and the account is named test.

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

  3. 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)
  4. Log on to and connect to the EMR Serverless StarRocks instance. For more information, see Connect to a StarRocks instance through a MySQL client.

  5. Run the following commands to create the test_cdc database, create a user named test with the password 1qaz!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.

Note

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.

      Note

      To 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

  1. 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. WITH parameters

    Parameter

    Required

    Description

    starrocks.create.table.properties

    Yes

    Additional clauses for the StarRocks CREATE TABLE statement, excluding column definitions. Examples include engine, key, and buckets.

    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.com is the internal endpoint of an FE node of the EMR Serverless StarRocks instance.

    Note

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

    Note

    To obtain the internal endpoint of the FE node, see View instance list and details.

    sink.semantic

    No

    The delivery guarantee. Set to exactly-once to 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.

    • simple mode: The default engine is olap and the key type is primary 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 the buckets parameter. Other parameters, such as properties, are optional.

    sink.properties.row_delimiter

    No

    A custom row delimiter.

    sink.properties.column_separator

    No

    A custom column delimiter.

    Note
    • Thesink.use.new-api parameter 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.

  2. Click Publish.

  3. On the Publish new version page, select a deployment target and click OK.

  4. On the Job O&M page, find the target job and click START in the Actions column.

Note

The Realtime Compute for Apache Flink console does not support debugging CTAS statements.

Step 4: Verify data synchronization

Data query

  1. Log on to and connect to the EMR Serverless StarRocks instance. For more information, see Connect to a StarRocks instance through a MySQL client.

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

  1. 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)
  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 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

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

  1. In the ApsaraDB RDS for MySQL window, run the following command to delete a record.

    DELETE FROM runoob_tbl WHERE runoob_id = 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 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

  1. 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;
  2. 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)
  3. 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.