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 objetRowDatapossè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
SchemaChangeEventpour 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 utiliseDataChangeEventpour 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 |
|
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 TABLEetDROP 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
Connectez-vous à la console Realtime Compute for Apache Flink.
Dans la colonne Actions de l’espace de travail cible, cliquez sur Console.
Dans le volet de navigation de gauche, choisissez .
Cliquez sur
, puis sur New Draft with Template.-
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.
Saisissez les informations du job, telles que OK, location et engine version, puis cliquez sur OK.
-
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
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
transformgé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-idde la source dans le job d’ingestion de données Flink CDC pour éviter les conflits.
Connectez-vous à la console Realtime Compute for Apache Flink.
Dans la colonne Actions de l’espace de travail cible, cliquez sur Console.
Dans le volet de navigation de gauche, choisissez .
-
Cliquez sur
, 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.
Saisissez les informations du job, telles que OK, location et engine version, puis cliquez sur OK.
À partir de Flink CDC open source
Connectez-vous à la console Realtime Compute for Apache Flink.
Dans la colonne Actions de l’espace de travail cible, cliquez sur Console.
Dans le volet de navigation de gauche, choisissez .
Cliquez sur
, puis sélectionnez New Draft. Saisissez le name et la engine version, puis cliquez sur Create.Collez le code YAML de votre job Flink CDC open source dans l’éditeur.
-
(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
Connectez-vous à la console Realtime Compute for Apache Flink.
Dans la colonne Actions de l’espace de travail cible, cliquez sur Console.
Dans le volet de navigation de gauche, choisissez .
Cliquez sur
, puis sélectionnez New Draft. Saisissez le name et la engine version, puis cliquez sur Create.-
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 PipelineRemarqueDans 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
routevous 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 -
-
(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
Une fois le développement d’un job Flink CDC terminé, vous devez le déployer. Pour plus d’informations, consultez Déployer un job.
Pour créer rapidement un job Flink CDC afin de synchroniser les données d’une base de données MySQL vers StarRocks, consultez Tutoriel : Créer un job d’ingestion de données Flink CDC.