Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Develop a Flink CDC data ingestion job

Dernière mise à jour :Aug 09, 2026

Realtime Compute for Apache Flink vous permet de créer des jobs d’ingestion de données Flink CDC à l’aide d’un fichier YAML pour synchroniser les données d’une source vers un puits. Cette rubrique décrit les étapes de développement d’un job d’ingestion de données Flink CDC.

Contexte

Les configurations YAML vous permettent de définir aisément des pipelines ETL complexes, automatiquement convertis en logique d’exécution Flink. L’ingestion de données Flink CDC propose une solution d’intégration puissante basée sur Flink CDC. Cette solution prend efficacement en charge la synchronisation complète de bases de données, la synchronisation de tables uniques, la synchronisation de bases et de tables partitionnées, la découverte automatique de nouvelles tables, la gestion des modifications de schéma ainsi que les colonnes calculées personnalisées. Elle permet également le traitement ETL, le filtrage via la clause WHERE et l’élagage de colonnes. Cette approche déclarative simplifie considérablement le processus d’intégration des données et améliore l’efficacité et la fiabilité.

Avantages de Flink CDC

Dans Realtime Compute for Apache Flink, vous pouvez développer un job d’ingestion de données Flink CDC, un job SQL ou un job DataStream pour synchroniser les données. Les sections suivantes détaillent les avantages de l’utilisation d’un job d’ingestion de données Flink CDC par rapport aux deux autres options.

Flink CDC par rapport à Flink SQL

Les jobs d’ingestion de données Flink CDC et les jobs SQL utilisent des types de données différents pour la transmission :

  • Les jobs SQL transmettent les données sous forme de RowData. Chaque objet RowData possède son propre type de modification, qui comprend quatre catégories principales : insertion (+I), mise à jour avant (-U), mise à jour après (+U) et suppression (-D).

  • Flink CDC utilise SchemaChangeEvent pour transmettre les informations relatives aux modifications de schéma, telles que la création d’une table, l’ajout d’une colonne ou la troncature d’une table. Il utilise DataChangeEvent pour transmettre les modifications de données, comme les insertions, mises à jour et suppressions. Un message de mise à jour contient à la fois le contenu avant et après, ce qui vous permet d’écrire les données de modification originales dans le puits.

Le tableau suivant répertorie les avantages des jobs d’ingestion de données Flink CDC par rapport aux jobs SQL.

Ingestion Flink CDC

Flink SQL

Détection automatique du schéma et synchronisation complète de la base de données

Nécessite des instructions CREATE TABLE et INSERT manuelles

Prend en charge plusieurs politiques de modification de schéma

Ne prend pas en charge les modifications de schéma

Préserve le journal des modifications original

Altère la structure du journal des modifications original

Prend en charge la lecture et l’écriture sur plusieurs tables

Lecture et écriture sur une seule table

Par rapport aux instructions CTAS/CDAS, les jobs Flink CDC offrent des fonctionnalités plus puissantes, notamment :

  • Synchronisation immédiate des modifications de schéma en amont, sans attendre que de nouvelles écritures de données déclenchent la synchronisation.

  • Préservation du journal des modifications original, garantissant que les messages de mise à jour ne sont pas fragmentés.

  • Synchronisation de davantage de types de modifications de schéma, tels que TRUNCATE TABLE et DROP TABLE.

  • Mappage flexible des tables et définition des noms de tables de destination.

  • Comportements d’évolution de schéma flexibles et configurables.

  • Filtrage des données à l’aide de clauses WHERE.

  • Prise en charge de l’élagage de colonnes.

Flink CDC par rapport à Flink DataStream

Le tableau suivant présente les avantages des jobs d’ingestion de données Flink CDC par rapport aux jobs DataStream.

Ingestion Flink CDC

Flink DataStream

Accessible aux utilisateurs de tous niveaux de compétence.

Nécessite une expertise en Java et en systèmes distribués.

Masque la complexité sous-jacente et simplifie le développement.

Nécessite une connaissance du framework Flink.

Format YAML facile à apprendre.

Nécessite la maîtrise d’outils comme Maven pour la gestion des dépendances.

Réutilisabilité élevée des jobs existants.

Réutilisation difficile du code existant.

