Cette rubrique explique comment utiliser le connecteur MySQL dans les jobs SQL.
Informations générales
Le connecteur MySQL prend en charge toutes les bases de données compatibles avec le protocole MySQL, telles qu'ApsaraDB RDS for MySQL, PolarDB for MySQL, OceanBase (mode MySQL) et les instances MySQL auto-gérées.
Lorsque vous utilisez le connecteur MySQL pour lire des données depuis OceanBase, assurez-vous que la journalisation binaire (binlog) est activée et correctement configurée. Pour plus d'informations, consultez la section Opérations liées aux binlogs. Cette fonctionnalité est en aperçu public. Utilisez-la avec prudence.
Le connecteur MySQL offre les fonctionnalités suivantes.
Catégorie | Détails |
Types pris en charge | Tables source, tables de dimension, tables de destination et sources de données pour l'ingestion de données |
Mode d'exécution | Seul le mode streaming est pris en charge. |
Format des données | Sans objet |
Métriques de surveillance spécifiques | |
Types d'API | DataStream, SQL et YAML d'ingestion de données |
Prise en charge de la mise à jour ou de la suppression des données dans les tables de destination | Oui |
Fonctionnalités
Une table source MySQL CDC (Change Data Capture), également appelée table source de streaming MySQL, commence par lire l'intégralité des données historiques de la base de données. Elle bascule ensuite de manière transparente vers la lecture des journaux binaires. Ce processus garantit qu'aucune donnée n'est manquante ou dupliquée. Même en cas de défaillance, les données sont traitées avec une sémantique exactly-once. Une table source MySQL CDC prend en charge la lecture simultanée des données complètes. Elle utilise un algorithme d'instantané incrémentiel pour implémenter une lecture sans verrou et un transfert de données reprenable. Pour plus d'informations, consultez la section À propos des tables sources MySQL CDC.
Traitement unifié par lots et en flux qui permet de lire à la fois les données complètes et incrémentielles, éliminant ainsi la nécessité de maintenir deux processus distincts.
Lecture simultanée des données complètes pour une mise à l'échelle horizontale des performances.
Basculement transparent de la lecture des données complètes vers la lecture des données incrémentielles et mise à l'échelle automatique vers le bas pour économiser les ressources de calcul.
Transfert de données reprenable pendant la phase de lecture des données complètes pour une stabilité accrue.
Lecture sans verrou des données complètes, qui n'affecte pas les services en ligne.
Prise en charge de la lecture des journaux de sauvegarde d'ApsaraDB RDS for MySQL.
Analyse parallèle des fichiers de journaux binaires pour réduire la latence de lecture.
Prérequis
Avant d'utiliser une table source MySQL CDC, vous devez effectuer les opérations préalables décrites dans la section Configurer MySQL.
ApsaraDB RDS for MySQL
Effectuez une sonde réseau pour garantir la connectivité réseau avec Realtime Compute for Apache Flink.
Version MySQL : 5.6, 5,7, 8,0.x ou 8,4.
La journalisation binaire doit être activée. Elle est activée par défaut.
Le format du journal binaire doit être ROW. Il s'agit du format par défaut.
Le paramètre
binlog_row_imagedoit être défini sur FULL. Il s'agit du paramètre par défaut.La compression des transactions de journal binaire doit être désactivée. Cette fonctionnalité a été introduite dans MySQL 8.0.20 et est désactivée par défaut.
Un utilisateur MySQL a été créé avec les autorisations SELECT, SHOW DATABASES, REPLICATION SLAVE et REPLICATION CLIENT.
Créez une base de données et une table MySQL. Pour plus d'informations, consultez la section Créer une base de données et un compte pour une instance ApsaraDB RDS for MySQL. Utilisez un compte privilégié pour créer la base de données MySQL afin d'éviter les échecs d'opération dus à des autorisations insuffisantes.
Configurez une liste blanche d'adresses IP. Pour plus d'informations, consultez la section Configurer une liste blanche d'adresses IP pour une instance ApsaraDB RDS for MySQL.
PolarDB for MySQL
Effectuez une sonde réseau pour garantir la connectivité réseau avec Realtime Compute for Apache Flink.
Version MySQL : 5.6, 5,7, 8,0.x ou 8,4.
La journalisation binaire doit être activée. Elle est désactivée par défaut.
Le format du journal binaire doit être ROW. Il s'agit du format par défaut.
Le paramètre
binlog_row_imagedoit être défini sur FULL. Il s'agit du paramètre par défaut.La compression des transactions de journal binaire doit être désactivée. Cette fonctionnalité a été introduite dans MySQL 8.0.20 et est désactivée par défaut.
Vous avez créé un utilisateur MySQL disposant des autorisations SELECT, SHOW DATABASES, REPLICATION SLAVE et REPLICATION CLIENT.
Créez une base de données et une table MySQL. Pour plus d'informations, consultez la section Créer une base de données et un compte pour un cluster PolarDB for MySQL. Utilisez un compte privilégié pour créer la base de données MySQL afin d'éviter les échecs d'opération dus à des autorisations insuffisantes.
Configurez une liste blanche d'adresses IP. Pour plus d'informations, consultez la section Configurer une liste blanche d'adresses IP pour un cluster PolarDB for MySQL.
Self-managed MySQL
Effectuez une sonde réseau pour garantir la connectivité réseau avec Realtime Compute for Apache Flink.
Version MySQL : 5.6, 5,7, 8,0.x ou 8,4.
La journalisation binaire doit être activée. Elle est désactivée par défaut.
Le format du journal binaire doit être ROW. Le format par défaut est STATEMENT.
Le paramètre
binlog_row_imagedoit être défini sur FULL. Il s'agit du paramètre par défaut.La compression des transactions de journal binaire doit être désactivée. Cette fonctionnalité a été introduite dans MySQL 8.0.20 et est désactivée par défaut.
Créez un utilisateur MySQL et accordez-lui les autorisations SELECT, SHOW DATABASES, REPLICATION SLAVE et REPLICATION CLIENT.
Créez une base de données et une table MySQL. Pour plus d'informations, consultez la section Créer une base de données et un compte pour une instance MySQL auto-gérée. Utilisez un compte privilégié pour créer la base de données MySQL afin d'éviter les échecs d'opération dus à des autorisations insuffisantes.
Configurez une liste blanche d'adresses IP. Pour plus d'informations, consultez la section Configurer une liste blanche d'adresses IP pour une instance MySQL auto-gérée.
Limites
Limites générales
Les tables sources MySQL CDC ne prennent pas en charge les définitions de watermark.
Dans les jobs Create Table As Select (CTAS) et Create Database As Select (CDAS), les tables sources MySQL CDC peuvent synchroniser certaines modifications de schéma. Pour plus d'informations sur les types de modification pris en charge, consultez la section Politiques de synchronisation de l'évolution du schéma.
Le connecteur MySQL CDC ne prend pas en charge la fonctionnalité de compression des transactions de journal binaire. Par conséquent, lorsque vous utilisez le connecteur MySQL CDC pour consommer des données incrémentielles, assurez-vous que la compression des transactions de journal binaire est désactivée. Sinon, le connecteur risque de ne pas pouvoir récupérer les données incrémentielles.
Limites d'**ApsaraDB RDS for MySQL**
Pour ApsaraDB RDS for MySQL, ne lisez pas les données depuis une base de données secondaire ou un réplica en lecture seule. En effet, la période de rétention par défaut des journaux binaires pour les bases de données secondaires et les réplicas en lecture seule est courte. Si les journaux binaires expirent et sont supprimés, le job peut échouer à consommer les données du journal binaire et signaler une erreur.
ApsaraDB RDS for MySQL active par défaut la synchronisation parallèle primaire/secondaire, mais ne garantit pas un ordre de transaction cohérent entre les instances primaires et secondaires. Cela peut entraîner une perte de données lors d'un basculement primaire/secondaire et d'une récupération de point de contrôle. Pour éviter ce problème, vous pouvez activer manuellement l'option
slave_preserve_commit_orderpour ApsaraDB RDS for MySQL.
Limites d'**PolarDB for MySQL**
Les tables sources MySQL CDC ne prennent pas en charge la lecture des données depuis les clusters d'architecture Multi-master Cluster de PolarDB for MySQL V1.0.19 et versions antérieures. Pour plus d'informations, consultez la section Qu'est-ce qu'un Multi-master Cluster ?. Les journaux binaires générés par ces clusters peuvent contenir des ID de table dupliqués. Cela peut provoquer des erreurs de mappage de schéma dans la table source CDC, entraînant des erreurs lors de l'analyse des données du journal binaire.
Limites d'**Open source MySQL**
Par défaut, MySQL maintient l'ordre des transactions lors de la réplication des journaux binaires primaire/secondaire. Si un réplica MySQL a la réplication parallèle activée (slave_parallel_workers > 1) mais n'a pas activé slave_preserve_commit_order=ON, l'ordre de validation des transactions peut être incohérent avec celui de la base de données primaire. Lorsque Flink CDC récupère à partir d'un point de contrôle, il peut manquer des données en raison du séquençage désordonné. Vous pouvez définir slave_preserve_commit_order = ON sur le réplica MySQL. Vous pouvez également définir slave_parallel_workers = 1, mais cela sacrifiera les performances de réplication.
Remarques d'utilisation
-
Table source
Pendant la phase de lecture complète des données, vous ne pouvez pas enregistrer de point de contrôle (savepoint), ni ajouter ou supprimer une table dans la table source, puis redémarrer le job à partir du point de contrôle. Si vous effectuez ces opérations, le job échouera à lire les données.
-
Table de destination
Clés primaires auto-incrémentées : ne déclarez pas de clés primaires auto-incrémentées dans le DDL. MySQL les renseigne automatiquement lors de l'écriture des données.
Vous devez déclarer au moins un champ qui n'est pas une clé primaire. Sinon, une erreur est signalée.
La contrainte
NOT ENFORCEDdans le DDL indique que Flink n'applique pas la validation de la clé primaire. Il vous incombe de garantir l'exactitude et l'intégrité de la clé primaire. Pour plus d'informations, consultez Vérification de validité.
-
Table de dimension
Si vous souhaitez utiliser un index pour accélérer les requêtes, l'ordre des champs dans la clause JOIN doit correspondre à l'ordre défini dans l'index. Cela repose sur la règle du préfixe gauche. Par exemple, si l'index est (a, b, c), la condition JOIN est
ON t.a = x AND t.b = y.Le SQL généré par Flink peut être réécrit par l'optimiseur. Cela peut empêcher l'utilisation de l'index lors de la requête réelle vers la base de données. Pour confirmer si l'index est utilisé, vérifiez le plan d'exécution (EXPLAIN) ou le journal des requêtes lentes dans MySQL afin de consulter l'instruction SELECT réellement exécutée.
SQL
Vous pouvez utiliser le connecteur MySQL dans les jobs SQL en tant que table source, table de dimension ou table de destination.
Syntaxe
CREATE TEMPORARY TABLE mysqlcdc_source (
order_id INT,
order_date TIMESTAMP(0),
customer_name STRING,
price DECIMAL(10, 5),
product_id INT,
order_status BOOLEAN,
PRIMARY KEY(order_id) NOT ENFORCED
) WITH (
'connector' = 'mysql',
'hostname' = '<yourHostname>',
'port' = '3306',
'username' = '<yourUsername>',
'password' = '<yourPassword>',
'database-name' = '<yourDatabaseName>',
'table-name' = '<yourTableName>'
);
-
Lors de l'écriture dans une table de destination, le connecteur construit et exécute une instruction SQL pour chaque enregistrement de données reçu. L'instruction est structurée comme suit :
Pour une table de destination sans clé primaire, une instruction
INSERT INTO table_name (column1, column2, ...) VALUES (value1, value2, ...);est exécutée.Pour une table de destination avec une clé primaire, une instruction
INSERT INTO table_name (column1, column2, ...) VALUES (value1, value2, ...) ON DUPLICATE KEY UPDATE column1 = VALUES(column1), column2 = VALUES(column2), ...;est exécutée. Remarque : si la table physique possède une contrainte d'index unique autre que la clé primaire, l'insertion de deux enregistrements ayant des clés primaires différentes mais la même valeur d'index unique provoque un conflit d'index unique. Cela entraîne l'écrasement et la perte des données.
Si une clé primaire auto-incrémentée est définie dans la base de données MySQL, ne déclarez pas le champ auto-incrémenté dans le DDL Flink. La base de données renseigne automatiquement ce champ lors de l'écriture des données. Le connecteur prend en charge l'écriture et la suppression de données avec des champs auto-incrémentés, mais ne prend pas en charge la mise à jour de ces données.
Paramètres WITH
-
Général
Paramètre
Description
Obligatoire
Type de données
Valeur par défaut
Notes
connector
Le type de table.
Oui
STRING
Aucune
Lorsqu'il est utilisé comme table source, vous pouvez définir ce paramètre sur
mysql-cdcoumysql. Ces deux valeurs sont équivalentes. Lorsqu'il est utilisé comme table de dimension ou table de destination, la valeur doit êtremysql.hostname
L'adresse IP ou le nom d'hôte de la base de données MySQL.
Oui
STRING
Aucune
Nous vous recommandons de spécifier une adresse VPC (Virtual Private Cloud).
RemarqueSi la base de données MySQL et Realtime Compute for Apache Flink ne se trouvent pas dans le même VPC, vous devez établir une connexion réseau inter-VPC ou utiliser un endpoint public pour accéder à la base de données. Pour plus d'informations, consultez les rubriques Gérer et exploiter les espaces de travail et Comment un cluster Flink entièrement géré peut-il accéder à Internet ?.
username
Le nom d'utilisateur du service de base de données MySQL.
Oui
STRING
Aucune
Aucune.
password
Le mot de passe du service de base de données MySQL.
Oui
STRING
Aucune
Aucune.
database-name
Le nom de la base de données MySQL.
Oui
STRING
Aucune
Lorsqu'une base de données est utilisée comme table source, vous pouvez utiliser une expression régulière pour le nom de la base de données afin de lire les données provenant de plusieurs bases.
Lorsque vous utilisez des expressions régulières, n'utilisez pas les symboles ^ et $ pour faire correspondre le début et la fin de la chaîne. Pour plus d'informations, consultez les notes relatives au paramètre table-name.
table-name
Le nom de la table MySQL.
Oui
STRING
Aucune
Vous pouvez utiliser une expression régulière pour le nom de la table source afin de lire les données provenant de plusieurs tables.
Lorsque vous lisez des données depuis plusieurs tables MySQL, soumettez plusieurs instructions CTAS dans le cadre d'une seule tâche. Cela évite d'activer plusieurs écouteurs de journaux binaires et améliore les performances et l'efficacité. Pour plus d'informations, consultez la rubrique Plusieurs instructions CTAS : Soumission en tant que tâche unique.
Lorsque vous utilisez des expressions régulières, n'utilisez pas les symboles ^ et $ pour faire correspondre le début et la fin de la chaîne. Pour plus d'informations, consultez la note suivante.
RemarqueLorsqu'une table source MySQL CDC fait correspondre des noms de table à l'aide d'une expression régulière, elle concatène les paramètres database-name et table-name que vous avez spécifiés avec la chaîne \\. pour former une expression régulière de chemin complet. Avant la version VVR 8.0.1, le caractère . était utilisé. Le connecteur utilise ensuite cette expression régulière pour faire correspondre les noms qualifiés complets des tables dans la base de données MySQL.
Par exemple, si vous définissez 'database-name'='db_.' et 'table-name'='tb_.+', le connecteur utilise l'expression régulière db_.\\.tb_.+ pour faire correspondre les noms de table qualifiés complets et déterminer quelles tables lire. Avant la version VVR 8.0.1, l'expression régulière était db_.*.tb_.+.
port
Le numéro de port du service de base de données MySQL.
Non
INTEGER
3306
Aucune.
-
Table source uniquement
Paramètre
Description
Obligatoire
Type de données
Valeur par défaut
Notes
server-id
Identifiant numérique du client de base de données.
Non
STRING
Une valeur aléatoire comprise entre 5400 et 6400 est générée.
Cet identifiant doit être globalement unique au sein du cluster MySQL. Définissez un identifiant différent pour chaque tâche qui se connecte à la même base de données.
Ce paramètre prend également en charge le format de plage d'identifiants, par exemple 5400-5408. Lorsque la lecture incrémentielle est activée, la lecture simultanée est prise en charge. Dans ce cas, définissez une plage d'identifiants afin que chaque lecteur simultané utilise un identifiant différent. Pour plus d'informations, consultez Utiliser l'ID de serveur.
scan.incremental.snapshot.enabled
Indique s'il faut activer les snapshots incrémentiels.
Non
BOOLEAN
true
Les snapshots incrémentiels sont activés par défaut. Le snapshot incrémentiel est un nouveau mécanisme de lecture des snapshots complets de données. Par rapport à l'ancienne méthode de lecture des snapshots, les snapshots incrémentiels présentent de nombreux avantages, notamment :
La source peut lire les données complètes en parallèle.
La source prend en charge les points de contrôle au niveau des fragments lors de la lecture des données complètes.
La source n'a pas besoin d'acquérir un verrou de lecture global (FLUSH TABLES WITH read lock) lors de la lecture des données complètes.
Si vous souhaitez que la source prenne en charge la lecture simultanée, chaque lecteur simultané doit disposer d'un ID de serveur unique. Par conséquent, server-id doit être une plage, telle que 5400-6400, et la taille de la plage doit être supérieure ou égale au niveau de simultanéité.
RemarqueCet élément de configuration est supprimé dans Ververica Runtime (VVR) 11.1 et versions ultérieures.
scan.incremental.snapshot.chunk.size
Taille de chaque fragment, exprimée en nombre de lignes.
Non
INTEGER
8096
Lorsque la lecture par snapshot incrémentiel est activée, la table est divisée en plusieurs fragments pour la lecture. Les données d'un fragment sont mises en cache en mémoire avant d'être entièrement lues.
Moins un fragment contient de lignes, plus le nombre total de fragments dans la table est élevé. Bien que cela réduise la granularité de la récupération après incident, cela peut entraîner des erreurs de dépassement de mémoire (OOM) et une diminution du débit global. Vous devez donc trouver un compromis et définir une taille de fragment raisonnable.
scan.snapshot.fetch.size
Nombre maximal d'enregistrements à extraire à la fois lors de la lecture des données complètes d'une table.
Non
INTEGER
1024
Aucun.
scan.startup.mode
Mode de démarrage pour la consommation de données.
Non
STRING
initial
Valeurs valides :
initial (par défaut) : Lors du premier démarrage ou d'un démarrage sans état, le connecteur analyse toutes les données historiques, puis lit les dernières données du journal binaire.
latest-offset : Lors du premier démarrage ou d'un démarrage sans état, le connecteur n'analyse pas les données historiques. Il commence la lecture à partir de la fin du journal binaire, ce qui signifie qu'il ne lit que les dernières modifications effectuées après le démarrage du connecteur.
earliest-offset : Le connecteur n'analyse pas les données historiques. Il commence la lecture à partir du journal binaire disponible le plus ancien.
specific-offset : Le connecteur n'analyse pas les données historiques. Il démarre à partir d'un décalage spécifique du journal binaire. Vous pouvez spécifier le décalage en configurant à la fois scan.startup.specific-offset.file et scan.startup.specific-offset.pos, ou en configurant uniquement scan.startup.specific-offset.gtid-set pour démarrer à partir d'un ensemble GTID spécifique.
timestamp : Le connecteur n'analyse pas les données historiques. Il commence à lire le journal binaire à partir d'un horodatage spécifié. L'horodatage est indiqué par scan.startup.timestamp-millis en millisecondes.
ImportantLors de l'utilisation du mode de démarrage earliest-offset, specific-offset ou timestamp, assurez-vous que le schéma de la table correspondante ne change pas entre la position de consommation du journal binaire spécifiée et l'heure de démarrage de la tâche. Cela permet d'éviter les erreurs dues à des incompatibilités de schéma.
scan.startup.specific-offset.file
Nom du fichier du journal binaire pour le décalage de départ lors de l'utilisation du mode de démarrage specific-offset.
Non
STRING
Aucun
Lorsque vous utilisez ce paramètre, vous devez définir scan.startup.mode sur specific-offset. Exemple de format de nom de fichier :
mysql-bin.000003.scan.startup.specific-offset.pos
Décalage au sein du fichier de journal binaire spécifié pour le décalage de départ lors de l'utilisation du mode de démarrage specific-offset.
Non
INTEGER
Aucun
Lorsque vous utilisez ce paramètre, vous devez définir scan.startup.mode sur specific-offset.
scan.startup.specific-offset.gtid-set
Ensemble GTID pour le décalage de départ lors de l'utilisation du mode de démarrage specific-offset.
Non
STRING
Aucun
Lorsque vous utilisez ce paramètre, vous devez définir scan.startup.mode sur specific-offset. Exemple de format d'ensemble GTID :
24DA167-0C0C-11E8-8442-00059A3C7B00:1-19.scan.startup.timestamp-millis
Horodatage en millisecondes pour le décalage de départ lors de l'utilisation du mode de démarrage timestamp.
Non
LONG
Aucun
Lorsque vous utilisez ce paramètre, vous devez définir scan.startup.mode sur timestamp. L'unité de l'horodatage est la milliseconde.
ImportantLorsque vous spécifiez une heure, MySQL CDC tente de lire l'événement initial de chaque fichier de journal binaire pour déterminer son horodatage. Il localise ensuite le fichier de journal binaire correspondant à l'heure spécifiée. Assurez-vous que le fichier de journal binaire correspondant à l'horodatage spécifié n'a pas été effacé de la base de données et qu'il peut être lu.
server-time-zone
Fuseau horaire de session utilisé par la base de données.
Non
STRING
Si vous ne spécifiez pas ce paramètre, le système utilise le fuseau horaire de l'environnement d'exécution de la tâche Flink comme fuseau horaire du serveur de base de données. Il s'agit du fuseau horaire de la zone que vous avez sélectionnée.
Exemple : Asia/Shanghai. Ce paramètre contrôle la conversion du type TIMESTAMP de MySQL en type STRING. Pour plus d'informations, consultez Valeurs temporelles Debezium.
debezium.min.row.count.to.stream.results
Lorsque le nombre de lignes d'une table est supérieur à cette valeur, le mode de lecture par lots est utilisé.
Non
INTEGER
1000
Flink lit les données d'une table source MySQL de l'une des manières suivantes :
Lecture complète : lit directement toutes les données de la table en mémoire. Cette méthode est rapide, mais consomme une quantité de mémoire correspondante. Si la table source est très volumineuse, il existe un risque d'erreurs OOM.
Lecture par lots : lit les données en plusieurs lots, avec un certain nombre de lignes par lot, jusqu'à ce que toutes les données soient lues. Cette méthode évite les risques d'erreur OOM lors de la lecture de tables volumineuses, mais est relativement lente.
connect.timeout
Délai maximal d'attente avant expiration de la connexion au serveur de base de données MySQL avant nouvelle tentative.
Non
DURATION
30s
Aucun.
connect.max-retries
Nombre maximal de tentatives après un échec de connexion au service de base de données MySQL.
Non
INTEGER
3
Aucun.
connection.pool.size
Taille du pool de connexions à la base de données.
Non
INTEGER
20
Le pool de connexions permet de réutiliser les connexions, ce qui réduit le nombre de connexions à la base de données.
jdbc.properties.*
Paramètres de connexion personnalisés dans l'URL JDBC.
Non
STRING
Aucun
Vous pouvez transmettre des paramètres de connexion personnalisés. Par exemple, pour ne pas utiliser le protocole SSL, configurez 'jdbc.properties.useSSL' = 'false'.
Pour plus d'informations sur les paramètres de connexion pris en charge, consultez Propriétés de configuration MySQL.
debezium.*
Paramètres personnalisés permettant à Debezium de lire les journaux binaires.
Non
STRING
Aucun
Vous pouvez transmettre des paramètres Debezium personnalisés. Par exemple, utilisez 'debezium.event.deserialization.failure.handling.mode'='ignore' pour spécifier la logique de gestion des erreurs d'analyse.
AvertissementNe modifiez pas les paramètres Debezium de manière arbitraire. Cela pourrait entraîner une lecture incorrecte des données par le connecteur. Par exemple, il n'est pas autorisé de configurer le paramètre debezium.binlog.buffer.size.
heartbeat.interval
Intervalle auquel la source fait avancer le décalage du journal binaire à l'aide d'événements de pulsation (heartbeat).
Non
DURATION
30 s
Les événements de pulsation permettent de faire avancer le décalage du journal binaire dans la source. Cette fonctionnalité est particulièrement utile pour les tables MySQL mises à jour rarement. Pour ces tables, le décalage du journal binaire ne peut pas avancer automatiquement. Les événements de pulsation font progresser ce décalage, évitant ainsi les problèmes liés à son expiration. Un décalage expiré peut provoquer l'échec irréversible du job, nécessitant un redémarrage sans état.
scan.incremental.snapshot.chunk.key-column
Spécifie une colonne à utiliser comme colonne de fractionnement pour le sharding pendant la phase de snapshot.
Consultez la colonne Notes.
STRING
Aucun
Requis pour les tables sans clé primaire. La colonne sélectionnée doit être de type non nul (NOT NULL).
Facultatif pour les tables avec une clé primaire. Une seule colonne de la clé primaire peut être sélectionnée.
rds.region-id
ID de région de l'instance ApsaraDB RDS for MySQL d'Alibaba Cloud.
Requis lors de l'utilisation de la fonctionnalité de lecture des journaux archivés depuis OSS.
STRING
Aucun
Pour plus d'informations sur les ID de région, consultez Régions et zones.
ImportantÉtant donné que la chaîne GTID pour MySQL CDC est générée de manière aléatoire et n'augmente pas de façon monotone comme les décalages de fichiers de journaux binaires, localiser un GTID dans un fichier nécessite de télécharger et d'analyser tous les journaux archivés depuis OSS. Ce processus consomme beaucoup de ressources et prend du temps, rendant impossibles les fonctionnalités reposant sur les décalages GTID. Par conséquent, la fonctionnalité de journaux archivés OSS prend uniquement en charge le démarrage à partir d'un horodatage spécifié ou d'un décalage de fichier de journal binaire spécifié. Elle ne prend pas en charge le démarrage à partir d'un GTID spécifié, ni les scénarios impliquant des basculements primaire/secondaire dans les journaux archivés, car les basculements MySQL primaire/secondaire reposent sur les GTID. Évaluez soigneusement cette fonctionnalité avant utilisation.
rds.access-key-id
AccessKey ID du compte ApsaraDB RDS for MySQL d'Alibaba Cloud.
Requis lors de l'utilisation de la fonctionnalité de lecture des journaux archivés depuis OSS.
STRING
Aucun
Pour plus d'informations, consultez Comment afficher l'AccessKey ID et l'AccessKey secret ?.
ImportantPour éviter la divulgation de vos informations AccessKey, utilisez la fonctionnalité de gestion des secrets pour spécifier l'AccessKey ID. Pour plus d'informations, consultez Gérer les variables.
rds.access-key-secret
AccessKey secret du compte ApsaraDB RDS for MySQL d'Alibaba Cloud.
Requis lors de l'utilisation de la fonctionnalité de lecture des journaux archivés depuis OSS.
STRING
Aucun
Pour plus d'informations, consultez Comment afficher l'AccessKey ID et l'AccessKey secret ?
ImportantPour éviter la divulgation de vos informations AccessKey, utilisez la fonctionnalité de gestion des secrets pour spécifier l'AccessKey secret. Pour plus d'informations, consultez Gérer les variables.
rds.db-instance-id
ID de l'instance ApsaraDB RDS for MySQL d'Alibaba Cloud.
Requis lors de l'utilisation de la fonctionnalité de lecture des journaux archivés depuis OSS.
STRING
Aucun
Aucun.
rds.main-db-id
Numéro de la base de données principale de l'instance ApsaraDB RDS for MySQL d'Alibaba Cloud.
Non
STRING
Aucun
Pour savoir comment obtenir le numéro de la base de données principale, consultez Sauvegarde des journaux ApsaraDB RDS for MySQL.
Pris en charge uniquement dans VVR 8.0.7 et versions ultérieures.
RemarqueSi ce paramètre n'est pas spécifié, VVR 11.7 et versions ultérieures interrogent automatiquement le numéro de la base de données principale en fonction des informations de connexion ApsaraDB RDS for MySQL.
rds.download.timeout
Délai d'expiration pour le téléchargement d'un seul journal archivé depuis OSS.
Non
DURATION
60 s
Aucun.
rds.endpoint
Endpoint de service permettant d'obtenir les informations sur les journaux binaires OSS.
Non
STRING
Aucun
Pour plus d'informations sur les valeurs valides, consultez Endpoints.
Pris en charge uniquement dans VVR 8.0.8 et versions ultérieures.
scan.incremental.close-idle-reader.enabled
Indique s'il faut fermer les lecteurs inactifs une fois le snapshot terminé.
Non
BOOLEAN
false
Pris en charge uniquement dans VVR 8.0.1 et versions ultérieures.
Pour que cette configuration prenne effet, vous devez définir execution.checkpointing.checkpoints-after-tasks-finish.enabled sur true.
scan.read-changelog-as-append-only.enabled
Indique s'il faut convertir le flux de données changelog en un flux de données append-only.
Non
BOOLEAN
false
Valeurs valides :
true : Tous les types de messages, y compris INSERT, DELETE, UPDATE_BEFORE et UPDATE_AFTER, sont convertis en messages INSERT. Activez cette option uniquement dans des scénarios spécifiques, par exemple lorsque vous devez conserver les messages de suppression de la table amont.
false (par défaut) : Tous les types de messages sont transmis en aval tels quels.
RemarquePris en charge uniquement dans VVR 8.0.8 et versions ultérieures.
scan.only.deserialize.captured.tables.changelog.enabled
En phase incrémentielle, indique s'il faut désérialiser uniquement les événements de modification des tables spécifiées.
Non
BOOLEAN
La valeur par défaut est false dans les versions VVR 8.x.
La valeur par défaut est true à partir de la version VVR 11.1.
Valeurs valides :
true : Désérialise uniquement les données de modification des tables cibles pour accélérer la lecture du journal binaire.
false (par défaut) : Désérialise les données de modification de toutes les tables.
RemarquePris en charge uniquement à partir de la version VVR 8.0.7.
Si vous utilisez une version VVR 8.0.8 ou antérieure, vous devez modifier le nom du paramètre en debezium.scan.only.deserialize.captured.tables.changelog.enable.
scan.parse.online.schema.changes.enabled
En phase incrémentielle, indique s'il faut tenter d'analyser les événements DDL de modification sans verrou RDS.
Non
BOOLEAN
false
Valeurs valides :
true : Analyse les événements DDL de modification sans verrou RDS.
false (par défaut) : N'analyse pas les événements DDL de modification sans verrou RDS.
Il s'agit d'une fonctionnalité expérimentale. Avant d'effectuer une modification en ligne sans verrou, prenez un snapshot du job Flink pour permettre la récupération.
RemarquePris en charge uniquement à partir de la version VVR 11.1.
scan.incremental.snapshot.backfill.skip
Indique s'il faut ignorer le remplissage rétroactif (backfill) pendant la phase de lecture du snapshot.
Non
BOOLEAN
false
Valeurs valides :
true : Ignore le remplissage rétroactif pendant la phase de lecture du snapshot.
false (par défaut) : N'ignore pas le remplissage rétroactif pendant la phase de lecture du snapshot.
Le remplissage rétroactif s'applique uniquement lors de l'interrogation du snapshot d'un seul fragment (chunk) et ne couvre pas l'intégralité de la phase de lecture complète. Lorsque le remplissage rétroactif est ignoré, l'interrogation du snapshot de chaque fragment lit les dernières données de la table à cet instant précis ; les mises à jour intervenant sur un fragment après sa lecture ne sont pas fusionnées pendant la phase de lecture complète et sont lues depuis le Binlog une fois la phase incrémentielle engagée. Par exemple, une mise à jour du fragment5 qui survient pendant la prise du snapshot du fragment5 est directement reflétée dans le snapshot du fragment5 ; si le fragment5 est mis à jour après que le lecteur soit passé au fragment80, la mise à jour sera appliquée ultérieurement depuis le Binlog pendant la phase incrémentielle.
ImportantLorsque cette option est activée, les modifications survenant pendant ou après l'analyse d'un fragment sont toujours transmises depuis le Binlog en phase incrémentielle et peuvent entraîner des doublons. Seule la sémantique « au moins une fois » (at-least-once) est garantie. Activez cette option uniquement si le puits (sink) en aval supporte les écritures idempotentes par clé primaire.
RemarquePris en charge uniquement à partir de la version VVR 11.1.
scan.incremental.snapshot.unbounded-chunk-first.enabled
Indique s'il faut distribuer en priorité les fragments non bornés (unbounded chunks) pendant la phase de lecture du snapshot.
Non
BOOELEAN
false
Valeurs valides :
true : Distribue en priorité les fragments non bornés pendant la phase de lecture du snapshot.
false (par défaut) : Ne distribue pas en priorité les fragments non bornés pendant la phase de lecture du snapshot.
Il s'agit d'une fonctionnalité expérimentale. Son activation peut réduire le risque d'erreurs OOM (Out Of Memory) sur le TaskManager lors de la synchronisation du dernier fragment pendant la phase de snapshot. Ajoutez ce paramètre avant le premier démarrage du job.
RemarquePris en charge uniquement à partir de la version VVR 11.1.
binlog.session.network.timeout
Le délai d'expiration réseau en lecture/écriture pour la connexion au journal binaire.
Non
DURATION
10m
Si la valeur est définie sur 0s, le délai d'expiration par défaut du serveur MySQL est utilisé.
RemarquePris en charge uniquement à partir de la version VVR 11.5.
scan.rate-limit.records-per-second
Limite le nombre maximal d'enregistrements envoyés par la source par seconde.
Non
LONG
Aucun
Ce paramètre s'applique aux scénarios nécessitant de limiter la lecture des données. Cette limite est effective tant en phase complète qu'en phase incrémentielle.
La métrique
numRecordsOutPerSecondde la source reflète le nombre d'enregistrements produits par l'ensemble du flux de données par seconde. Vous pouvez ajuster ce paramètre en fonction de cette métrique.En phase de lecture complète des données, il est généralement nécessaire de réduire le nombre de lignes lues dans chaque lot. Pour ce faire, vous pouvez diminuer la valeur du paramètre
scan.incremental.snapshot.chunk.size.RemarquePris en charge uniquement à partir de la version VVR 11.5.
scan.binlog.tolerate.gtid-holes
L'activation de ce paramètre ignore les lacunes dans la séquence GTID, permettant au job de contourner les événements discontinus et de poursuivre son exécution.
Non
BOOLEAN
false
Avant d'activer ce paramètre, vous devez vous assurer que l'offset de démarrage du job n'a pas expiré. Si le job démarre à partir d'un offset GTID effacé ou expiré, le moteur ignorera continuellement les journaux manquants, ce qui entraînera une perte de données.
RemarqueCe paramètre est pris en charge uniquement à partir de la version VVR 11.6.
-
Paramètres spécifiques aux tables de dimension
Paramètre
Description
Obligatoire
Type de données
Valeur par défaut
Notes
url
L'URL JDBC MySQL.
Non
STRING
Aucune
Le format de l'URL est :
jdbc:mysql://<endpoint>:<port>/<database_name>.lookup.max-retries
Le nombre maximal de tentatives après un échec de lecture des données.
Non
INTEGER
3
Pris en charge uniquement dans VVR 6.0.7 et versions ultérieures.
lookup.cache.strategy
La politique de cache.
Non
STRING
Aucune
Les politiques de cache prises en charge sont None, LRU et ALL. Pour plus d'informations sur les valeurs, consultez Instructions JOIN pour les tables de dimension.
RemarqueLorsque vous utilisez la politique de cache LRU, vous devez également configurer le paramètre lookup.cache.max-rows.
lookup.cache.max-rows
Le nombre maximal de lignes mises en cache.
Non
INTEGER
100000
Si vous sélectionnez la politique de cache LRU, vous devez définir la taille du cache.
Si vous sélectionnez la politique de cache ALL, il n'est pas nécessaire de définir la taille du cache.
lookup.cache.ttl
La durée de vie (TTL) du cache.
Non
DURATION
10 s
La configuration de lookup.cache.ttl dépend de lookup.cache.strategy :
Si lookup.cache.strategy est défini sur None, vous n'avez pas besoin de configurer lookup.cache.ttl. Cela signifie que le cache n'expire pas.
Si lookup.cache.strategy est défini sur LRU, lookup.cache.ttl correspond à la durée de vie du cache. Par défaut, le cache n'expire pas.
Si lookup.cache.strategy est défini sur ALL, lookup.cache.ttl correspond au temps de chargement du cache. Par défaut, le cache n'est pas rechargé.
Utilisez un format temporel, tel que 1min ou 10s.
lookup.max-join-rows
Le nombre maximal de résultats renvoyés lorsqu'un enregistrement de la table principale correspond à des enregistrements de la table de dimension.
Non
INTEGER
1024
Aucune.
lookup.filter-push-down.enabled
Indique s'il faut activer le pushdown de filtre pour la table de dimension.
Non
BOOLEAN
false
Valeurs valides :
true : Active le pushdown de filtre pour la table de dimension. Lors du chargement des données depuis la table de base de données MySQL, la table de dimension filtre les données à l'avance en fonction des conditions définies dans la tâche SQL.
false (par défaut) : Désactive le pushdown de filtre pour la table de dimension. Lors du chargement des données depuis la table de base de données MySQL, la table de dimension charge toutes les données.
RemarquePris en charge uniquement dans VVR 8.0.7 et versions ultérieures.
ImportantLe pushdown de table de dimension ne doit être activé que lorsqu'une table Flink est utilisée comme table de dimension. Les tables source MySQL ne prennent pas en charge l'activation du pushdown de filtre. Si une table Flink est utilisée à la fois comme table source et comme table de dimension, et que le pushdown de filtre est activé pour la table de dimension, vous devez explicitement définir cet élément de configuration sur false pour la table source à l'aide de SQL Hints. Sinon, la tâche peut s'exécuter de manière anormale.
-
Uniquement pour les tables sink
Paramètre
Description
Obligatoire
Type de données
Valeur par défaut
Notes
url
L'URL JDBC MySQL.
Non
STRING
Aucune
Le format de l'URL est :
jdbc:mysql://<endpoint>:<port>/<database_name>.sink.max-retries
Nombre maximal de tentatives après un échec d'écriture des données.
Non
INTEGER
3
Aucune.
sink.buffer-flush.batch-size
Nombre de lignes dans une écriture par lot unique.
Non
INTEGER
4096
Aucune.
sink.buffer-flush.max-rows
Nombre de lignes de données mises en cache en mémoire.
Non
INTEGER
10000
Ce paramètre ne prend effet que si une clé primaire est spécifiée.
sink.buffer-flush.interval
Intervalle de vidage du cache. Si les données du cache ne remplissent pas les conditions de sortie après le temps d'attente spécifié, le système sort automatiquement toutes les données du cache.
Non
DURATION
1s
Aucune.
sink.ignore-delete
Indique s'il faut ignorer les opérations de suppression (DELETE) des données.
Non
BOOLEAN
false
Lorsque le flux généré par Flink SQL inclut des enregistrements delete ou update-before, des incohérences de données peuvent survenir si plusieurs tâches de sortie mettent à jour simultanément différents champs de la même table.
Par exemple, après la suppression d'un enregistrement, une autre tâche met à jour uniquement certains champs. Les champs non mis à jour deviennent alors null ou prennent leurs valeurs par défaut, ce qui entraîne des erreurs de données.
En définissant sink.ignore-delete sur true, vous pouvez ignorer les opérations DELETE et UPDATE_BEFORE en amont pour éviter ces problèmes.
RemarqueUPDATE_BEFORE fait partie du mécanisme de rétractation de Flink, utilisé pour « rétracter » l'ancienne valeur lors d'une opération de mise à jour.
Lorsque ignoreDelete = true, tous les enregistrements de type DELETE et UPDATE_BEFORE sont ignorés. Seuls les enregistrements INSERT et UPDATE_AFTER sont traités.
sink.ignore-delete-mode
Stratégie de gestion des enregistrements de type delete après l'ignorance des opérations DELETE.
Non
STRING
ALL
Valeurs valides :
ALL : Ignore les enregistrements -D et -U.
REAL_DELETE : Ignore uniquement les enregistrements -D.
UPDATE_BEFORE : Ignore uniquement les enregistrements -U.
RemarqueCette option est prise en charge uniquement dans le moteur Realtime Compute VVR 11.8 et versions ultérieures.
Effectif uniquement lorsque sink.ignore-delete=true. Une configuration isolée de ce paramètre entraîne une erreur.
sink.ignore-null-when-update
Lors de la mise à jour des données, indique s'il faut mettre à jour le champ correspondant avec la valeur null ou ignorer la mise à jour de ce champ si la valeur du champ de données entrant est null.
Non
BOOLEAN
false
Valeurs valides :
true : Ne met pas à jour le champ. Ce paramètre ne peut être défini sur true que si une clé primaire est configurée pour la table Flink. Lorsque la valeur est true :
Pour les versions VVR 8.0.6 et antérieures, la table de destination (sink) ne prend pas en charge l'écriture par lot.
Pour les versions VVR 8.0.7 et ultérieures, la table de destination (sink) prend en charge l'écriture par lot.
L'écriture par lot peut améliorer considérablement l'efficacité d'écriture et le débit global, mais elle introduit une latence des données et un risque d'erreurs OOM (Out Of Memory). Vous devez donc effectuer un compromis en fonction de votre scénario métier.
false : Met à jour le champ avec la valeur null.
RemarqueCe paramètre est pris en charge uniquement dans les versions VVR 8.0.5 et ultérieures.
Mappage des types
-
Tables source CDC
Type de champ MySQL CDC
Type de champ Flink
TINYINT
TINYINT
SMALLINT
SMALLINT
TINYINT UNSIGNED
TINYINT UNSIGNED ZEROFILL
INT
INT
MEDIUMINT
SMALLINT UNSIGNED
SMALLINT UNSIGNED ZEROFILL
BIGINT
BIGINT
INT UNSIGNED
INT UNSIGNED ZEROFILL
MEDIUMINT UNSIGNED
MEDIUMINT UNSIGNED ZEROFILL
BIGINT UNSIGNED
DECIMAL(20, 0)
BIGINT UNSIGNED ZEROFILL
SERIAL
FLOAT [UNSIGNED] [ZEROFILL]
FLOAT
DOUBLE [UNSIGNED] [ZEROFILL]
DOUBLE
DOUBLE PRECISION [UNSIGNED] [ZEROFILL]
REAL [UNSIGNED] [ZEROFILL]
NUMERIC(p, s) [UNSIGNED] [ZEROFILL]
DECIMAL(p, s)
DECIMAL(p, s) [UNSIGNED] [ZEROFILL]
BOOLEAN
BOOLEAN
TINYINT(1)
DATE
DATE
TIME [(p)]
TIME [(p)] [WITHOUT TIME ZONE]
DATETIME [(p)]
TIMESTAMP [(p)] [WITHOUT TIME ZONE]
TIMESTAMP [(p)]
TIMESTAMP [(p)]
TIMESTAMP [(p)] WITH LOCAL TIME ZONE
CHAR(n)
STRING
VARCHAR(n)
TEXT
BINARY
BYTES
VARBINARY
BLOB
ImportantN'utilisez pas le type TINYINT(1) dans MySQL pour stocker des valeurs autres que 0 et 1. Lorsque la propriété property-version=0, la table source MySQL CDC mappe par défaut TINYINT(1) au type BOOLEAN dans Flink. Cela peut entraîner des inexactitudes dans les données. Pour utiliser le type TINYINT(1) afin de stocker des valeurs autres que 0 et 1, consultez le paramètre de configuration catalog.table.treat-tinyint1-as-boolean.
-
Tables de dimension et tables de destination (sink)
Type de champ MySQL
Type de champ Flink
TINYINT
TINYINT
SMALLINT
SMALLINT
TINYINT UNSIGNED
INT
INT
MEDIUMINT
SMALLINT UNSIGNED
BIGINT
BIGINT
INT UNSIGNED
BIGINT UNSIGNED
DECIMAL(20, 0)
FLOAT
FLOAT
DOUBLE
DOUBLE
DOUBLE PRECISION
NUMERIC(p, s)
DECIMAL(p, s)
Remarqueoù p <= 38.
DECIMAL(p, s)
BOOLEAN
BOOLEAN
TINYINT(1)
DATE
DATE
TIME [(p)]
TIME [(p)] [WITHOUT TIME ZONE]
DATETIME [(p)]
TIMESTAMP [(p)] [WITHOUT TIME ZONE]
TIMESTAMP [(p)]
CHAR(n)
CHAR(n)
VARCHAR(n)
VARCHAR(n)
BIT(n)
BINARY(⌈n/8⌉)
BINARY(n)
BINARY(n)
VARBINARY(N)
VARBINARY(N)
TINYTEXT
STRING
TEXT
MEDIUMTEXT
LONGTEXT
TINYBLOB
BYTES
ImportantFlink prend uniquement en charge les enregistrements de type BLOB MySQL dont la taille est inférieure ou égale à 2 147 483 647 (2^31 - 1) octets.
BLOB
MEDIUMBLOB
LONGBLOB
Exemples d'utilisation
-
Table source CDC
CREATE TEMPORARY TABLE mysqlcdc_source ( order_id INT, order_date TIMESTAMP(0), customer_name STRING, price DECIMAL(10, 5), product_id INT, order_status BOOLEAN, PRIMARY KEY(order_id) NOT ENFORCED ) WITH ( 'connector' = 'mysql', 'hostname' = '<yourHostname>', 'port' = '3306', 'username' = '<yourUsername>', 'password' = '<yourPassword>', 'database-name' = '<yourDatabaseName>', 'table-name' = '<yourTableName>' ); CREATE TEMPORARY TABLE blackhole_sink( order_id INT, customer_name STRING ) WITH ( 'connector' = 'blackhole' ); INSERT INTO blackhole_sink SELECT order_id, customer_name FROM mysqlcdc_source; -
Table de dimension
CREATE TEMPORARY TABLE datagen_source( a INT, b BIGINT, c STRING, `proctime` AS PROCTIME() ) WITH ( 'connector' = 'datagen' ); CREATE TEMPORARY TABLE mysql_dim ( a INT, b VARCHAR, c VARCHAR ) WITH ( 'connector' = 'mysql', 'hostname' = '<yourHostname>', 'port' = '3306', 'username' = '<yourUsername>', 'password' = '<yourPassword>', 'database-name' = '<yourDatabaseName>', 'table-name' = '<yourTableName>' ); CREATE TEMPORARY TABLE blackhole_sink( a INT, b STRING ) WITH ( 'connector' = 'blackhole' ); INSERT INTO blackhole_sink SELECT T.a, H.b FROM datagen_source AS T JOIN mysql_dim FOR SYSTEM_TIME AS OF T.`proctime` AS H ON T.a = H.a; -
Table de destination
CREATE TEMPORARY TABLE datagen_source ( `name` VARCHAR, `age` INT ) WITH ( 'connector' = 'datagen' ); CREATE TEMPORARY TABLE mysql_sink ( `name` VARCHAR, `age` INT ) WITH ( 'connector' = 'mysql', 'hostname' = '<yourHostname>', 'port' = '3306', 'username' = '<yourUsername>', 'password' = '<yourPassword>', 'database-name' = '<yourDatabaseName>', 'table-name' = '<yourTableName>' ); INSERT INTO mysql_sink SELECT * FROM datagen_source; -
Source de données pour l'ingestion de données
source: type: mysql name: MySQL Source hostname: ${mysql.hostname} port: ${mysql.port} username: ${mysql.username} password: ${mysql.password} tables: ${mysql.source.table} server-id: 7601-7604 sink: type: values name: Values Sink print.enabled: true sink.print.logger: true
À propos des tables sources MySQL CDC
-
Fonctionnement
Lorsqu'une table source MySQL CDC démarre, elle analyse l'intégralité de la table, la divise en plusieurs segments selon la clé primaire et enregistre le décalage actuel du journal binaire. La table source utilise ensuite un algorithme d'instantané incrémentiel pour lire les données de chaque segment à l'aide d'instructions SELECT. Le job effectue régulièrement des points de contrôle pour enregistrer les segments terminés. En cas de basculement, le job poursuit la lecture des données depuis les segments inachevés. Une fois tous les segments lus, le job commence à lire les enregistrements de modification incrémentiels à partir du décalage du journal binaire précédemment enregistré. Le job Flink continue d'effectuer des points de contrôle périodiques pour enregistrer le décalage du journal binaire. Si le job subit un basculement, il reprend le traitement à partir du dernier décalage du journal binaire enregistré, ce qui permet d'atteindre une sémantique exactement une fois.
Pour une explication plus détaillée de l'algorithme d'instantané incrémentiel, consultez Connecteur MySQL CDC.
-
Métadonnées
Les métadonnées sont utiles dans les scénarios où les données provenant de bases de données et de tables fragmentées sont fusionnées et synchronisées. En effet, après la fusion, les entreprises souhaitent souvent distinguer la base de données et la table source pour chaque enregistrement de données. Les colonnes de métadonnées permettent d'accéder aux informations sur le nom de la base de données et de la table source. Vous pouvez ainsi facilement fusionner plusieurs tables fragmentées en une seule table de destination à l'aide des colonnes de métadonnées.
La source MySQL CDC prend en charge la syntaxe des colonnes de métadonnées. Vous pouvez accéder aux métadonnées suivantes via ces colonnes.
Clé de métadonnée
Type de métadonnée
Description
database_name
STRING NOT NULL
Nom de la base de données contenant la ligne.
table_name
STRING NOT NULL
Nom de la table contenant la ligne.
op_ts
TIMESTAMP_LTZ(3) NOT NULL
Heure à laquelle la ligne a été modifiée dans la base de données. Si l'enregistrement provient des données historiques de la table plutôt que du journal binaire, cette valeur est toujours 0.
RemarqueCe champ n'est précis qu'à la seconde près.
op_type
STRING NOT NULL
Type de modification de la ligne.
+I : message INSERT
-D : message DELETE
-U : message UPDATE_BEFORE
+U : message UPDATE_AFTER
RemarquePris en charge uniquement dans VVR 8.0.7 et versions ultérieures.
query_log
STRING NOT NULL
Vous pouvez lire l'enregistrement du journal des requêtes MySQL pour cette ligne.
RemarqueMySQL doit avoir le paramètre binlog_rows_query_log_events activé pour enregistrer les journaux des requêtes.
L'exemple de code suivant montre comment fusionner et synchroniser plusieurs tables de commandes provenant de plusieurs bases de données fragmentées dans une instance MySQL vers une table holo_orders dans Hologres.
CREATE TEMPORARY TABLE mysql_orders ( db_name STRING METADATA FROM 'database_name' VIRTUAL, -- Read the database name. table_name STRING METADATA FROM 'table_name' VIRTUAL, -- Read the table name. operation_ts TIMESTAMP_LTZ(3) METADATA FROM 'op_ts' VIRTUAL, -- Read the change time. op_type STRING METADATA FROM 'op_type' VIRTUAL, -- Read the change type. order_id INT, order_date TIMESTAMP(0), customer_name STRING, price DECIMAL(10, 5), product_id INT, order_status BOOLEAN, PRIMARY KEY(order_id) NOT ENFORCED ) WITH ( 'connector' = 'mysql-cdc', 'hostname' = 'localhost', 'port' = '3306', 'username' = 'flinkuser', 'password' = 'flinkpw', 'database-name' = 'mydb_.*', -- Regular expression to match multiple sharded databases. 'table-name' = 'orders_.*' -- Regular expression to match multiple sharded tables. ); INSERT INTO holo_orders SELECT * FROM mysql_orders;Sur la base du code ci-dessus, si le paramètre
scan.read-changelog-as-append-only.enabledest défini sur true dans la clause WITH, le résultat de sortie varie en fonction du paramétrage de la clé primaire de la table en aval :Si la clé primaire de la table en aval est
order_id, le résultat de sortie contient uniquement la dernière modification pour chaque clé primaire de la table en amont. Pour les données dont la dernière modification pour une clé primaire était une opération de suppression, vous pouvez voir un enregistrement dans la table en aval avec la même clé primaire et unop_typede -D.Si la clé primaire de la table en aval est
order_id,operation_tsetop_type, le résultat de sortie contient les modifications complètes pour chaque clé primaire de la table en amont.
-
Prise en charge des expressions régulières
La table source MySQL CDC permet d'utiliser des expressions régulières dans le nom de la table ou de la base de données pour faire correspondre plusieurs tables ou bases de données. L'exemple de code suivant montre comment spécifier plusieurs tables à l'aide d'une expression régulière.
CREATE TABLE products ( db_name STRING METADATA FROM 'database_name' VIRTUAL, table_name STRING METADATA FROM 'table_name' VIRTUAL, operation_ts TIMESTAMP_LTZ(3) METADATA FROM 'op_ts' VIRTUAL, order_id INT, order_date TIMESTAMP(0), customer_name STRING, price DECIMAL(10, 5), product_id INT, order_status BOOLEAN, PRIMARY KEY(order_id) NOT ENFORCED ) WITH ( 'connector' = 'mysql-cdc', 'hostname' = 'localhost', 'port' = '3306', 'username' = 'root', 'password' = '123456', 'database-name' = '(^(test).*|^(tpc).*|txc|.*[p$]|t{2})', -- Regular expression to match multiple databases. 'table-name' = '(t[5-8]|tt)' -- Regular expression to match multiple tables. );Les expressions régulières de l'exemple sont expliquées comme suit :
^(test).*est un exemple de mise en correspondance par préfixe. Cette expression peut correspondre à des noms de bases de données commençant par « test », tels que « test1 » ou « test2 »..*[p$]est un exemple de mise en correspondance par suffixe. Cette expression peut correspondre à des noms de bases de données se terminant par « p », tels que « cdcp » ou « edcp ».txcest une correspondance spécifique. Elle peut correspondre à un nom de base de données exactement égal à « txc ».
Lorsque MySQL CDC fait correspondre un nom de table qualifié complet, il utilise le modèle
database-name.table-namepour identifier de manière unique une table. Par exemple, le modèle(^(test).|^(tpc).|txc|.*[p$]|t{2}).(t[ 5-8]|tt)peut correspondre à des tables telles quetxc.ttettest2.test5dans la base de données.ImportantDans la configuration d'un job SQL, les paramètres
table-nameetdatabase-namene prennent pas en charge l'utilisation d'une virgule (,) pour spécifier plusieurs tables ou bases de données.Pour faire correspondre plusieurs tables ou utiliser plusieurs expressions régulières, reliez-les avec une barre verticale (|) et placez-les entre parenthèses. Par exemple, pour lire les tables
useretproduct, vous pouvez définirtable-namesur(user|product).Si une expression régulière contient une virgule, vous devez la réécrire à l'aide de l'opérateur barre verticale (|). Par exemple, l'expression régulière
mytable_\d{1, 2}doit être réécrite sous la forme équivalente(mytable_\d{1}|mytable_\d{2})pour éviter d'utiliser une virgule.
-
Contrôle de la concurrence
Le connecteur MySQL prend en charge la lecture multithread des données complètes, ce qui peut améliorer l'efficacité du chargement des données. Conjuguée à la fonctionnalité de réglage automatique Autopilot dans la console Realtime Compute for Apache Flink, le connecteur peut automatiquement réduire la montée en puissance pendant la phase incrémentielle après la fin de la lecture multithread afin d'économiser des ressources de calcul.
Dans la console de développement de Realtime Compute for Apache Flink, vous pouvez définir la concurrence d'un job en mode basique ou en mode expert sur la page Resource Configuration.
-
La concurrence définie en mode basique correspond à la concurrence globale pour l'ensemble du job.
Par exemple, lorsque le parallelism est défini sur
8en basic mode, l'server-iddans la clause SQL WITH doit être configuré comme une plage continue (telle que'404-412'). Le mode expert permet de définir la concurrence pour un VERTEX spécifique selon les besoins.
Pour plus d'informations sur la configuration des ressources, consultez Configurer les informations de déploiement pour un job.
ImportantQue vous soyez en mode basique ou en mode expert, lors de la définition de la concurrence, la plage d'ID de serveur déclarée dans la table doit être supérieure ou égale à la concurrence du job. Par exemple, si la plage d'ID de serveur est
5404-5412, il y a neuf ID de serveur uniques. Par conséquent, la concurrence du job peut être définie sur un maximum de 9. Différents jobs pour la même instance MySQL ne doivent pas avoir de plages d'ID de serveur qui se chevauchent. Cela signifie que chaque job doit être explicitement configuré avec un ID de serveur ou une plage d'ID de serveur différent. -
-
Réduction automatique de la montée en puissance Autopilot
La phase de données complètes accumule une grande quantité de données historiques. Pour améliorer l'efficacité de la lecture, les données historiques sont généralement lues en parallèle. Dans la phase incrémentielle du journal binaire, comme la quantité de données du journal binaire est faible et afin de garantir l'ordre global, la lecture monothread est généralement suffisante. La fonctionnalité de réglage automatique permet d'équilibrer les différentes exigences en matière de ressources des phases complète et incrémentielle pour optimiser les performances et les ressources.
Le réglage automatique surveille le trafic de chaque tâche de la source MySQL CDC. Lors du passage à la phase de lecture des journaux binaires, si une seule tâche est chargée de cette lecture tandis que les autres sont inactives, le réglage automatique réduit automatiquement le nombre d'unités de calcul (CU) et le niveau de parallélisme de la source. Pour activer le réglage automatique, définissez le mode de réglage automatique sur Active dans la page O&M du job.
RemarqueL'intervalle minimal par défaut pour déclencher une réduction du parallélisme est de 24 heures. Pour plus d'informations sur les paramètres et les détails du réglage automatique, consultez Configurer le réglage automatique.
-
Modes de démarrage
Utilisez l'option de configuration
scan.startup.modepour spécifier le mode de démarrage de la table source MySQL CDC. Les options disponibles sont les suivantes :initial (par défaut) : Lors du premier démarrage ou d'un démarrage sans état, effectue une lecture complète de la table de base de données, puis passe en mode incrémentiel pour lire le journal binaire.
earliest-offset : Ignore la phase d'instantané et commence la lecture à partir du décalage de journal binaire le plus ancien disponible.
latest-offset : Ignore la phase d'instantané et commence la lecture à partir de la fin du journal binaire. Dans ce mode, la table source ne peut lire que les modifications de données survenues après le démarrage du job.
specific-offset : Ignore la phase d'instantané et commence la lecture à partir d'un décalage de journal binaire spécifié. Le décalage peut être indiqué par le nom de fichier et la position du journal binaire, ou par un ensemble GTID.
timestamp : Ignore la phase d'instantané et commence la lecture des événements du journal binaire à partir d'un horodatage spécifié.
Un démarrage sans état ne réutilise aucun état précédent. Le connecteur source le traite comme un premier démarrage, de sorte que
scan.startup.modes'applique à nouveau. Pour plus d'informations sur les modes de démarrage des déploiements, consultez Démarrer un déploiement.Exemple d'utilisation :
CREATE TABLE mysql_source (...) WITH ( 'connector' = 'mysql-cdc', 'scan.startup.mode' = 'earliest-offset', -- Start from the earliest offset. 'scan.startup.mode' = 'latest-offset', -- Start from the latest offset. 'scan.startup.mode' = 'specific-offset', -- Start from a specific offset. 'scan.startup.mode' = 'timestamp', -- Start from a specific timestamp. 'scan.startup.specific-offset.file' = 'mysql-bin.000003', -- Specify the binary log filename in specific-offset mode. 'scan.startup.specific-offset.pos' = '4', -- Specify the binary log position in specific-offset mode. 'scan.startup.specific-offset.gtid-set' = '24DA167-0C0C-11E8-8442-00059A3C7B00:1-19', -- Specify the GTID set in specific-offset mode. 'scan.startup.timestamp-millis' = '1667232000000' -- Specify the startup timestamp in timestamp mode. ... )ImportantLa source MySQL imprime le décalage actuel dans le journal au niveau INFO lors d'un point de contrôle. Le préfixe du journal est
Binlog offset on checkpoint {checkpoint-id}. Ce journal vous aide à démarrer un job à partir d'un décalage de point de contrôle spécifique.Si la table en cours de lecture a subi des modifications de schéma, un démarrage à partir de
earliest-offset,specific-offsetoutimestamppeut provoquer une erreur. En effet, le lecteur Debezium conserve en interne le dernier schéma de table connu, et les données antérieures dont le schéma ne correspond pas ne peuvent pas être analysées correctement.
-
À propos des tables sources CDC sans clé primaire
L'utilisation d'une table sans clé primaire nécessite de définir
scan.incremental.snapshot.chunk.key-column, et seule une colonne non nulle peut être sélectionnée.-
La sémantique de traitement d'une table source CDC sans clé primaire dépend du comportement de la colonne spécifiée par
scan.incremental.snapshot.chunk.key-column:Si la colonne spécifiée n'est pas mise à jour, la sémantique exactly-once est garantie.
Si la colonne spécifiée est mise à jour, seule la sémantique at-least-once est garantie. Toutefois, vous pouvez garantir l'exactitude des données en combinant cette approche avec le système en aval, en spécifiant une clé primaire en aval et en utilisant des opérations idempotentes.
-
Lire les journaux de sauvegarde d'Alibaba Cloud ApsaraDB RDS for MySQL
La table source MySQL CDC prend en charge la lecture des journaux de sauvegarde d'Alibaba Cloud ApsaraDB RDS for MySQL. Cette fonctionnalité est utile lorsque la phase de données complètes prend beaucoup de temps et que les fichiers de journaux binaires locaux ont été automatiquement supprimés, mais que les fichiers de sauvegarde téléchargés automatiquement ou manuellement sont toujours disponibles.
Exemple d'utilisation :
CREATE TABLE mysql_source (...) WITH ( 'connector' = 'mysql-cdc', 'rds.region-id' = 'cn-beijing', 'rds.access-key-id' = 'xxxxxxxxx', 'rds.access-key-secret' = 'xxxxxxxxx', 'rds.db-instance-id' = 'rm-xxxxxxxxxxxxxxxxx', 'rds.main-db-id' = '12345678', 'rds.download.timeout' = '60s' ... ) -
Activer la réutilisation de la source CDC
Dans un même job, plusieurs tables sources MySQL CDC démarrent plusieurs clients de journaux binaires. Si toutes les tables sources se trouvent sur la même instance, cela augmente la charge sur la base de données. Pour plus d'informations, consultez FAQ MySQL CDC.
Solution
Les versions VVR 8.0.7 et ultérieures prennent en charge la réutilisation des sources MySQL CDC. Cette fonctionnalité fusionne les tables sources MySQL CDC qui peuvent l'être. La fusion se produit lorsque les configurations des tables sources sont identiques, à l'exception du nom de la base de données, du nom de la table et de
server-id. Le moteur fusionne automatiquement les sources MySQL CDC au sein d'un même job.Procédure
-
Utilisez la commande
SETdans votre job SQL :SET 'table.optimizer.source-merge.enabled' = 'true'; # (For VVR 8.0.8 and 8.0.9) Also set this item: SET 'sql-gateway.exec-plan.enabled' = 'false';La réutilisation est activée par défaut dans les versions VVR 11.1 et ultérieures.
Démarrez le job sans état. Étant donné que la modification de la configuration de réutilisation de la source change la topologie du job, vous devez démarrer le job sans état. Sinon, le démarrage du job peut échouer ou entraîner une perte de données. Si une source est fusionnée, vous verrez un nœud
MergetableSourceScandans la topologie.
ImportantAprès avoir activé la réutilisation, ne désactivez pas le chaînage des opérateurs. Si vous définissez
pipeline.operator-chainingsurfalse, cela augmente la surcharge liée à la sérialisation et à la désérialisation des données. Plus il y a de sources fusionnées, plus la surcharge est importante.Dans la version VVR 8.0.7, la désactivation du chaînage des opérateurs provoque des problèmes de sérialisation.
-
Accélérer la lecture des journaux binaires
Lorsque vous utilisez le connecteur MySQL comme table source ou comme source d'ingestion de données, il analyse les fichiers de journaux binaires pour générer divers messages de modification pendant la phase incrémentielle. Les fichiers de journaux binaires enregistrent toutes les modifications de table au format binaire. Vous pouvez accélérer l'analyse des fichiers de journaux binaires de plusieurs manières.
-
Activez l'analyse parallèle et les filtres d'analyse (cette fonctionnalité nécessite Realtime Compute for Apache Flink avec Ververica Runtime (VVR) 8.0.7 ou une version ultérieure. Elle n'est pas disponible dans l'édition communautaire du connecteur MySQL CDC.)
Activez l'option
scan.only.deserialize.captured.tables.changelog.enabledpour analyser uniquement les événements de modification des tables spécifiées.Activez l'option
scan.parallel-deserialize-changelog.enabledpour utiliser plusieurs threads afin d'analyser le fichier de journal binaire et de transmettre les événements à la file d'attente du consommateur dans l'ordre. Lorsque vous activez cette option, vous devez généralement augmenter également leTaskManager CPU.
-
Optimisez les paramètres Debezium
debezium.max.queue.size: 162580 debezium.max.batch.size: 40960 debezium.poll.interval.ms: 50debezium.max.queue.size: Nombre maximal d'enregistrements que la file d'attente bloquante peut contenir. Lorsque Debezium lit un flux d'événements depuis la base de données, il place les événements dans une file d'attente bloquante avant de les écrire en aval. La valeur par défaut est 8192.debezium.max.batch.size: Nombre maximal d'événements traités par le connecteur à chaque itération. La valeur par défaut est 2048.debezium.poll.interval.ms: Nombre de millisecondes pendant lesquelles le connecteur doit attendre avant de demander de nouveaux événements de modification. La valeur par défaut est de 1000 millisecondes, soit 1 seconde.
Exemple d'utilisation :
CREATE TABLE mysql_source (...) WITH (
'connector' = 'mysql-cdc',
-- Debezium configuration
'debezium.max.queue.size' = '162580',
'debezium.max.batch.size' = '40960',
'debezium.poll.interval.ms' = '50',
-- Enable parsing filter
'scan.only.deserialize.captured.tables.changelog.enabled' = 'true', -- Parse only the change events of specified tables.
...
)
source:
type: mysql
name: MySQL Source
hostname: ${mysql.hostname}
port: ${mysql.port}
username: ${mysql.username}
password: ${mysql.password}
tables: ${mysql.source.table}
server-id: 7601-7604
# Debezium configuration
debezium.max.queue.size: 162580
debezium.max.batch.size: 40960
debezium.poll.interval.ms: 50
# Enable parsing filter
scan.only.deserialize.captured.tables.changelog.enabled: true
La capacité de consommation des journaux binaires de l'édition Enterprise de MySQL CDC est de 85 Mo/s, soit environ le double de celle de la version open source communautaire. Lorsque la vitesse de génération des fichiers de journaux binaires dépasse 85 Mo/s (c'est-à-dire un fichier de 512 Mo toutes les 6 secondes), la latence du job Flink continue d'augmenter. La latence de traitement diminue progressivement une fois que la vitesse de génération des fichiers de journaux binaires ralentit. Si un fichier de journal binaire contient une transaction volumineuse, la latence de traitement peut augmenter temporairement. La latence diminue après la lecture du journal correspondant à cette transaction.
Diagnostiquer la latence des données pour optimiser le débit du job
Si vous rencontrez une latence des données pendant la phase incrémentielle, analysez le problème en suivant ces étapes :
-
Vérifiez les métriques
currentFetchEventTimeLagetcurrentEmitEventTimeLagsur la page Overview. La métriquecurrentFetchEventTimeLagreprésente la latence de lecture des données depuis le journal binaire. La métriquecurrentEmitEventTimeLagreprésente la latence de lecture des données pour les tables pertinentes pour le job depuis le journal binaire.Scénario
Description
currentFetchEventTimeLagest faible, tandis quecurrentEmitEventTimeLagest élevé et rarement mis à jour.Une valeur faible de
currentFetchEventTimeLagindique que l'extraction du journal binaire depuis la base de données est efficace. Cependant, le journal binaire contient peu de données pour les tables que le job doit lire. Par conséquent,currentEmitEventTimeLagest rarement mis à jour. Il s'agit d'un comportement attendu.currentFetchEventTimeLagetcurrentEmitEventTimeLagsont tous deux élevés.Cela indique que la table source présente de mauvaises performances de lecture. Vous pouvez passer aux étapes suivantes de cette section pour l'optimisation.
La contre-pression peut réduire la vitesse à laquelle la source envoie des données aux opérateurs en aval. Vous pouvez observer que
sourceIdleTimeaugmente périodiquement, et quecurrentFetchEventTimeLagetcurrentEmitEventTimeLagaugmentent continuellement. Pour résoudre ce problème, augmentez le parallélisme du nœud à l'origine de la contre-pression.Vérifiez la métrique TM CPU Usage sur la page CPU et la métrique TM GC Time sur la page JVM pour déterminer s'il existe une insuffisance de ressources CPU ou mémoire. Vous pouvez augmenter les ressources du job pour optimiser les performances de lecture. Vous pouvez également activer les paramètres mini-batch pour améliorer le débit. Pour plus d'informations, consultez Techniques d'optimisation haute performance de Flink SQL.
Si un opérateur SinkUpsertMaterializer avec un état volumineux existe dans le job, il peut affecter les performances de lecture. Envisagez d'augmenter le parallélisme du job ou d'éviter l'opérateur SinkUpsertMaterializer. Pour plus d'informations, consultez Éviter l'utilisation de SinkUpsertMaterializer. La suppression de l'opérateur SinkUpsertMaterializer d'un job existant nécessite un redémarrage sans état. En effet, la topologie du job change et un démarrage à partir de l'état existant peut entraîner l'échec du job ou une perte de données.
Définissez l'ID de serveur pour éviter les conflits de journaux binaires
Chaque client qui synchronise des données depuis une base de données possède un identifiant unique appelé ID de serveur. Si différents jobs utilisent le même ID de serveur, des conflits peuvent survenir et entraîner l'échec des jobs. Nous vous recommandons d'attribuer un ID de serveur différent à chaque source de données MySQL CDC.
-
Comment configurer l'ID de serveur
Vous pouvez spécifier l'ID de serveur dans une instruction DDL de table Flink ou en utilisant des indices SQL (SQL Hints).
Nous vous recommandons d'utiliser des indices SQL pour configurer l'ID de serveur plutôt que de le spécifier dans la clause WITH de la DDL de la table. Pour plus d'informations, consultez Indices SQL.
-
Configuration de l'ID de serveur pour différents scénarios
-
Instantané incrémentiel désactivé ou parallélisme égal à 1
Si l'instantané incrémentiel est désactivé ou si le parallélisme est égal à 1, vous pouvez spécifier un seul ID de serveur.
SELECT * FROM source_table /*+ OPTIONS('server-id'='123456') */ ; -
Instantané incrémentiel activé et parallélisme supérieur à 1
Lorsque l'instantané incrémentiel est activé et que le parallélisme est supérieur à 1, vous devez spécifier une plage d'IDs de serveur. Le nombre d'IDs de serveur disponibles dans la plage doit être au moins égal au parallélisme. Par exemple, si le parallélisme est de 3, vous pouvez utiliser la configuration suivante :
SELECT * FROM source_table /*+ OPTIONS('server-id'='123456-123458') */ ; -
Synchronisation des données avec CTAS
Lorsque vous synchronisez des données en utilisant CREATE TABLE AS (CTAS), les sources de données CDC ayant des configurations identiques sont automatiquement fusionnées et réutilisées. Dans ce cas, vous pouvez attribuer le même ID de serveur à plusieurs sources de données CDC. Pour plus d'informations, consultez Exemple 4 : Plusieurs instructions CTAS.
-
Un job avec plusieurs tables source MySQL CDC (hors CTAS)
Si un job contient plusieurs tables source MySQL CDC, n'utilise pas d'instructions CTAS et a la réutilisation de source désactivée, vous devez fournir un ID de serveur différent pour chaque table source CDC. De même, si l'instantané incrémentiel est activé et que le parallélisme est supérieur à 1, vous devez spécifier une plage d'IDs de serveur.
select * from source_table1 /*+ OPTIONS('server-id'='123456-123457') */ left join source_table2 /*+ OPTIONS('server-id'='123458-123459') */ on source_table1.id=source_table2.id;
-
Définissez les paramètres de chunk pour optimiser l'utilisation de la mémoire
Lorsqu'une table source MySQL CDC démarre, elle effectue une analyse complète de la table, divise la table en plusieurs chunks (fragments) en fonction de la clé primaire et enregistre la position actuelle du journal binaire. Le job utilise ensuite un algorithme d'instantané incrémentiel pour lire les données de chaque chunk séquentiellement à l'aide d'instructions SELECT. Le job effectue régulièrement un point de contrôle (checkpoint) pour enregistrer les chunks terminés. En cas de basculement, il reprend la lecture à partir du premier chunk inachevé. Une fois tous les chunks lus, le job passe à la lecture des modifications incrémentielles à partir de la position du journal binaire précédemment enregistrée. Le job Flink effectue des points de contrôle périodiques pour sauvegarder la position du journal binaire. Si un basculement se produit, le job reprend le traitement à partir de la dernière position sauvegardée, atteignant ainsi une sémantique exactement une fois (exactly-once).
Pour plus de détails sur l'algorithme d'instantané incrémentiel, consultez Connecteur MySQL CDC.
Pour les tables avec une clé primaire à colonne unique, les chunks sont divisés par défaut en fonction de cette clé. Pour les tables avec une clé primaire composite, la première colonne de la clé primaire est utilisée par défaut pour la division. Ververica Runtime (VVR) 6.0.7 et versions ultérieures prend en charge la lecture de tables source sans clé primaire. Vous devez définir le paramètre scan.incremental.snapshot.chunk.key-column pour spécifier une colonne non nullable pour la division.
Optimisation des paramètres de chunk
Les données des chunks et leurs métadonnées sont stockées en mémoire, ce qui peut parfois entraîner des erreurs d'épuisement de la mémoire (OOM). Vous pouvez ajuster les paramètres en fonction du composant qui rencontre l'erreur OOM :
-
JobManager
Le JobManager stocke les métadonnées de tous les chunks. Un nombre excessif de chunks peut provoquer une erreur OOM. Pour résoudre ce problème, augmentez la valeur de
scan.incremental.snapshot.chunk.sizeafin de réduire le nombre de chunks. Vous pouvez également augmenter la mémoire heap du JobManager en définissantjobmanager.memory.heap.sizedans votre configuration d'exécution. Pour plus d'informations, consultez Configuration des paramètres Flink. -
TaskManager
Le TaskManager lit les données de chaque chunk. Si un chunk contient trop de lignes, une erreur OOM peut se produire. Pour résoudre ce problème, diminuez la valeur de
scan.incremental.snapshot.chunk.sizeafin de réduire le nombre de lignes par chunk. Vous pouvez également augmenter la mémoire heap du TaskManager en augmentant la valeur deTaskManager Memorydans votre configuration d'exécution.Dans VVR 8.0.8 et versions antérieures, le dernier chunk peut contenir une grande quantité de données, ce qui peut provoquer une erreur OOM sur le TaskManager. Nous vous recommandons de mettre à niveau vers VVR 8.0.9 ou une version ultérieure pour éviter ce problème.
Pour une table source MySQL CDC avec une clé primaire composite, les chunks sont divisés par défaut en fonction de la première colonne de la clé. Si les données sont fortement asymétriques, avec de nombreuses lignes partageant la même valeur dans cette colonne, le chunk correspondant à cette valeur peut devenir très volumineux et provoquer une erreur OOM sur le TaskManager. Vous pouvez définir
scan.incremental.snapshot.chunk.key-columnpour spécifier une autre colonne de la clé primaire pour la division.
Accélérez les lectures pendant la phase d'instantané
Pendant la phase d'instantané, la table source MySQL lit les données d'instantané via une connexion JDBC. Utilisez les méthodes suivantes pour accélérer les lectures durant cette phase.
Augmentez le parallélisme de la source pour accélérer les lectures pendant la phase d'instantané.
Augmentez la valeur de
scan.incremental.snapshot.chunk.sizepour récupérer plus de données dans un seul chunk.Si la table de résultat en aval possède une clé primaire et prend en charge les écritures idempotentes, vous pouvez activer
scan.incremental.snapshot.backfill.skippour ignorer la lecture du journal binaire pour la partie de remplissage initial (backfill). Cela accélère le traitement pendant la phase d'instantané.
Activez la réutilisation de la source pour réduire les connexions aux journaux binaires
Lorsqu'un job inclut plusieurs tables source MySQL, vous pouvez activer la réutilisation de la source pour réduire la charge de la base de données en partageant une seule connexion au journal binaire. Cette fonctionnalité est disponible uniquement dans Realtime Compute for Apache Flink et n'est pas prise en charge dans l'édition communautaire du connecteur MySQL CDC.
Activez la fonctionnalité de réutilisation de la source dans un job SQL en utilisant la commande SET :
SET 'table.optimizer.source-merge.enabled' = 'true';
Nous vous recommandons d'activer la réutilisation de la source uniquement pour les nouveaux jobs. Si vous activez la réutilisation de la source pour un job existant, vous devez effectuer un redémarrage sans état. En effet, la réutilisation de la source modifie la topologie du job, et le démarrage à partir d'un état existant peut entraîner l'échec du job ou une perte de données.
Après avoir activé la réutilisation de la source, les tables source MySQL ayant les mêmes paramètres de configuration sont fusionnées. Si toutes les tables source de votre job partagent la même configuration, le nombre de connexions aux journaux binaires est calculé comme suit :
Pendant la phase d'instantané, le nombre de connexions aux journaux binaires est égal au parallélisme de la source.
Pendant la phase incrémentielle, le nombre de connexions aux journaux binaires est de 1.
Dans VVR 8.0.8 et 8.0.9, vous devez également définir
SET 'sql-gateway.exec-plan.enabled' = 'false';lorsque vous activez la réutilisation de la source CDC.Après avoir activé la réutilisation de la source CDC, ne définissez pas l'option de job
pipeline.operator-chainingsur false. La rupture de la chaîne d'opérateurs ajoute une surcharge de sérialisation et de désérialisation pour les données envoyées de la source aux opérateurs en aval. Plus il y a de sources fusionnées, plus la surcharge est importante.Dans Ververica Runtime (VVR) 8.0.7, le fait de définir
pipeline.operator-chainingsur false provoque un problème de sérialisation.
Lire les journaux binaires archivés depuis OSS
Lorsque vous utilisez une instance ApsaraDB RDS for MySQL comme source de données, vous pouvez lire les sauvegardes de journaux stockées dans OSS. Si le fichier correspondant à l'horodatage ou à la position du journal binaire spécifié est stocké dans OSS, Flink extrait automatiquement le fichier journal depuis OSS vers le cluster. Si le fichier est stocké localement dans la base de données, Flink bascule automatiquement vers la lecture via une connexion à la base de données. Cette fonctionnalité est disponible uniquement dans Realtime Compute for Apache Flink et n'est pas prise en charge dans l'édition communautaire du connecteur MySQL CDC.
Pour activer la lecture depuis les sauvegardes de journaux OSS, vous devez configurer les paramètres de connexion ApsaraDB RDS for MySQL. Exemple :
CREATE TABLE mysql_source (...) WITH (
'connector' = 'mysql-cdc',
'rds.region-id' = 'cn-beijing',
'rds.access-key-id' = 'your_access_key_id',
'rds.access-key-secret' = 'your_access_key_secret',
'rds.db-instance-id' = 'rm-xxxxxxxx', // The ID of the database instance.
'rds.main-db-id' = '12345678', // The ID of the primary database.
'rds.endpoint' = 'rds.aliyuncs.com'
...
)
FAQ
Pour plus d'informations sur les problèmes que vous pourriez rencontrer lors de l'utilisation des tables source CDC, consultez FAQ sur CDC.