Tous les produits
Search
Centre de documentation

Data Lake Formation:Connect Flink CDC to DLF

Dernière mise à jour :Aug 11, 2026

Synchronisez des données vers un catalogue DLF via Iceberg REST en utilisant Flink CDC sur Realtime Compute for Apache Flink.

Prérequis

  • Vous disposez d'un espace de travail entièrement géré pour Realtime Compute for Apache Flink. Si vous n'en avez pas encore créé, consultez la rubrique Activer Realtime Compute for Apache Flink.

  • Vérifiez que votre espace de travail Realtime Compute for Apache Flink et DLF se trouvent dans la même région. Vous devez également ajouter le VPC de l'espace de travail à la liste d'autorisation de DLF. Pour plus d'informations, consultez la section Configurer une liste d'autorisation VPC.

Limites

La connectivité Iceberg REST à DLF nécessite le moteur Realtime Compute for Apache Flink VVR 11.6.0 ou une version ultérieure.

Enregistrer le catalogue DLF dans Flink

L'enregistrement d'un catalogue dans Flink crée un mappage vers votre catalogue DLF. La création ou la suppression du catalogue dans Flink n'affecte pas les données réelles stockées dans DLF. Toutes les tables créées dans le catalogue DLF via Iceberg REST sont des tables Iceberg.
  1. Connectez-vous à la console de gestion Realtime Compute for Apache Flink.

  2. Dans la colonne Actions de votre espace de travail, cliquez sur Console.

  3. Dans le volet de navigation de gauche, cliquez sur Development > Scripts.

  4. Créez un nouveau script, puis collez l'instruction SQL suivante dans l'éditeur SQL.

    CREATE CATALOG `catalog_name`
     WITH (
        'type' = 'iceberg',
        'catalog-type' = 'rest',
        'uri' = 'http://{region-id}-vpc.dlf.aliyuncs.com/iceberg',
        'warehouse' = 'iceberg_test',
        'rest.signing-region' = '{region-id}',
        'io-impl' = 'org.apache.iceberg.rest.DlfFileIO'
    );

    Remplacez {region-id} par l'ID de région de votre catalogue DLF, par exemple cn-hangzhou ou ap-southeast-1. Pour connaître toutes les régions prises en charge et leurs valeurs d'endpoint, consultez la rubrique Régions et endpoints.

  5. Dans le coin inférieur droit, cliquez sur Environment, sélectionnez un cluster de session exécutant VVR 11.2.0 ou une version ultérieure, puis exécutez l'instruction SQL.

Le tableau suivant décrit les options de configuration.

Option Description Required Example
type Type de catalogue. Définissez la valeur sur iceberg. Oui iceberg
catalog-type Type de catalogue. Définissez la valeur sur rest. Oui rest
token.provider Fournisseur de jetons pour l'authentification DLF. Définissez la valeur sur dlf. Oui dlf
uri Endpoint Iceberg REST de votre catalogue DLF. Utilisez le format http://{region-id}-vpc.dlf.aliyuncs.com/iceberg. Pour les valeurs spécifiques à chaque région, consultez la rubrique Régions et endpoints. Oui http://ap-southeast-1-vpc.dlf.aliyuncs.com/iceberg
warehouse Nom de votre catalogue DLF. Oui iceberg_test
rest.signing-region ID de région de votre catalogue DLF. Pour obtenir les ID de région, consultez la page Régions et endpoints. Oui ap-southeast-1
io-impl Implémentation FileIO pour DLF. Définissez la valeur sur org.apache.iceberg.rest.DlfFileIO. Oui org.apache.iceberg.rest.DlfFileIO

Configurer Flink CDC pour utiliser le catalogue

Créez une tâche d'ingestion de données : Développer une tâche d'ingestion de données Flink CDC.

Si vous disposez déjà d'un mappage de catalogue Flink, Réutilisez un catalogue existant pour obtenir les informations de connexion et ajoutez cette configuration Sink :

sink:
  type: iceberg
  using.built-in-catalog: catalog_name

Exemples de configuration

Modèles courants de synchronisation de données utilisant des tâches YAML Flink CDC :

Synchroniser une base de données MySQL complète vers DLF

Cette tâche synchronise l'intégralité d'une base de données MySQL vers DLF :

