Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Flink CDC : synchroniser une base de données MySQL entière vers Kafka

Dernière mise à jour :Aug 13, 2026

Cette rubrique explique comment synchroniser une base de données MySQL entière vers Kafka. Cette approche réduit la charge que plusieurs jobs imposent à la base de données MySQL.

Contexte

Une table source MySQL CDC capture les données de MySQL et synchronise les modifications en temps réel de la table. Ce mécanisme est courant dans les scénarios de calcul complexes, par exemple lorsqu'une table sert de table de dimension dans une opération JOIN avec d'autres tables de données. Une seule table MySQL peut constituer une dépendance pour plusieurs jobs. Lorsque plusieurs jobs traitent des données provenant de la même table MySQL, la base de données ouvre plusieurs connexions, ce qui exerce une pression importante sur le serveur MySQL et le réseau.

Fonctionnement

Pour alléger la pression sur la base de données MySQL en amont, Realtime Compute for Apache Flink permet de synchroniser une base de données MySQL entière vers Kafka. Cette solution introduit Kafka comme couche intermédiaire et utilise un job d'ingestion de données Flink CDC pour synchroniser les données vers Kafka.

Un seul job synchronise en temps réel les données d'une base de données MySQL en amont vers Kafka. Chaque table MySQL est écrite dans un topic Kafka correspondant en mode upsert. Les jobs en aval utilisent ensuite un connecteur Kafka upsert pour lire les données des topics au lieu d'accéder directement aux tables MySQL. Cette méthode réduit efficacement la pression exercée par plusieurs jobs sur la base de données MySQL.

图片 1

Limites

  • Chaque table MySQL que vous synchronisez doit posséder une clé primaire.

  • Vous pouvez utiliser des clusters Kafka auto-gérés, des clusters EMR Kafka ou ApsaraMQ for Kafka. Si vous utilisez ApsaraMQ for Kafka, la connexion s'effectue uniquement via un endpoint par défaut.

  • L'espace de stockage du cluster Kafka doit être supérieur à celui des tables sources. Dans le cas contraire, des données risquent d'être perdues en raison d'un espace insuffisant. Les topics créés pour la synchronisation de bases de données sont des topics compactés. Dans un topic compacté, seul le dernier message de chaque clé de message est conservé, mais les données n'expirent jamais. Cela signifie que le topic compacté stocke un volume de données globalement équivalent à la taille de la table source.

Scénario d'exemple

Prenons l'exemple d'un scénario d'analyse de vérification de commandes en temps réel comportant trois tables : une table utilisateur (user), une table de commandes (order) et une table de retours utilisateurs (feedback). Ces tables contiennent les données illustrées dans la figure suivante.mysql database

Pour afficher les informations de commande et les avis des utilisateurs, vous devez joindre la table user afin de récupérer les noms d'utilisateur depuis le champ name. L'exemple SQL suivant illustre cette opération.

-- Join order information with the user table to display the username and product name for each order.
SELECT order.id as order_id, product, user.name as user_name
FROM order LEFT JOIN user
ON order.user_id = user.id;
-- Join reviews with the user table to display the content of each review and the corresponding username.
SELECT feedback.id as feedback_id, comment, user.name as user_name
FROM feedback LEFT JOIN user
ON feedback.user_id = user.id;

Les deux jobs SQL précédents utilisent tous deux la table user. Lors de l'exécution, chaque job lit les données complètes et incrémentielles depuis MySQL. Une lecture complète nécessite la création d'une connexion MySQL, tandis qu'une lecture incrémentielle requiert la création d'un client Binlog. À mesure que le nombre de jobs augmente, la demande en ressources de connexion MySQL et de clients Binlog croît également, exerçant une pression considérable sur la base de données en amont. Pour soulager cette pression, utilisez un job d'ingestion de données Flink CDC afin de synchroniser en temps réel les données de la base de données MySQL en amont vers Kafka, permettant ainsi leur consommation par plusieurs jobs en aval.

Prérequis