Limites

  • Utilisez Ververica Runtime (VVR) 11.1 ou version ultérieure pour développer des jobs d’ingestion de données Flink CDC. Si vous devez utiliser VVR 8.x, utilisez VVR 8.0.11.

  • Chaque job prend en charge une seule source et un seul puits. Pour lire depuis plusieurs sources ou écrire dans plusieurs puits, vous devez créer plusieurs jobs Flink CDC.

  • Les jobs Flink CDC ne peuvent pas être déployés sur un cluster de session.

  • Les jobs d’ingestion de données Flink CDC ne prennent pas en charge le réglage automatique.

Connecteurs d’ingestion de données Flink CDC

Pour obtenir des détails sur les connecteurs source et puits pris en charge pour l’ingestion de données Flink CDC, consultez Connecteurs pris en charge.

Créer un job d’ingestion de données Flink CDC

À partir d’un modèle

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

  2. Dans la colonne Actions de l’espace de travail cible, cliquez sur Console.

  3. Dans le volet de navigation de gauche, choisissez Development > Data Ingestion.

  4. Cliquez sur image, puis sur New Draft with Template.

  5. Sélectionnez un modèle de synchronisation de données.

    Actuellement, seuls les modèles MySQL vers StarRocks, MySQL vers Paimon et MySQL vers Hologres sont disponibles.

  6. Saisissez les informations du job, telles que OK, location et engine version, puis cliquez sur OK.

  7. Configurez les informations de source et de puits pour le job Flink CDC.

    Pour plus de détails sur la configuration des paramètres, consultez la documentation relative aux connecteurs concernés.

À partir d’un job CTAS/CDAS

Important
  • Si un job contient plusieurs instructions CXAS, Flink détecte et convertit uniquement la première.

  • En raison des différences de prise en charge des fonctions intégrées entre Flink SQL et Flink CDC, les règles transform générées peuvent ne pas fonctionner telles quelles. Vous devez les examiner et les ajuster si nécessaire.

  • Si la source est MySQL et qu’un job CTAS/CDAS original est toujours en cours d’exécution, vous devez ajuster le server-id de la source dans le job d’ingestion de données Flink CDC pour éviter les conflits.

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

  2. Dans la colonne Actions de l’espace de travail cible, cliquez sur Console.

  3. Dans le volet de navigation de gauche, choisissez Development > Data Ingestion.

  4. Cliquez sur image, puis sur New Draft from CTAS/CDAS. Sélectionnez le job CTAS ou CDAS cible, puis cliquez sur OK.

    Sur la page de sélection, le système affiche uniquement les jobs CTAS et CDAS valides. Les jobs ETL standard et les brouillons contenant des erreurs de syntaxe ne sont pas affichés.

  5. Saisissez les informations du job, telles que OK, location et engine version, puis cliquez sur OK.

À partir de Flink CDC open source

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

  2. Dans la colonne Actions de l’espace de travail cible, cliquez sur Console.

  3. Dans le volet de navigation de gauche, choisissez Development > Data Ingestion.

  4. Cliquez sur image, puis sélectionnez New Draft. Saisissez le name et la engine version, puis cliquez sur Create.

  5. Collez le code YAML de votre job Flink CDC open source dans l’éditeur.

  6. (Facultatif) Cliquez sur Validate.

    Cette option vérifie les erreurs de syntaxe, les problèmes de connectivité réseau et les problèmes d’autorisation.

