Tous les produits
Search
Centre de documentation

MaxCompute:Flink CDC open source pour l'écriture quasi temps réel dans les tables Delta

Dernière mise à jour :Aug 10, 2026

MaxCompute propose un connecteur Flink Change Data Capture (CDC) qui synchronise en temps réel les modifications de données provenant de sources telles que MySQL vers des tables standard MaxCompute ou des tables Delta. Utilisez le connecteur de pipeline Flink CDC pour configurer un pipeline d'extraction, de transformation et de chargement (ETL) en continu de MySQL vers MaxCompute.

Le connecteur prend en charge :

  • La création automatique de tables basée sur la structure de la table source

  • L'évolution du schéma : les modifications apportées au schéma de la table source sont propagées à MaxCompute

  • La synchronisation complète de la base de données sur plusieurs tables

  • Le routage des tables : mappez les tables sources vers différents noms de tables cibles, y compris la fusion de bases de données fragmentées

Prérequis

Avant de commencer, assurez-vous de disposer des éléments suivants :

  • Un projet MaxCompute. Pour en créer un, consultez la rubrique Créer un projet MaxCompute.

  • Docker et Docker Compose installés sur votre machine

  • Un environnement d'exécution Java (JRE) compatible avec Flink 1.18.0

Notes d'utilisation

  • Sélection du type de table : si la table source possède une clé primaire, une table Delta est créée automatiquement. En l'absence de clé primaire, une table standard MaxCompute est créée.

  • Comportement d'écriture pour les tables standard : les opérations DELETE sont ignorées. Les opérations UPDATE sont traitées comme des opérations INSERT.

  • Sémantique de livraison : seule la sémantique « au moins une fois » (at-least-once) est prise en charge. Pour mettre en œuvre des écritures idempotentes, utilisez une table Delta ; la clé primaire de la table Delta gère la déduplication.

  • Évolution du schéma : les modifications de schéma de la table source sont propagées à MaxCompute. Les nouvelles colonnes sont toujours ajoutées en dernière position. Les types de données des colonnes ne peuvent être modifiés que vers un type compatible. Pour plus de détails, consultez la rubrique Modifier le type de données d'une colonne.

  • Fusion de bases de données fragmentées : vous pouvez fusionner plusieurs tables sources en une seule table de destination à l'aide d'expressions régulières. La fusion de tables présentant des valeurs de clé primaire en double entre les sources n'est pas encore prise en charge.

Premiers pas

Ce tutoriel configure un travail ETL en continu complet qui synchronise les modifications de données de MySQL vers MaxCompute, couvrant la synchronisation complète de la base de données, l'évolution du schéma et le routage des tables.

Configurer les environnements

Configurer un cluster Flink en mode autonome

  1. Téléchargez le fichier flink-1.18.0-bin-scala_2.12.tgz et décompressez-le pour obtenir le répertoire flink-1.18.0. Dans le répertoire flink-1.18.0, exécutez la commande suivante pour définir la variable FLINK_HOME :

    export FLINK_HOME=$(pwd)
  2. Ouvrez le fichier $FLINK_HOME/conf/flink-conf.yaml avec un éditeur de texte et ajoutez les paramètres suivants :

    # Enable checkpointing. Run a checkpoint every 3 seconds.
    # For production, set the checkpoint interval to at least 30 seconds.
    execution.checkpointing.interval: 3000
    
    # The Flink CDC pipeline connector relies on Flink's communication mechanism
    # for data synchronization. Increase the timeout to prevent connection drops.
    pekko.ask.timeout: 60s
  3. Démarrez le cluster Flink :

    ./bin/start-cluster.sh

    Après le démarrage, ouvrez http://localhost:8081/ dans un navigateur pour accéder à l'interface web Flink. Exécutez start-cluster.sh plusieurs fois pour ajouter des TaskManagers.

Configurer un environnement MySQL