Préparation

Préparer la source de données MySQL

  1. Créez une base de données ApsaraDB RDS for MySQL. Pour plus d'informations, consultez Créer une base de données.

    Créez une base de données nommée order_dw pour l'instance cible.

  2. Préparez la source de données MySQL CDC.

    1. Sur la page des détails de l'instance, cliquez sur Log on to Database en haut de la page.

    2. Dans la boîte de dialogue de connexion DMS qui s'affiche, saisissez le nom d'utilisateur et le mot de passe du compte de base de données créé, puis cliquez sur Login.

    3. Après la connexion, double-cliquez sur la base de données order_dw dans le volet de gauche pour changer de base de données.

    4. Dans la console SQL, saisissez les instructions DDL pour créer les trois tables métier ainsi que les instructions d'insertion de données.

      CREATE TABLE `user` (
        id bigint not null primary key,
        name varchar(50) not null
      );
      CREATE TABLE `order` (
        id bigint not null primary key,
        product varchar(50) not null,
        user_id bigint not null
      );
      CREATE TABLE `feedback` (
        id bigint not null primary key,
        user_id bigint not null,
        comment varchar(50) not null
      );
      -- Prepare data
      INSERT INTO `user` VALUES(1, 'Tom'),(2, 'Jerry');
      INSERT INTO `order` VALUES
      (1, 'Football', 2),
      (2, 'Basket', 1);
      INSERT INTO `feedback` VALUES
      (1, 1, 'Good.'),
      (2, 2, 'Very good');
  3. Cliquez sur Execute, puis cliquez sur Execute.

