Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Flink CDC source and sink modules

Dernière mise à jour :Aug 09, 2026

Cette rubrique décrit les modules source et sink pour les travaux d'ingestion de données Flink Change Data Capture (CDC) et répertorie les connecteurs pris en charge.

Connecteurs pris en charge

Connecteur

Type pris en charge

Source

Sink

MySQL

Remarque

Prend en charge les connexions à ApsaraDB RDS for MySQL, PolarDB for MySQL et aux instances MySQL autogérées.

×

Paimon

×

Kafka

Remarque

Pris en charge uniquement à partir de VVR 8.0.10.

Upsert Kafka

×

StarRocks

×

Hologres

×

Simple Log Service

Remarque

Pris en charge uniquement à partir de VVR 11.1.

×

MongoDB

Remarque

Pris en charge uniquement à partir de VVR 11.2.

×

MaxCompute

×

Remarque

Pris en charge uniquement à partir de VVR 11.1.

SelectDB

×

Remarque

Pris en charge uniquement à partir de VVR 11.1.

PostgreSQL CDC (Public Preview)

Remarque

Pris en charge uniquement à partir de VVR 11.4.

×

Print

×

Iceberg

×

Remarque

Pris en charge uniquement à partir de VVR 11.6.

Configuration du connecteur

Configurez les paramètres des connecteurs source et sink dans un travail d'ingestion de données Flink CDC. Pour plus d'informations sur les connecteurs pris en charge et leurs paramètres, consultez les sections suivantes.

# Source module
source:
  type: mysql # Or another connector identifier
  name: MySQL Source
  # Other parameters. Use key: value pairs.

# Sink module
sink:
  type: paimon # Or another connector identifier
  name: Paimon Sink
  # Other parameters. Use key: value pairs.

Configuration générale

Paramètre

Description

Obligatoire

Type

Valeur par défaut

Remarques

type

Type de connecteur pour la source ou le sink.

Oui

String

Aucun

Aucun

name

Nom du nœud.

Non

String

Aucun

Aucun

using.built-in-catalog

Réutilise un catalogue intégré.

Non

String

Aucun

Réutilisation des informations de connexion d'un catalogue

À partir de VVR 11.5, référencez directement un catalogue intégré créé sur la page Data Management pour votre travail d'ingestion de données Flink CDC. Cette fonctionnalité permet de réutiliser les propriétés de connexion telles que l'URL, le nom d'utilisateur et le mot de passe, ce qui simplifie la configuration manuelle.

Syntaxe
source:
  type: mysql
  using.built-in-catalog: mysql_rds_catalog
  
sink:
  type: paimon
  using.built-in-catalog: paimon_dlf_catalog

Utilisez le paramètre using.built-in-catalog dans les modules source et sink pour faire référence à un catalogue intégré existant.

Par exemple, dans l'exemple ci-dessus, les métadonnées du catalogue pour mysql_rds_catalog incluent déjà les paramètres requis tels que hostname, username et password. Il n'est donc pas nécessaire de fournir à nouveau ces paramètres dans le travail YAML.

Notes d'utilisation

Les connecteurs suivants prennent en charge la réutilisation des informations de connexion d'un catalogue :

  • MySQL (source)

  • Kafka (source)

  • Upsert Kafka (sink)

  • StarRocks (sink)

  • Hologres (sink)

  • Paimon (sink)

  • Simple Log Service (source)

  • Iceberg (sink)

Remarque

Les paramètres de catalogue incompatibles avec la configuration CDC YAML sont ignorés. Pour plus d'informations, consultez la liste des paramètres de chaque connecteur.

Configuration de la source

Paramètre

Description

Obligatoire

Type

Valeur par défaut

Remarques

source-expand

Spécifie la stratégie de distribution des données émises par la source.

Non

Consultez la section de syntaxe pour les détails de configuration.

Aucun

  • Disponible à partir de VVR 11.6.

source-expand

Le paramètre source-expand distribue ou réplique les données avant leur traitement par les modules Transform et Route.

Syntaxe
source:                                                                                                                                                                                                                                                                                                          
  type: mysql                                                                                                                                                                                                                                                                                                     
  host: localhost                                                                                                                                                                                                                                                                                                  
  port: 3306                                                                                                                                                                                                                                                                                                       
  username: admin
  password: pass
  tables: mydb.orders
  source-expand:
  # Replicate mydb.orders into three tables, each processed differently before writing to the sink.
    - input-table: mydb.orders
      output-table: [ dwd.orders_full, dws.orders_summary, ads.orders_report ]

# Apply different transformations to each of the three expanded tables
transform:
  # DWD layer: Retain all fields, add a calculated column
  - source-table: dwd.orders_full
    projection: "*, amount * discount as final_price"
  # DWS layer: Retain only key fields, filter out small orders
  - source-table: dws.orders_summary
    projection: order_id, user_id, amount, order_status
    filter: amount > 100
    primary-keys: order_id
  # ADS layer: Retain only fields needed for reporting, filter out canceled orders
  - source-table: ads.orders_report
    projection: order_id, user_id, amount, TO_UPPER(order_status) as status
    filter: order_status <> 'CANCELLED'
    primary-keys: order_id

# Route the three expanded tables to different sink tables
route:
  - source-table: dwd.orders_full
    sink-table: starrocks_dwd.orders_full_detail
  - source-table: dws.orders_summary
    sink-table: starrocks_dws.orders_summary
  - source-table: ads.orders_report
    sink-table: starrocks_ads.orders_report

sink:
  type: starrocks
  name: sink-starrocks
  jdbc-url: jdbc:mysql://localhost:9030
  load-url: localhost:8030
  username: root
  password: pass

Comme illustré dans l'exemple, cette configuration réplique les données de la table mydb.orders vers trois tables logiques : dwd.orders_full, dws.orders_summary et ads.orders_report. Chacune de ces tables est ensuite traitée différemment avant d'être écrite dans une table sink distincte.

Notes d'utilisation
  • Disponible à partir de VVR 11.6.

  • Par défaut, la table d'origine n'est pas conservée après la distribution des données. Dans l'exemple de syntaxe, la table mydb.orders est supprimée du flux de données. Pour conserver la table d'origine, ajoutez-la à la liste output-table.

    source-expand:
      - input-table: mydb.orders
        output-table: [ mydb.orders, db1.orders, db2.orders ]
  • Les paramètres input-table et output-table ne prennent pas en charge les expressions régulières.

Configuration du sink

Paramètre

Description

Obligatoire

Type

Valeur par défaut

Remarques

include.schema.changes

Spécifie les types de modifications de schéma à appliquer.

Non

List<String>

Aucun

Par défaut, toutes les modifications de schéma sont synchronisées.

exclude.schema.changes

Spécifie les types de modifications de schéma à exclure.

Non

List<String>

Aucun

Priorité supérieure à include.schema.changes.

Pour plus d'informations sur l'utilisation de include.schema.changes et exclude.schema.changes, consultez la section Configuration de la synchronisation des modifications de schéma.