Ce tutoriel utilise Docker Compose pour exécuter MySQL.

  1. Créez un fichier nommé docker-compose.yaml avec le contenu suivant :

    Paramètre

    description

    version

    Version de Docker Compose

    image

    Version de l'image Docker. Définissez-la sur debezium/example-mysql:1.1

    ports

    Mappage du port MySQL

    environment

    Mot de passe root MySQL et identifiants utilisateur

    version: '2.1'
    services:
      mysql:
        image: debezium/example-mysql:1.1
        ports:
          - "3306:3306"
        environment:
          - MYSQL_ROOT_PASSWORD=123456
          - MYSQL_USER=mysqluser
          - MYSQL_PASSWORD=mysqlpw
  2. Dans le répertoire contenant le fichier docker-compose.yaml, démarrez les conteneurs :

    docker-compose up -d

    Exécutez docker ps pour confirmer que tous les conteneurs sont en cours d'exécution.

Préparer les données dans MySQL

  1. Connectez-vous au conteneur MySQL :

    docker-compose exec mysql mysql -uroot -p123456
  2. Créez la base de données app_db et renseignez trois tables :

    1. Créez la base de données :

      CREATE DATABASE app_db;
      USE app_db;
    2. Créez et renseignez la table orders :

      CREATE TABLE `orders` (
        `id` INT NOT NULL,
        `price` DECIMAL(10,2) NOT NULL,
        PRIMARY KEY (`id`)
      );
      
      INSERT INTO `orders` (`id`, `price`) VALUES (1, 4.00);
      INSERT INTO `orders` (`id`, `price`) VALUES (2, 100.00);
    3. Créez et renseignez la table shipments :

      CREATE TABLE `shipments` (
        `id` INT NOT NULL,
        `city` VARCHAR(255) NOT NULL,
        PRIMARY KEY (`id`)
      );
      
      INSERT INTO `shipments` (`id`, `city`) VALUES (1, 'beijing');
      INSERT INTO `shipments` (`id`, `city`) VALUES (2, 'xian');
    4. Créez et renseignez la table products :

      CREATE TABLE `products` (
        `id` INT NOT NULL,
        `product` VARCHAR(255) NOT NULL,
        PRIMARY KEY (`id`)
      );
      
      INSERT INTO `products` (`id`, `product`) VALUES (1, 'Beer');
      INSERT INTO `products` (`id`, `product`) VALUES (2, 'Cap');
      INSERT INTO `products` (`id`, `product`) VALUES (3, 'Peanut');