Procédure

  1. Créez et démarrez un job d'ingestion de données Flink CDC pour synchroniser en temps réel les données de la base de données MySQL en amont vers Kafka, afin qu'elles soient consommées par plusieurs jobs en aval. Le job de synchronisation de base de données crée automatiquement les topics. Vous pouvez définir les noms des topics à l'aide du module route. Les topics utilisent les paramètres par défaut du cluster Kafka pour le nombre de partitions et de réplicas, et le paramètre cleanup.policy est défini sur compact.

    Noms de topics par défaut

    Par défaut, les topics Kafka créés par le job de synchronisation de base de données suivent le format de nommage {database_name}.{table_name}. Le job suivant crée trois topics : order_dw.user, order_dw.order et order_dw.feedback.

    1. Sur la page Development > Data Ingestion, créez un job d'ingestion de données Flink CDC et copiez le code suivant dans l'éditeur YAML.

      source:
        type: mysql
        name: MySQL Source
        hostname: #{hostname}
        port: 3306
        username: #{usernmae}
        password: #{password}
        tables: order_dw.\.*
        server-id: 28601-28604
        # (Optional) Synchronize data from newly created tables during the incremental phase.
        scan.binlog.newly-added-table.enabled: true
        # (Optional) Synchronize table and field comments.
        include-comments.enabled: true
        # (Optional) Prioritize dispatching unbounded splits to avoid potential TaskManager OutOfMemory issues.
        scan.incremental.snapshot.unbounded-chunk-first.enabled: true
        # (Optional) Enable parsing filters to speed up reads.
        scan.only.deserialize.captured.tables.changelog.enabled: true
      sink:
        type: upsert-kafka
        name: upsert-kafka Sink
        properties.bootstrap.servers: xxxx.alikafka.aliyuncs.com:9092
        # The following parameters are required for ApsaraMQ for Kafka.
        aliyun.kafka.accessKeyId: #{ak}
        aliyun.kafka.accessKeySecret: #{sk}
        aliyun.kafka.instanceId: #{instanceId}
        aliyun.kafka.endpoint: #{endpoint}
        aliyun.kafka.regionId: #{regionId}
    2. Dans le coin supérieur droit, cliquez sur Deploy pour déployer le job.

    3. Dans la barre de navigation de gauche, cliquez sur O&M > Deployments. Dans la colonne Actions du job cible, cliquez sur Start, sélectionnez Initial Mode, puis cliquez sur Start.

    Noms de topics par table

    Le module route permet de spécifier un nom de topic pour chaque table. Le job suivant crée trois topics : user1, order2 et feedback3.

    1. Sur la page Development > Data Ingestion, créez un job d'ingestion de données Flink CDC et copiez le code suivant dans l'éditeur YAML.

      source:
        type: mysql
        name: MySQL Source
        hostname: #{hostname}
        port: 3306
        username: #{usernmae}
        password: #{password}
        tables: order_dw.\.*
        server-id: 28601-28604
        # (Optional) Synchronize data from newly created tables during the incremental phase.
        scan.binlog.newly-added-table.enabled: true
        # (Optional) Synchronize table and field comments.
        include-comments.enabled: true
        # (Optional) Prioritize dispatching unbounded splits to avoid potential TaskManager OutOfMemory issues.
        scan.incremental.snapshot.unbounded-chunk-first.enabled: true
        # (Optional) Enable parsing filters to speed up reads.
        scan.only.deserialize.captured.tables.changelog.enabled: true
      route:
        - source-table: order_dw.user
          sink-table: user1
        - source-table: order_dw.order
          sink-table: order2
        - source-table: order_dw.feedback
          sink-table: feedback3
      sink:
        type: upsert-kafka
        name: upsert-kafka Sink
        properties.bootstrap.servers: xxxx.alikafka.aliyuncs.com:9092
        # The following parameters are required for ApsaraMQ for Kafka.
        aliyun.kafka.accessKeyId: #{ak}
        aliyun.kafka.accessKeySecret: #{sk}
        aliyun.kafka.instanceId: #{instanceId}
        aliyun.kafka.endpoint: #{endpoint}
        aliyun.kafka.regionId: #{regionId}
    2. Dans le coin supérieur droit, cliquez sur Deploy pour déployer le job.

    3. Dans la barre de navigation de gauche, sélectionnez O&M > Deployments, cliquez sur Start dans la colonne Actions du job cible, sélectionnez Initial Mode, puis cliquez sur Start.

    Noms de topics par lot

    Le module route permet également de définir un modèle pour les noms des topics générés. Le job suivant crée trois topics : topic_user, topic_order et topic_feedback.

    1. Sur la page Development > Data Ingestion, créez un job d'ingestion de données Flink CDC et copiez le code suivant dans l'éditeur YAML.

      source:
        type: mysql
        name: MySQL Source
        hostname: #{hostname}
        port: 3306
        username: #{usernmae}
        password: #{password}
        tables: order_dw.\.*
        server-id: 28601-28604
        # (Optional) Synchronize data from newly created tables during the incremental phase.
        scan.binlog.newly-added-table.enabled: true
        # (Optional) Synchronize table and field comments.
        include-comments.enabled: true
        # (Optional) Prioritize dispatching unbounded splits to avoid potential TaskManager OutOfMemory issues.
        scan.incremental.snapshot.unbounded-chunk-first.enabled: true
        # (Optional) Enable parsing filters to speed up reads.
        scan.only.deserialize.captured.tables.changelog.enabled: true
      route:
        - source-table: order_dw.\.*
          sink-table: topic_<>
          replace-symbol: <>
      sink:
        type: upsert-kafka
        name: upsert-kafka Sink
        properties.bootstrap.servers: xxxx.alikafka.aliyuncs.com:9092
        # The following parameters are required for ApsaraMQ for Kafka.
        aliyun.kafka.accessKeyId: #{ak}
        aliyun.kafka.accessKeySecret: #{sk}
        aliyun.kafka.instanceId: #{instanceId}
        aliyun.kafka.endpoint: #{endpoint}
        aliyun.kafka.regionId: #{regionId}
    2. Dans le coin supérieur droit, cliquez sur Deploy pour déployer le job.

    3. Dans la barre de navigation de gauche, cliquez sur O&M > Deployments. Cliquez sur Start dans la colonne Actions du job cible, sélectionnez Initial Mode, puis cliquez sur Start.

  1. Consommez les données Kafka en temps réel.

    Le job d'ingestion écrit les données de la base de données MySQL en amont vers Kafka au format JSON. Plusieurs jobs en aval peuvent alors consommer les données d'un seul topic pour récupérer l'état le plus récent des tables de la base de données. Vous pouvez consommer les données des tables synchronisées vers Kafka selon l'une des méthodes suivantes :

    Via un catalog

    Lisez les données d'un topic Kafka en l'utilisant comme table source.

    1. Sur la page Development > ETL, créez un job SQL en streaming et copiez le code suivant dans l'éditeur SQL.

      CREATE TEMPORARY TABLE print_user_proudct(
        order_id BIGINT,
        product STRING,
        user_name STRING
      ) WITH (
        'connector'='print',
        'logger'='true'
      );
      CREATE TEMPORARY TABLE print_user_feedback(
        feedback_id BIGINT,
        `comment` STRING,
        user_name STRING
      ) WITH (
        'connector'='print',
        'logger'='true'
      );
      BEGIN STATEMENT SET;      -- Required when writing to multiple sinks.
      -- Join order information with the user table in the Kafka JSON Catalog to display the username and product name for each order.
      INSERT INTO print_user_proudct
      SELECT `order`.key_id as order_id, value_product as product, `user`.value_name as user_name
      FROM `kafka-catalog`.`kafka`.`order`/*+OPTIONS('properties.group.id'='<yourGroupName>', 'scan.startup.mode'='earliest-offset')*/  as `order` -- Specify the group and startup mode.
      LEFT JOIN `kafka-catalog`.`kafka`.`user`/*+OPTIONS('properties.group.id'='<yourGroupName>', 'scan.startup.mode'='earliest-offset')*/ as `user` -- Specify the group and startup mode.
      ON `order`.value_user_id = `user`.key_id;
      -- Join reviews with the user table to display the content of each review and the corresponding username.
      INSERT INTO print_user_feedback
      SELECT feedback.key_id as feedback_id, value_comment as `comment`, `user`.value_name as user_name
      FROM `kafka-catalog`.`kafka`.feedback/*+OPTIONS('properties.group.id'='<yourGroupName>', 'scan.startup.mode'='earliest-offset')*/  as feedback  -- Specify the group and startup mode.
      LEFT JOIN `kafka-catalog`.`kafka`.`user`/*+OPTIONS('properties.group.id'='<yourGroupName>', 'scan.startup.mode'='earliest-offset')*/ as `user` -- Specify the group and startup mode.
      ON feedback.value_user_id = `user`.key_id;
      END;      -- Required when writing to multiple sinks.

      Cet exemple utilise le connecteur Print pour imprimer directement les résultats. Vous pouvez également diriger les résultats vers une table de résultats utilisant un autre connecteur pour une analyse plus poussée. Pour plus d'informations sur la syntaxe d'écriture vers plusieurs sinks, consultez Instruction INSERT INTO.

      Remarque

      Lorsque vous utilisez directement cette méthode, des modifications de schéma peuvent survenir. Par conséquent, le schéma analysé par le catalog JSON Kafka peut différer du schéma de la table MySQL correspondante. Par exemple, des champs supprimés peuvent encore apparaître et certains champs peuvent présenter des valeurs nulles.

      Le schéma lu depuis le catalog se compose des champs issus des données consommées. Si un champ est supprimé mais que ses messages n'ont pas expiré, le champ peut encore apparaître avec une valeur nulle. Aucun traitement particulier n'est nécessaire dans ce cas.

    2. Dans le coin supérieur droit, cliquez sur Deploy pour déployer le job.

    3. Dans la barre de navigation de gauche, cliquez sur O&M > Deployments, cliquez sur Start dans la colonne Actions du job cible, sélectionnez Initial Mode, puis cliquez sur Start.

    Via une table temporaire

    Définissez un schéma personnalisé et lisez les données depuis une table temporaire.

    1. Sur la page Development > ETL, créez un job SQL en streaming et copiez le code suivant dans l'éditeur SQL.

      CREATE TEMPORARY TABLE user_source (
        key_id BIGINT,
        value_name STRING
      ) WITH (
        'connector' = 'kafka',
        'topic' = 'user',
        'properties.bootstrap.servers' = '<yourKafkaBrokers>',
        'scan.startup.mode' = 'earliest-offset',
        'key.format' = 'json',
        'value.format' = 'json',
        'key.fields' = 'key_id',
        'key.fields-prefix' = 'key_',
        'value.fields-prefix' = 'value_',
        'value.fields-include' = 'EXCEPT_KEY',
        'value.json.infer-schema.flatten-nested-columns.enable' = 'false',
        'value.json.infer-schema.primitive-as-string' = 'false'
      );
      CREATE TEMPORARY TABLE order_source (
        key_id  BIGINT,
        value_product STRING,
        value_user_id BIGINT  
      ) WITH (
        'connector' = 'kafka',
        'topic' = 'order',
        'properties.bootstrap.servers' = '<yourKafkaBrokers>',
        'scan.startup.mode' = 'earliest-offset',
        'key.format' = 'json',
        'value.format' = 'json',
        'key.fields' = 'key_id',
        'key.fields-prefix' = 'key_',
        'value.fields-prefix' = 'value_',
        'value.fields-include' = 'EXCEPT_KEY',
        'value.json.infer-schema.flatten-nested-columns.enable' = 'false',
        'value.json.infer-schema.primitive-as-string' = 'false'
      );
      CREATE TEMPORARY TABLE feedback_source (
        key_id  BIGINT,
        value_user_id BIGINT,
        value_comment STRING
      ) WITH (
        'connector' = 'kafka',
        'topic' = 'feedback',
        'properties.bootstrap.servers' = '<yourKafkaBrokers>',
        'scan.startup.mode' = 'earliest-offset',
        'key.format' = 'json',
        'value.format' = 'json',
        'key.fields' = 'key_id',
        'key.fields-prefix' = 'key_',
        'value.fields-prefix' = 'value_',
        'value.fields-include' = 'EXCEPT_KEY',
        'value.json.infer-schema.flatten-nested-columns.enable' = 'false',
        'value.json.infer-schema.primitive-as-string' = 'false'
      );
      CREATE TEMPORARY TABLE print_user_proudct(
        order_id BIGINT,
        product STRING,
        user_name STRING
      ) WITH (
        'connector'='print',
        'logger'='true'
      );
      CREATE TEMPORARY TABLE print_user_feedback(
        feedback_id BIGINT,
        `comment` STRING,
        user_name STRING
      ) WITH (
        'connector'='print',
        'logger'='true'
      );
      BEGIN STATEMENT SET;      -- Required when writing to multiple sinks.
      -- Join order information with the user table from the Kafka JSON Catalog to display the username and product name for each order.
      INSERT INTO print_user_proudct
      SELECT order_source.key_id as order_id, value_product as product, user_source.value_name as user_name
      FROM order_source LEFT JOIN user_source
      ON order_source.value_user_id = user_source.key_id;
      -- Join reviews with the user table to display the content of each review and the corresponding username.
      INSERT INTO print_user_feedback
      SELECT feedback_source.key_id as feedback_id, value_comment as `comment`, user_source.value_name as user_name
      FROM feedback_source  LEFT JOIN user_source
      ON feedback_source.value_user_id = user_source.key_id;
      END;      -- Required when writing to multiple sinks.

      Cet exemple utilise le connecteur Print pour imprimer directement les résultats. Vous pouvez également diriger les résultats vers une table de résultats d'un autre connecteur pour une analyse plus poussée. Pour plus d'informations sur la syntaxe d'écriture vers plusieurs sinks, consultez Instruction INSERT INTO.

      Le tableau suivant décrit les paramètres de configuration de la table temporaire.

      Paramètre

      Description

      Notes

      connector

      Type de connecteur.

      Définissez la valeur sur kafka.

      topic

      Nom du topic correspondant.

      Doit correspondre à la description du catalog JSON Kafka.

      properties.bootstrap.servers

      Adresses des brokers Kafka.

      Le format est host:port,host:port,host:port, séparé par des virgules (,).

      scan.startup.mode

      Position de départ pour la lecture des données depuis Kafka.

      Valeurs valides :

      • earliest-offset : commence la lecture au premier offset disponible.

      • latest-offset : commence la lecture au dernier offset.

      • group-offsets (par défaut) : lit à partir de l'offset validé pour le groupe spécifié par properties.group.id.

      • timestamp : lit à partir de l'horodatage spécifié par scan.startup.timestamp-millis.

      • specific-offsets : commence la lecture aux offsets spécifiés dans scan.startup.specific-offsets.

      Note

      Ce paramètre prend effet lorsque le job démarre sans état sauvegardé. Lorsque le job redémarre ou récupère à partir d'un checkpoint, il privilégie la lecture depuis l'état sauvegardé.

      key.format

      Format utilisé par le connecteur Flink Kafka pour sérialiser ou désérialiser la clé du message Kafka.

      Définissez la valeur sur json.

      key.fields

      Champs de la table source ou de résultat correspondant à la clé du message Kafka.

      Utilisez des points-virgules (;) pour séparer plusieurs noms de champs. Par exemple, field1;field2.

      key.fields-prefix

      Préfixe personnalisé pour tous les champs de clé de message Kafka afin d'éviter les conflits de noms avec les champs de valeur de message ou les champs de métadonnées.

      Cette valeur doit correspondre à celle du paramètre key.fields-prefix du catalog JSON Kafka.

      value.format

      Format utilisé par le connecteur Flink Kafka pour sérialiser ou désérialiser la valeur du message Kafka.

      Définissez la valeur sur json.

      value.fields-prefix

      Préfixe personnalisé pour tous les champs de valeur de message Kafka afin d'éviter les conflits de noms avec les champs de clé de message ou les champs de métadonnées.

      Doit correspondre à la valeur du paramètre value.fields-prefix du catalog JSON Kafka.

      value.fields-include

      Politique de gestion des champs de clé de message dans la valeur du message.

      Définissez la valeur sur EXCEPT_KEY. Cela indique que la valeur du message n'inclut pas les champs de la clé du message.

      value.json.infer-schema.flatten-nested-columns.enable

      Indique s'il faut développer récursivement les colonnes JSON imbriquées dans la valeur du message Kafka.

      Valeur du paramètre infer-schema.flatten-nested-columns.enable du catalog correspondant.

      value.json.infer-schema.primitive-as-string

      Indique s'il faut déduire tous les types primitifs comme String dans la valeur du message Kafka.

      Valeur du paramètre infer-schema.primitive-as-string du catalog correspondant.

    2. Dans le coin supérieur droit, cliquez sur Deploy pour déployer le job.

    3. Dans la barre de navigation de gauche, cliquez sur O&M > Deployments, cliquez sur Start dans la colonne Actions du job cible, sélectionnez Initial Mode, puis cliquez sur Start.

  2. Consultez les résultats du job.

    1. Dans la barre de navigation de gauche, cliquez sur O&M > Deployments, puis cliquez sur le job cible.

    2. Sous l'onglet Logs, dans l'onglet Running Task Managers, cliquez sur la tâche dont vous souhaitez consulter le Path, ID.

    3. Cliquez sur Logs et recherchez les informations de journal liées à PrintSinkOutputWriter.

      Recherchez dans les journaux la sortie de PrintSinkOutputWriter. La sortie contient quatre enregistrements de données jointes : +I[1, Good., Tom], +I[2, Very good, Jerry], +I[2, Basket, Tom] et +I[1, Football, Jerry]. Cela indique que les jointures entre la table utilisateur et les tables de commandes et de retours ont réussi.

Documents associés