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 | |
Remarque Prend en charge les connexions à ApsaraDB RDS for MySQL, PolarDB for MySQL et aux instances MySQL autogérées. | √ | × |
× | √ | |
√ Remarque Pris en charge uniquement à partir de VVR 8.0.10. | √ | |
× | √ | |
× | √ | |
× | √ | |
√ Remarque Pris en charge uniquement à partir de VVR 11.1. | × | |
√ Remarque Pris en charge uniquement à partir de VVR 11.2. | × | |
× | √ Remarque Pris en charge uniquement à partir de VVR 11.1. | |
× | √ Remarque Pris en charge uniquement à partir de VVR 11.1. | |
√ Remarque Pris en charge uniquement à partir de VVR 11.4. | × | |
× | √ | |
× | √ 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)
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 |
|
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.ordersest supprimée du flux de données. Pour conserver la table d'origine, ajoutez-la à la listeoutput-table.source-expand: - input-table: mydb.orders output-table: [ mydb.orders, db1.orders, db2.orders ] Les paramètres
input-tableetoutput-tablene 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 à |
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.