Soumettre un travail de pipeline YAML

  1. Téléchargez les packages JAR requis :

    • Package Flink CDC : téléchargez le fichier flink-cdc-3.1.1-bin.tar.gz et décompressez-le pour obtenir le répertoire flink-cdc-3.1.1. Déplacez le contenu des répertoires bin, lib, log et conf vers les répertoires correspondants sous flink-1.18.0.

    • Packages de connecteurs : téléchargez les fichiers JAR suivants et placez-les dans le répertoire flink-1.18.0/lib : > Remarque : les liens de téléchargement ci-dessus pointent uniquement vers des versions publiées. Pour utiliser une version SNAPSHOT, compilez le code source à partir de la branche master ou release localement.

    • Pilote JDBC : téléchargez le fichier mysql-connector-java-8.0.27.jar. Transmettez-le à l'interface CLI Flink CDC avec l'option --jar, ou placez-le dans $FLINK_HOME/lib et redémarrez le cluster. Les connecteurs CDC n'incluent plus les pilotes JDBC.

  2. Créez le fichier YAML du pipeline. L'exemple suivant (mysql-to-maxcompute.yaml) synchronise toutes les tables de la base de données app_db :

    Espace réservé

    description

    ${your_accessId}

    Votre ID AccessKey. Obtenez-le depuis la page Paire de clés d'accès.

    ${your_accessKey}

    Le secret AccessKey correspondant à votre ID AccessKey.

    ${your_maxcompute_endpoint}

    L'endpoint MaxCompute pour votre région et votre réseau. Consultez la rubrique Endpoints.

    ${your_project}

    Le nom de votre projet MaxCompute. Retrouvez-le dans la console MaxCompute sous Workspace > Projects.

    ################################################################################
    # Description: Sync MySQL all tables to MaxCompute
    ################################################################################
    source:
      type: mysql
      hostname: localhost
      port: 3306
      username: root
      password: 123456
      tables: app_db\.*
      server-id: 5400-5404
      server-time-zone: UTC
    
    # Configure the accessId, accessKey, endpoint, and project parameters.
    sink:
      type: maxcompute
      name: MaxComputeSink
      accessId: ${your_accessId}
      accessKey: ${your_accessKey}
      endpoint: ${your_maxcompute_endpoint}
      project: ${your_project}
      bucketsNum: 8
    
    pipeline:
      name: Sync MySQL Database to MaxCompute
      parallelism: 1

    Remplacez les espaces réservés suivants : Pour connaître tous les paramètres source disponibles, consultez la rubrique Connecteur MySQL. Pour connaître tous les paramètres de destination disponibles, consultez la section Configuration du connecteur.

  3. Soumettez le travail de pipeline au cluster Flink autonome :

    ./bin/flink-cdc.sh mysql-to-maxcompute.yaml

    Une soumission réussie renvoie une sortie similaire à :

    Pipeline has been submitted to cluster.
    Job ID: f9f9689866946e25bf151ecc179ef46f
    Job Description: Sync MySQL Database to MaxCompute

    L'interface web Flink affiche un travail en cours d'exécution nommé Sync MySQL Database to MaxCompute.

  4. Dans MaxCompute, vérifiez que les tables sont créées et que les données sont écrites :

    -- Query the orders table.
    read orders;
    
    -- Expected result:
    +------------+------------+
    | id         | price      |
    +------------+------------+
    | 1          | 4          |
    | 2          | 100        |
    +------------+------------+
    
    -- Query the shipments table.
    read shipments;
    
    -- Expected result:
    +------------+------------+
    | id         | city       |
    +------------+------------+
    | 1          | beijing    |
    | 2          | xian       |
    +------------+------------+
    
    -- Query the products table.
    read products;
    
    -- Expected result:
    +------------+------------+
    | id         | product    |
    +------------+------------+
    | 3          | Peanut     |
    | 1          | Beer       |
    | 2          | Cap        |
    +------------+------------+

Synchroniser les modifications en temps réel

Lorsque le pipeline est en cours d'exécution, les modifications effectuées dans MySQL sont reflétées dans MaxCompute en temps réel. Les exemples suivants utilisent la table orders pour illustrer les opérations d'insertion, de modification du schéma, de mise à jour et de suppression.

  1. Connectez-vous au conteneur MySQL :

    docker-compose exec mysql mysql -uroot -p123456
  2. Insérez une ligne dans la table orders :

    INSERT INTO app_db.orders (id, price) VALUES (3, 100.00);

    Exécutez read orders; dans MaxCompute. La nouvelle ligne apparaît :

    +------------+------------+
    | id         | price      |
    +------------+------------+
    | 3          | 100        |
    | 1          | 4          |
    | 2          | 100        |
    +------------+------------+
  3. Ajoutez une colonne à la table orders dans MySQL :

    ALTER TABLE app_db.orders ADD amount varchar(100) NULL;

    Exécutez read orders; dans MaxCompute. La nouvelle colonne amount est propagée automatiquement :

    +------------+------------+------------+
    | id         | price      | amount     |
    +------------+------------+------------+
    | 3          | 100        | NULL       |
    | 1          | 4          | NULL       |
    | 2          | 100        | NULL       |
    +------------+------------+------------+
  4. Mettez à jour une ligne dans la table orders :

    UPDATE app_db.orders SET price=100.00, amount=100.00 WHERE id=1;

    Exécutez read orders; dans MaxCompute :

    +------------+------------+------------+
    | id         | price      | amount     |
    +------------+------------+------------+
    | 3          | 100        | NULL       |
    | 1          | 100        | 100.00     |
    | 2          | 100        | NULL       |
    +------------+------------+------------+
  5. Supprimez une ligne de la table orders :

    DELETE FROM app_db.orders WHERE id=2;

    Exécutez read orders; dans MaxCompute :

    +------------+------------+------------+
    | id         | price      | amount     |
    +------------+------------+------------+
    | 3          | 100        | NULL       |
    | 1          | 100        | 100.00     |
    +------------+------------+------------+

