Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Tâche d'ingestion de données Flink CDC

Dernière mise à jour :Aug 09, 2026

Realtime Compute for Apache Flink propose une fonctionnalité d'ingestion de données puissante basée sur Flink CDC. Ce guide explique comment créer une tâche d'ingestion de données Flink CDC pour synchroniser l'intégralité d'une base de données MySQL vers une base de données StarRocks.

Prérequis

Contexte

Supposons que votre instance ApsaraDB RDS for MySQL contient une base de données nommée order_dw_mysql avec trois tables métier : orders, orders_pay et product_catalog. Pour synchroniser ces tables et leurs données vers la base de données order_dw_sr dans StarRocks, procédez comme suit :

  1. Étape 1 : Préparer les données de test dans ApsaraDB RDS for MySQL

  2. Étape 2 : Développer une tâche d'ingestion de données Flink CDC

  3. Étape 3 : Démarrer la tâche d'ingestion de données Flink CDC

  4. Étape 4 : Vérifier les résultats de synchronisation dans StarRocks

Étape 1 : Préparer les données de test MySQL

  1. Créez une base de données et un compte.

    Créez une base de données nommée order_dw_mysql ainsi qu'un compte standard disposant des autorisations de lecture et d'écriture sur cette base. Pour plus d'informations, consultez les rubriques (Obsolète, redirigé vers « Étape 1 ») Créer une base de données et un compte et Gérer les bases de données.

  2. Connectez-vous à l'instance ApsaraDB RDS for MySQL à l'aide de Data Management (DMS).

  3. Dans la fenêtre SQL Console, saisissez les commandes suivantes, puis cliquez sur Execute pour créer trois tables métier et y insérer des données.

    CREATE TABLE `orders` (
      order_id bigint not null primary key,
      user_id varchar(50) not null,
      shop_id bigint not null,
      product_id bigint not null,
      buy_fee numeric(20,2) not null,   
      create_time timestamp not null,
      update_time timestamp not null default now(),
      state int not null 
    );
    CREATE TABLE `orders_pay` (
      pay_id bigint not null primary key,
      order_id bigint not null,
      pay_platform int not null, 
      create_time timestamp not null
    );
    CREATE TABLE `product_catalog` (
      product_id bigint not null primary key,
      catalog_name varchar(50) not null
    );
    -- Prepare data
    INSERT INTO product_catalog VALUES(1, 'phone_aaa'),(2, 'phone_bbb'),(3, 'phone_ccc'),(4, 'phone_ddd'),(5, 'phone_eee');
    INSERT INTO orders VALUES
    (100001, 'user_001', 12345, 1, 5000.05, '2023-02-15 16:40:56', '2023-02-15 18:42:56', 1),
    (100002, 'user_002', 12346, 2, 4000.04, '2023-02-15 15:40:56', '2023-02-15 18:42:56', 1),
    (100003, 'user_003', 12347, 3, 3000.03, '2023-02-15 14:40:56', '2023-02-15 18:42:56', 1),
    (100004, 'user_001', 12347, 4, 2000.02, '2023-02-15 13:40:56', '2023-02-15 18:42:56', 1),
    (100005, 'user_002', 12348, 5, 1000.01, '2023-02-15 12:40:56', '2023-02-15 18:42:56', 1),
    (100006, 'user_001', 12348, 1, 1000.01, '2023-02-15 11:40:56', '2023-02-15 18:42:56', 1),
    (100007, 'user_003', 12347, 4, 2000.02, '2023-02-15 10:40:56', '2023-02-15 18:42:56', 1);
    INSERT INTO orders_pay VALUES
    (2001, 100001, 1, '2023-02-15 17:40:56'),
    (2002, 100002, 1, '2023-02-15 17:40:56'),
    (2003, 100003, 0, '2023-02-15 17:40:56'),
    (2004, 100004, 0, '2023-02-15 17:40:56'),
    (2005, 100005, 0, '2023-02-15 18:40:56'),
    (2006, 100006, 0, '2023-02-15 18:40:56'),
    (2007, 100007, 0, '2023-02-15 18:40:56');