source:
  type: mysql
  name: MySQL Source
  hostname: ${mysql.hostname}
  port: ${mysql.port}
  username: ${mysql.username}
  password: ${mysql.password}
  tables: mysql_test.\.*
  server-id: 8601-8604
  # (Optional) Sync tables created during incremental phase.
  scan.binlog.newly-added-table.enabled: true
  # (Optional) Sync table and field comments.
  include-comments.enabled: true
  # (Optional) Prioritize unbounded chunks to prevent TaskManager OOM errors.
  scan.incremental.snapshot.unbounded-chunk-first.enabled: true
  # (Optional) Speed up reads by parsing only matched tables.
  scan.only.deserialize.captured.tables.changelog.enabled: true

sink:
  type: iceberg
  using.built-in-catalog: catalog_name
Remarque

Paramètres optionnels recommandés pour les sources MySQL (MySQL) :

  1. Paramètre : scan.binlog.newly-added-table.enabled

    Fonction : Synchronise les tables créées pendant la phase incrémentielle.

  2. Paramètre : include-comments.enabled

    Fonction : Synchronise les commentaires de table et de champ.

  3. Paramètre : scan.incremental.snapshot.unbounded-chunk-first.enabled

    Fonction : Empêche les erreurs OOM (Out Of Memory) du TaskManager.

  4. Paramètre : scan.only.deserialize.captured.tables.changelog.enabled: true

    Fonction : Accélère les lectures en analysant uniquement les tables correspondantes.

Écrire dans une table DLF partitionnée