Router les tables

Flink CDC permet de router les schémas et les données des tables sources vers différents noms de tables de destination. Utilisez cette fonctionnalité pour les migrations de base de données, les renommages de tables ou les fusions de bases de données fragmentées.

L'exemple YAML suivant route chaque table source vers une table cible portant un nom différent :

################################################################################
# Description: Sync MySQL all tables to MaxCompute
################################################################################
source:
  type: mysql
  hostname: localhost
  port: 3306
  username: root
  password: 123456
  tables: app_db\.*
  server-id: 5400-5404
  server-time-zone: UTC

# Configure the accessId, accessKey, endpoint, and project parameters.
sink:
  type: maxcompute
  name: MaxComputeSink
  accessId: ${your_accessId}
  accessKey: ${your_accessKey}
  endpoint: ${your_maxcompute_endpoint}
  project: ${your_project}
  bucketsNum: 8

route:
  - source-table: app_db.orders
    sink-table: ods_db.ods_orders
  - source-table: app_db.shipments
    sink-table: ods_db.ods_shipments
  - source-table: app_db.products
    sink-table: ods_db.ods_products

pipeline:
  name: Sync MySQL Database to MaxCompute
  parallelism: 1

Le champ source-table prend en charge les expressions régulières. Pour fusionner plusieurs tables fragmentées en une seule table cible :

route:
  - source-table: app_db.order\.*
    sink-table: ods_db.ods_orders

Cette configuration agrège les données de app_db.order01, app_db.order02, app_db.order03 et des tables similaires dans ods_db.ods_orders.

Les scénarios dans lesquels des valeurs de clé primaire en double existent entre plusieurs tables sources ne sont pas pris en charge et le seront dans les prochaines versions du connecteur.

Pour connaître tous les paramètres de routage, consultez la rubrique Route.

Nettoyer l'environnement

Une fois le tutoriel terminé, arrêtez tous les composants en cours d'exécution.

  1. Arrêtez les conteneurs Docker. Exécutez la commande suivante dans le répertoire contenant le fichier docker-compose.yaml :

    docker-compose down
  2. Arrêtez le cluster Flink. Exécutez la commande suivante dans le répertoire flink-1.18.0 :

    ./bin/stop-cluster.sh

Annexe

Configuration du connecteur

Le connecteur de destination MaxCompute prend en charge les paramètres de configuration suivants.

Les données sont vidées vers MaxCompute lorsque le tampon d'une partition ou d'un compartiment dépasse le seuil de taille configuré, ou lorsqu'un point de contrôle (checkpoint) est déclenché.

Paramètre

Obligatoire

Valeur par défaut

Type

description

type

Oui

String

Le type de connecteur. Définissez-le sur maxcompute.

name

Non

String

Le nom d'affichage de la destination.

accessId

Oui

String

Votre ID AccessKey. Obtenez-le depuis la page Paire de clés d'accès.

accessKey

Oui

String

Le secret AccessKey correspondant à votre ID AccessKey.

endpoint

Oui

String

L'endpoint MaxCompute. Configurez-le en fonction de votre région et de votre méthode de connexion réseau. Consultez la rubrique Endpoints.

project

Oui

String

Le nom du projet MaxCompute. Retrouvez-le dans la console MaxCompute sous Workspace > Projects.