Étape 2 : Développer une tâche Flink CDC

  1. Connectez-vous à la console de gestion Realtime Compute for Apache Flink.

  2. Cliquez sur Console pour accéder à l'espace de travail du projet.

  3. Dans le volet de navigation de gauche, sélectionnez Development > Data Ingestion.

  4. Cliquez sur l'icône image, choisissez New Draft with Template, sélectionnez MySQL to StarRocks data synchronization, puis cliquez sur Next.

  5. Saisissez un Job Name et une Location, spécifiez une Engine Version, puis cliquez sur OK.

  6. Configurez le code YAML de la tâche.

    L'exemple de code suivant illustre la synchronisation de toutes les tables de la base de données order_dw_mysql de MySQL vers la base de données order_dw_sr de StarRocks.

    source:
      type: mysql
      hostname: rm-bp1rk934iidc3****.mysql.rds.aliyuncs.com
      port: 3306
      username: ${secret_values.mysqlusername}
      password: ${secret_values.mysqlpassword}
      tables: order_dw_mysql.\.*
      server-id: 8601-8604
      # (Optional) Synchronize data from tables that are newly created during the incremental phase.
      scan.binlog.newly-added-table.enabled: true
      # (Optional) Synchronize table and column comments.
      include-comments.enabled: true
      # (Optional) Prioritize the distribution of unbounded chunks to prevent potential TaskManager OutOfMemory errors.
      scan.incremental.snapshot.unbounded-chunk-first.enabled: true
      # (Optional) Enable parsing filters to accelerate reading.
      scan.only.deserialize.captured.tables.changelog.enabled: true 
    sink:
      type: starrocks
      name: StarRocks Sink
      jdbc-url: jdbc:mysql://fe-c-b76b6aa51807****-internal.starrocks.aliyuncs.com:9030
      load-url: fe-c-b76b6aa51807****-internal.starrocks.aliyuncs.com:8030
      username: ${secret_values.starrocksusername}
      password: ${secret_values.starrockspassword}
      table.create.properties.replication_num: 1
      sink.buffer-flush.interval-ms: 5000 # Flush data every 5 seconds.
    route:
      - source-table: order_dw_mysql.\.*
        sink-table: order_dw_sr.<>
        replace-symbol: <>
        description: route all tables in source_db to sink_db
    pipeline:
      name: Sync MySQL Database to StarRocks

    Le tableau ci-dessous décrit les paramètres de configuration requis pour cet exemple. Pour en savoir plus sur les paramètres d'ingestion de données, consultez les rubriques MySQL et StarRocks.

    Remarque

    Les tâches YAML prennent uniquement en charge les variables de projet. Utilisez des variables pour éviter l'affichage en clair d'informations sensibles telles que les mots de passe. Pour plus d'informations, consultez la rubrique Gestion des variables.

    Catégorie

    Paramètre

    Description

    Valeur d'exemple

    source

    hostname

    Adresse IP ou nom d'hôte de la base de données MySQL.

    Nous recommandons d'utiliser le endpoint interne.

    rm-bp1rk934iidc3****.mysql.rds.aliyuncs.com

    port

    Numéro de port du service de base de données MySQL.

    3306

    username

    Nom d'utilisateur et mot de passe de la base de données MySQL. Utilisez les identifiants du compte créé à l'Étape 1 : Préparer les données de test dans ApsaraDB RDS for MySQL.

    ${secret_values.mysqlusername}

    password

    ${secret_values.mysqlpassword}

    tables

    Noms des tables MySQL. Vous pouvez utiliser des expressions régulières pour lire les données de plusieurs tables.

    Dans cette rubrique, toutes les tables et données de la base de données order_dw_mysql sont synchronisées.

    order_dw_mysql.\.*

    server-id

    Identifiant numérique unique pour la connexion client à la base de données.

    8601-8604

    sink

    jdbc-url

    URL de connexion JDBC.

    Spécifiez l'adresse IP et le port de requête du Frontend (FE) au format jdbc:mysql://ip:port.

    Dans l'onglet Instance Details de la console E-MapReduce, vous pouvez consulter le internal endpoint et le query port du FE pour l'instance cible.

    jdbc:mysql://fe-c-b76b6aa51807****-internal.starrocks.aliyuncs.com:9030

    load-url

    URL du service HTTP utilisée pour se connecter au nœud FE.

    Dans l'onglet Instance Details de la console E-MapReduce, vous pouvez consulter le internal endpoint et le HTTP port du FE pour l'instance cible.

    fe-c-b76b6aa51807****-internal.starrocks.aliyuncs.com:8030

    username

    Identifiants de connexion à StarRocks.

    Utilisez les identifiants configurés lors de la création de l'instance StarRocks.

    Remarque

    Cet exemple utilise des variables pour éviter d'exposer les identifiants en clair. Pour plus d'informations, consultez la rubrique Gérer les variables.

    ${secret_values.starrocksusername}

    password

    ${secret_values.starrockspassword}

    sink.buffer-flush.interval-ms

    Intervalle de vidage du tampon interne.

    Un intervalle court (5 secondes) est utilisé car cet exemple traite un faible volume de données, ce qui permet de visualiser rapidement les résultats.

    5000

    route

    source-table

    Table(s) source à router.

    Vous pouvez utiliser une expression régulière pour faire correspondre plusieurs tables. Par exemple, order_dw_mysql.\.* route toutes les tables de la base de données order_dw_mysql.

    order_dw_mysql.\.*

    sink-table

    Modèle de table de destination pour les données routées.

    Vous pouvez utiliser le symbole défini dans le paramètre replace-symbol comme espace réservé pour chaque nom de table source afin d'effectuer un routage plusieurs-à-plusieurs.

    Pour plus d'informations sur les règles de routage, consultez la rubrique Module de routage.

    order_dw_sr.<>

    replace-symbol

    Espace réservé pour le nom de la table source utilisé dans la mise en correspondance de motifs.

    <>

  7. Cliquez sur Deploy.