Spécifiez les clés de partition avec le paramètre partition-keys (Référence de développement des tâches d'ingestion de données Flink CDC) :

source:
  type: mysql
  name: MySQL Source
  hostname: ${mysql.hostname}
  port: ${mysql.port}
  username: ${mysql.username}
  password: ${mysql.password}
  tables: mysql_test.\.*
  server-id: 8601-8604
  # (Optional) Sync tables created during incremental phase.
  scan.binlog.newly-added-table.enabled: true
  # (Optional) Sync table and field comments.
  include-comments.enabled: true
  # (Optional) Prioritize unbounded chunks to prevent TaskManager OOM errors.
  scan.incremental.snapshot.unbounded-chunk-first.enabled: true
  # (Optional) Speed up reads by parsing only matched tables.
  scan.only.deserialize.captured.tables.changelog.enabled: true

sink:
  type: iceberg
  using.built-in-catalog: catalog_name

transform:
  - source-table: mysql_test.tbl1
    # (Optional) Set partition keys.
    partition-keys: id,pt
  - source-table: mysql_test.tbl2
    partition-keys: id,pt

Écrire dans une table DLF en ajout seul (append-only)

Pour implémenter des suppressions logicielles (convertir les opérations DELETE en INSERT en aval), configurez la tâche comme suit (Référence de développement des tâches d'ingestion de données Flink CDC) :

source:
  type: mysql
  name: MySQL Source
  hostname: ${mysql.hostname}
  port: ${mysql.port}
  username: ${mysql.username}
  password: ${mysql.password}
  tables: mysql_test.\.*
  server-id: 8601-8604
  # (Optional) Sync tables created during incremental phase.
  scan.binlog.newly-added-table.enabled: true
  # (Optional) Sync table and field comments.
  include-comments.enabled: true
  # (Optional) Prioritize unbounded chunks to prevent TaskManager OOM errors.
  scan.incremental.snapshot.unbounded-chunk-first.enabled: true
  # (Optional) Speed up reads by parsing only matched tables.
  scan.only.deserialize.captured.tables.changelog.enabled: true

sink:
  type: iceberg
  using.built-in-catalog: catalog_name

transform:
  - source-table: mysql_test.tbl1
    # (Optional) Set partition keys.
    partition-keys: id,pt
    # (Optional) Implement soft delete.
    projection: \*, __data_event_type__ AS op_type
    converter-after-transform: SOFT_DELETE
  - source-table: mysql_test.tbl2
    # (Optional) Set partition keys.
    partition-keys: id,pt
    # (Optional) Implement soft delete.
    projection: \*, __data_event_type__ AS op_type
    converter-after-transform: SOFT_DELETE
Remarque

Synchroniser des données Kafka CDC en temps réel vers DLF

Cet exemple synchronise les données de modification de deux tables (clients et produits) depuis un topic d'inventaire Kafka au format Debezium JSON vers DLF :

source:
  type: kafka
  name: Kafka Source
  properties.bootstrap.servers: ${kafka.bootstrap.servers}
  topic: inventory
  scan.startup.mode: earliest-offset
  value.format: debezium-json
  debezium-json.distributed-tables: true

sink:
  type: iceberg
  using.built-in-catalog: catalog_name

# Debezium JSON lacks primary key info. Add it manually.
transform:
  - source-table: \.*.\.*
    projection: \*
    primary-keys: id
Remarque
  • La source Kafka prend en charge les formats canal-json, debezium-json (par défaut) et json.

  • Lorsque vous utilisez debezium-json, ajoutez manuellement une clé primaire à l'aide d'une règle de transformation, car les messages Debezium JSON ne contiennent pas d'informations sur la clé primaire :

    transform:
      - source-table: \.*.\.*
        projection: \*
        primary-keys: id
  • Si les données d'une seule table s'étendent sur plusieurs partitions, ou si des tables situées dans différentes partitions doivent être fusionnées, définissez debezium-json.distributed-tables ou canal-json.distributed-tables sur true.

  • La source Kafka prend en charge plusieurs stratégies d'inférence de schéma via le paramètre schema.inference.strategy. Kafka.

Synchroniser des journaux Kafka en temps réel vers DLF

Pour les données JSON personnalisées dans Kafka, Flink CDC gère automatiquement l'inférence des types de données, l'inférence de schéma et l'évolution de schéma.

Cet exemple synchronise une seule table de journaux JSON depuis le topic d'inventaire vers DLF :

source:
  type: kafka
  name: Kafka Source
  properties.bootstrap.servers: ${kafka.bootstrap.servers}
  topic: inventory
  scan.startup.mode: earliest-offset
  value.format: json
  # (Optional) Recursively flatten nested columns in JSON data.
  json.infer-schema.flatten-nested-columns.enable: true
  # (Optional) Skip first 100 parsing errors. Job fails if errors exceed 100.
  ingestion.ignore-errors: true
  ingestion.error-tolerance.max-count: 100

sink:
  type: iceberg
  using.built-in-catalog: catalog_name

# Add primary key to table.
transform:
  - source-table: \.*.\.*
    projection: \*
    primary-keys: id

# Write all inventory topic data to test_database.inventory.
route:
  - source-table: inventory
    sink-table: test_database.inventory

pipeline:
  # (Optional) Log dirty data that causes processing exceptions.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Cet exemple synchronise plusieurs tables de journaux JSON depuis le topic d'inventaire, en utilisant les champs databaseName et tableName pour identifier les tables :

source:
  type: kafka
  name: Kafka Source
  properties.bootstrap.servers: ${kafka.bootstrap.servers}
  topic: inventory
  scan.startup.mode: earliest-offset
  value.format: json
  # (Optional) Recursively flatten nested columns in JSON data.
  json.infer-schema.flatten-nested-columns.enable: true
  # Use databaseName and tableName field values as database and table names.
  json.decode.parser-table-id.fields: databaseName,tableName
  # (Optional) Skip first 100 parsing errors. Job fails if errors exceed 100.
  ingestion.ignore-errors: true
  ingestion.error-tolerance.max-count: 100

sink:
  type: iceberg
  using.built-in-catalog: catalog_name

# Add primary key to tables.
transform:
  - source-table: \.*.\.*
    projection: \*
    primary-keys: id

# Write ods.inventory, ods.customer, and ods.user to test_database.inventory, test_database.customer, and test_database.user respectively.
route:
  - source-table: ods.inventory
    sink-table: test_database.inventory
  - source-table: ods.customer
    sink-table: test_database.customer
  - source-table: ods.user
    sink-table: test_database.user

pipeline:
  # (Optional) Log dirty data that causes processing exceptions.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Stratégies d'inférence et d'évolution de schéma JSON : Stratégies d'analyse de schéma et de synchronisation des modifications.

Référence de configuration complète : Référence de développement des tâches d'ingestion de données Flink CDC.