tunnelEndpoint

Non

String

L'endpoint MaxCompute Tunnel. Le routage automatique est pris en charge dans la plupart des cas. Configurez ce paramètre uniquement lors de l'utilisation d'un environnement réseau spécial, tel qu'un proxy.

quotaName

Non

String

Le nom du groupe de ressources exclusif pour MaxCompute Tunnel. Si ce paramètre n'est pas défini, un groupe de ressources partagé est utilisé.

stsToken

Non

String

Requis lors de l'utilisation d'un jeton STS émis par un rôle RAM pour l'authentification.

bucketsNum

Non

16

Integer

Le nombre de compartiments pour la création automatique de table Delta. Consultez la rubrique Présentation de l'entrepôt de données quasi temps réel.

compressAlgorithm

Non

zlib

String

L'algorithme de compression pour les écritures. Valeurs valides : raw (aucune compression), zlib, snappy.

totalBatchSize

Non

64MB

String

La taille maximale du tampon par partition ou table non partitionnée. Lorsque ce seuil est dépassé, les données sont vidées vers MaxCompute. Chaque partition ou table non partitionnée conserve un tampon indépendant.

bucketBatchSize

Non

4MB

String

La taille maximale du tampon par compartiment. S'applique uniquement aux écritures de table Delta. Lorsque ce seuil est dépassé, les données sont vidées vers MaxCompute. Chaque compartiment conserve un tampon indépendant.

numCommitThreads

Non

16

Integer

Le nombre maximal de partitions ou de tables traitées simultanément lors d'un point de contrôle.

numFlushConcurrent

Non

4

Integer

Le nombre maximal de compartiments écrits simultanément dans MaxCompute. S'applique uniquement aux écritures de table Delta.

retryTimes

Non

3

Integer

Le nombre maximal de tentatives en cas d'erreurs réseau.

sleepMillis

Non

Long

La durée d'attente entre les tentatives en cas d'erreurs réseau. Unité : millisecondes.

Mappages d'emplacement de table

Lorsque le connecteur crée automatiquement une table, il mappe les informations d'emplacement de la source vers MaxCompute comme suit.

Important

Si le projet MaxCompute ne prend pas en charge les modèles de schéma, chaque tâche de synchronisation ne peut synchroniser qu'une seule base de données MySQL. Le connecteur ignore les informations TableId.namespace pour les autres sources de données.

Objet dans Flink CDC

Emplacement dans MaxCompute

Emplacement dans MySQL

Projet (issu du fichier de configuration)

Projet

TableId.namespace

Schéma (nécessite la prise en charge du modèle de schéma ; ignoré si le projet ne prend pas en charge les modèles de schéma)

Base de données

TableId.tableName

Table

Table

Mappages de types de données

Le tableau suivant montre comment les types de données Flink sont mappés vers les types de données MaxCompute.

Type de données Flink

Type de données MaxCompute

notes

CHAR/VARCHAR

STRING

MaxCompute stocke tous les types de chaîne sous forme de STRING.

BOOLEAN

BOOLEAN

BINARY/VARBINARY

BINARY

DECIMAL

DECIMAL

TINYINT

TINYINT

SMALLINT

SMALLINT

INTEGER

INTEGER

BIGINT

BIGINT

FLOAT

FLOAT

DOUBLE

DOUBLE

TIME_WITHOUT_TIME_ZONE

STRING

MaxCompute ne dispose pas de type TIME natif. Les valeurs temporelles sont stockées sous forme de chaînes.

DATE

DATE

TIMESTAMP_WITHOUT_TIME_ZONE

TIMESTAMP_NTZ

Stocké sans information de fuseau horaire.

TIMESTAMP_WITH_LOCAL_TIME_ZONE

TIMESTAMP

TIMESTAMP_WITH_TIME_ZONE

TIMESTAMP

ARRAY

ARRAY

MAP

MAP

ROW

STRUCT