Étape 3 : Démarrer la tâche Flink CDC

  1. Sur la page Data Ingestion, cliquez sur Deploy, puis sur OK dans la boîte de dialogue qui s'affiche.

  2. Sur la page O&M > Deployments, localisez la tâche YAML cible et cliquez sur Start dans la colonne Actions.

  3. Cliquez sur Start.

    Dans cet exemple, sélectionnez Initial Mode. Pour plus d'informations sur les paramètres, consultez la rubrique Démarrer une tâche. Une fois la tâche démarrée, vous pouvez surveiller son état sur la page Deployments.

Étape 4 : Vérifier les résultats dans StarRocks

Une fois que la tâche atteint l'état RUNNING, vous pouvez vérifier les données dans StarRocks.

  1. Se connecter à une instance EMR Serverless StarRocks via EMR StarRocks Manager.

  2. Dans le volet de navigation de gauche, cliquez sur SQL Editor. Sous l'onglet Database, cliquez sur l'icône d'actualisation image.

    Une base de données nommée order_dw_sr apparaît sous default_catalog.

  3. Sous l'onglet Query List, cliquez sur + File pour créer un fichier Scripts. Saisissez les instructions SQL suivantes, puis cliquez sur Run.

    SELECT * FROM default_catalog.order_dw_sr.orders order by order_id;
    SELECT * FROM default_catalog.order_dw_sr.orders_pay order by pay_id;
    SELECT * FROM default_catalog.order_dw_sr.product_catalog order by product_id;
  4. Consultez les résultats affichés sous les commandes.

    Les résultats indiquent que les tables et les données de la base de données MySQL existent désormais dans StarRocks.

    Les tables synchronisées incluent default_catalog.order_dw_sr.orders, default_catalog.order_dw_sr.orders_pay et default_catalog.order_dw_sr.product_catalog. Exécutez des instructions SELECT sur chaque table pour vérifier l'intégrité des données.

Documentation connexe