À partir de zéro

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

  2. Dans la colonne Actions de l’espace de travail cible, cliquez sur Console.

  3. Dans le volet de navigation de gauche, choisissez Development > Data Ingestion.

  4. Cliquez sur image, puis sélectionnez New Draft. Saisissez le name et la engine version, puis cliquez sur Create.

  5. Configurez le job Flink CDC à l’aide de YAML. Exemple :

    # Required
    source:
      # The type of the data source.
      type: <Replace with your source connector type>
      # The configuration for the data source. For details about configuration items, see the documentation for the corresponding connector.
      ...
    # Required
    sink:
      # The type of the sink.
      type: <Replace with your sink connector type>
      # The configuration for the sink. For details about configuration items, see the documentation for the corresponding connector.
      ...
    # Optional
    transform:
      # A transform rule for the flink_test.customers table.
      - source-table: flink_test.customers
        # Projection configuration. Specifies the columns to synchronize and performs data transformations.
        projection: id, username, UPPER(username) as username1, age, (age + 1) as age1, test_col1, __schema_name__ || '.' || __table_name__ identifier_name
        # Filter condition. Synchronizes only data where id is greater than 10.
        filter: id > 10
        # Description of the transform rule.
        description: append calculated columns based on source table
    # Optional
    route:
      # A route rule that specifies the mapping between the source table and the sink table.
      - source-table: flink_test.customers
        sink-table: db.customers_o
        # Description of the route rule.
        description: sync customers table
      - source-table: flink_test.customers_suffix
        sink-table: db.customers_s
        # Description of the route rule.
        description: sync customers_suffix table
    # Optional
    pipeline:
      # The name of the job.
      name: MySQL to Hologres Pipeline
    Remarque

    Dans un job Flink CDC, séparez une clé et sa valeur par deux-points et un espace. Le format est Key: Value.

    Le tableau suivant décrit les blocs de code.

    Obligatoire

    Module

    Description

    Oui

    source

    Point de départ du pipeline de données. Flink CDC capture les données de modification depuis la source.

    Remarque
    • Actuellement, MySQL est la seule source prise en charge. Pour les détails de configuration spécifiques, consultez Connecteur MySQL.

    • Vous pouvez utiliser des variables pour gérer les informations sensibles. Pour plus d’informations, consultez Gestion des variables.

    sink

    Point de terminaison du pipeline de données. Flink CDC transmet les modifications de données capturées au système de destination.

    Remarque
    • Pour obtenir des informations sur les puits pris en charge, consultez Connecteurs d’ingestion de données Flink CDC. Pour les détails de configuration, consultez la documentation du connecteur spécifique.

    • Vous pouvez utiliser des variables pour gérer les informations sensibles. Pour plus d’informations, consultez Gestion des variables.

    Non

    pipeline

    (pipeline de données)

    Définit les configurations de base pour l’ensemble du job de pipeline de données, telles que le nom du pipeline.

    transform

    Spécifie les règles de transformation des données pour opérer sur les données lorsqu’elles circulent dans le pipeline Flink. Il prend en charge le traitement ETL, le filtrage via la clause WHERE, l’élagage de colonnes et les colonnes calculées.

    Si vous devez transformer les données brutes de modification capturées par Flink CDC pour les adapter à des systèmes en aval spécifiques, utilisez le bloc transform.

    route

    Si ce module n’est pas configuré, le job effectue par défaut une synchronisation complète de la base de données ou de la table cible.

    Dans certains cas, vous devrez peut-être envoyer les données de modification capturées vers différentes destinations selon des règles spécifiques. Le module route vous permet de spécifier de manière flexible la relation de mappage entre la source et la destination, en envoyant les données vers différentes cibles.

    Pour plus d’informations sur la syntaxe et la configuration de chaque module, consultez Référence pour les jobs d’ingestion de données Flink CDC.

    Le code suivant fournit un exemple de synchronisation de toutes les tables de la base de données app_db dans MySQL vers une base de données dans Hologres.

    source:
      type: mysql
      hostname: <hostname>
      port: 3306
      username: ${secret_values.mysqlusername}
      password: ${secret_values.mysqlpassword}
      tables: app_db.\.*
      server-id: 5400-5404
      # (Optional) Synchronize data from tables newly created in the incremental phase.
      scan.binlog.newly-added-table.enabled: true
      # (Optional) Synchronize table and field comments.
      include-comments.enabled: true
      # (Optional) Prioritize dispatching unbounded chunks to prevent potential TaskManager OutOfMemory issues.
      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: hologres
      name: Hologres Sink
      endpoint: <endpoint>
      dbname: <database-name>
      username: ${secret_values.holousername}
      password: ${secret_values.holopassword}
    pipeline:
      name: Sync MySQL Database to Hologres
  6. (Facultatif) Cliquez sur Validate.

    Cette option vérifie les erreurs de syntaxe, les problèmes de connectivité réseau et les problèmes d’autorisation.

Rubriques connexes