Cette rubrique explique comment synchroniser une base de données MySQL entière vers Kafka. Cette approche réduit la charge que plusieurs jobs imposent à la base de données MySQL.
Contexte
Une table source MySQL CDC capture les données de MySQL et synchronise les modifications en temps réel de la table. Ce mécanisme est courant dans les scénarios de calcul complexes, par exemple lorsqu'une table sert de table de dimension dans une opération JOIN avec d'autres tables de données. Une seule table MySQL peut constituer une dépendance pour plusieurs jobs. Lorsque plusieurs jobs traitent des données provenant de la même table MySQL, la base de données ouvre plusieurs connexions, ce qui exerce une pression importante sur le serveur MySQL et le réseau.
Fonctionnement
Pour alléger la pression sur la base de données MySQL en amont, Realtime Compute for Apache Flink permet de synchroniser une base de données MySQL entière vers Kafka. Cette solution introduit Kafka comme couche intermédiaire et utilise un job d'ingestion de données Flink CDC pour synchroniser les données vers Kafka.
Un seul job synchronise en temps réel les données d'une base de données MySQL en amont vers Kafka. Chaque table MySQL est écrite dans un topic Kafka correspondant en mode upsert. Les jobs en aval utilisent ensuite un connecteur Kafka upsert pour lire les données des topics au lieu d'accéder directement aux tables MySQL. Cette méthode réduit efficacement la pression exercée par plusieurs jobs sur la base de données MySQL.

Limites
Chaque table MySQL que vous synchronisez doit posséder une clé primaire.
Vous pouvez utiliser des clusters Kafka auto-gérés, des clusters EMR Kafka ou ApsaraMQ for Kafka. Si vous utilisez ApsaraMQ for Kafka, la connexion s'effectue uniquement via un endpoint par défaut.
L'espace de stockage du cluster Kafka doit être supérieur à celui des tables sources. Dans le cas contraire, des données risquent d'être perdues en raison d'un espace insuffisant. Les topics créés pour la synchronisation de bases de données sont des topics compactés. Dans un topic compacté, seul le dernier message de chaque clé de message est conservé, mais les données n'expirent jamais. Cela signifie que le topic compacté stocke un volume de données globalement équivalent à la taille de la table source.
Scénario d'exemple
Prenons l'exemple d'un scénario d'analyse de vérification de commandes en temps réel comportant trois tables : une table utilisateur (user), une table de commandes (order) et une table de retours utilisateurs (feedback). Ces tables contiennent les données illustrées dans la figure suivante.
Pour afficher les informations de commande et les avis des utilisateurs, vous devez joindre la table user afin de récupérer les noms d'utilisateur depuis le champ name. L'exemple SQL suivant illustre cette opération.
-- Join order information with the user table to display the username and product name for each order.
SELECT order.id as order_id, product, user.name as user_name
FROM order LEFT JOIN user
ON order.user_id = user.id;
-- Join reviews with the user table to display the content of each review and the corresponding username.
SELECT feedback.id as feedback_id, comment, user.name as user_name
FROM feedback LEFT JOIN user
ON feedback.user_id = user.id;
Les deux jobs SQL précédents utilisent tous deux la table user. Lors de l'exécution, chaque job lit les données complètes et incrémentielles depuis MySQL. Une lecture complète nécessite la création d'une connexion MySQL, tandis qu'une lecture incrémentielle requiert la création d'un client Binlog. À mesure que le nombre de jobs augmente, la demande en ressources de connexion MySQL et de clients Binlog croît également, exerçant une pression considérable sur la base de données en amont. Pour soulager cette pression, utilisez un job d'ingestion de données Flink CDC afin de synchroniser en temps réel les données de la base de données MySQL en amont vers Kafka, permettant ainsi leur consommation par plusieurs jobs en aval.
Prérequis
Realtime Compute for Apache Flink est activé. Pour plus d'informations, consultez Activer Realtime Compute for Apache Flink.
ApsaraMQ for Kafka est activé. Pour plus d'informations, consultez Déployer une instance ApsaraMQ for Kafka.
ApsaraDB RDS for MySQL est activé. Pour plus d'informations, consultez Créer une instance ApsaraDB RDS for MySQL.
Vos services Realtime Compute for Apache Flink, ApsaraDB RDS for MySQL et ApsaraMQ for Kafka se trouvent dans le même VPC. S'ils résident dans des VPC différents, vous devez activer l'accès réseau inter-VPC ou utiliser des endpoints publics. Pour plus d'informations, consultez Comment accéder à d'autres services entre VPC ? et Comment accéder à Internet ?.
Si vous accédez aux ressources en tant qu'utilisateur RAM ou via un rôle RAM, vous devez disposer des permissions requises.
Préparation
Préparer la source de données MySQL
-
Créez une base de données ApsaraDB RDS for MySQL. Pour plus d'informations, consultez Créer une base de données.
Créez une base de données nommée
order_dwpour l'instance cible. -
Préparez la source de données MySQL CDC.
Sur la page des détails de l'instance, cliquez sur Log on to Database en haut de la page.
Dans la boîte de dialogue de connexion DMS qui s'affiche, saisissez le nom d'utilisateur et le mot de passe du compte de base de données créé, puis cliquez sur Login.
Après la connexion, double-cliquez sur la base de données
order_dwdans le volet de gauche pour changer de base de données.-
Dans la console SQL, saisissez les instructions DDL pour créer les trois tables métier ainsi que les instructions d'insertion de données.
CREATE TABLE `user` ( id bigint not null primary key, name varchar(50) not null ); CREATE TABLE `order` ( id bigint not null primary key, product varchar(50) not null, user_id bigint not null ); CREATE TABLE `feedback` ( id bigint not null primary key, user_id bigint not null, comment varchar(50) not null ); -- Prepare data INSERT INTO `user` VALUES(1, 'Tom'),(2, 'Jerry'); INSERT INTO `order` VALUES (1, 'Football', 2), (2, 'Basket', 1); INSERT INTO `feedback` VALUES (1, 1, 'Good.'), (2, 2, 'Very good');
Cliquez sur Execute, puis cliquez sur Execute.
Procédure
-
Créez et démarrez un job d'ingestion de données Flink CDC pour synchroniser en temps réel les données de la base de données MySQL en amont vers Kafka, afin qu'elles soient consommées par plusieurs jobs en aval. Le job de synchronisation de base de données crée automatiquement les topics. Vous pouvez définir les noms des topics à l'aide du module
route. Les topics utilisent les paramètres par défaut du cluster Kafka pour le nombre de partitions et de réplicas, et le paramètrecleanup.policyest défini surcompact.Noms de topics par défaut
Par défaut, les topics Kafka créés par le job de synchronisation de base de données suivent le format de nommage
{database_name}.{table_name}. Le job suivant crée trois topics :order_dw.user,order_dw.orderetorder_dw.feedback.-
Sur la page , créez un job d'ingestion de données Flink CDC et copiez le code suivant dans l'éditeur YAML.
source: type: mysql name: MySQL Source hostname: #{hostname} port: 3306 username: #{usernmae} password: #{password} tables: order_dw.\.* server-id: 28601-28604 # (Optional) Synchronize data from newly created tables during the incremental phase. scan.binlog.newly-added-table.enabled: true # (Optional) Synchronize table and field comments. include-comments.enabled: true # (Optional) Prioritize dispatching unbounded splits to avoid potential TaskManager OutOfMemory issues. scan.incremental.snapshot.unbounded-chunk-first.enabled: true # (Optional) Enable parsing filters to speed up reads. scan.only.deserialize.captured.tables.changelog.enabled: true sink: type: upsert-kafka name: upsert-kafka Sink properties.bootstrap.servers: xxxx.alikafka.aliyuncs.com:9092 # The following parameters are required for ApsaraMQ for Kafka. aliyun.kafka.accessKeyId: #{ak} aliyun.kafka.accessKeySecret: #{sk} aliyun.kafka.instanceId: #{instanceId} aliyun.kafka.endpoint: #{endpoint} aliyun.kafka.regionId: #{regionId} Dans le coin supérieur droit, cliquez sur Deploy pour déployer le job.
Dans la barre de navigation de gauche, cliquez sur . Dans la colonne Actions du job cible, cliquez sur Start, sélectionnez Initial Mode, puis cliquez sur Start.
Noms de topics par table
Le module
routepermet de spécifier un nom de topic pour chaque table. Le job suivant crée trois topics :user1,order2etfeedback3.-
Sur la page , créez un job d'ingestion de données Flink CDC et copiez le code suivant dans l'éditeur YAML.
source: type: mysql name: MySQL Source hostname: #{hostname} port: 3306 username: #{usernmae} password: #{password} tables: order_dw.\.* server-id: 28601-28604 # (Optional) Synchronize data from newly created tables during the incremental phase. scan.binlog.newly-added-table.enabled: true # (Optional) Synchronize table and field comments. include-comments.enabled: true # (Optional) Prioritize dispatching unbounded splits to avoid potential TaskManager OutOfMemory issues. scan.incremental.snapshot.unbounded-chunk-first.enabled: true # (Optional) Enable parsing filters to speed up reads. scan.only.deserialize.captured.tables.changelog.enabled: true route: - source-table: order_dw.user sink-table: user1 - source-table: order_dw.order sink-table: order2 - source-table: order_dw.feedback sink-table: feedback3 sink: type: upsert-kafka name: upsert-kafka Sink properties.bootstrap.servers: xxxx.alikafka.aliyuncs.com:9092 # The following parameters are required for ApsaraMQ for Kafka. aliyun.kafka.accessKeyId: #{ak} aliyun.kafka.accessKeySecret: #{sk} aliyun.kafka.instanceId: #{instanceId} aliyun.kafka.endpoint: #{endpoint} aliyun.kafka.regionId: #{regionId} Dans le coin supérieur droit, cliquez sur Deploy pour déployer le job.
Dans la barre de navigation de gauche, sélectionnez , cliquez sur Start dans la colonne Actions du job cible, sélectionnez Initial Mode, puis cliquez sur Start.
Noms de topics par lot
Le module
routepermet également de définir un modèle pour les noms des topics générés. Le job suivant crée trois topics :topic_user,topic_orderettopic_feedback.-
Sur la page , créez un job d'ingestion de données Flink CDC et copiez le code suivant dans l'éditeur YAML.
source: type: mysql name: MySQL Source hostname: #{hostname} port: 3306 username: #{usernmae} password: #{password} tables: order_dw.\.* server-id: 28601-28604 # (Optional) Synchronize data from newly created tables during the incremental phase. scan.binlog.newly-added-table.enabled: true # (Optional) Synchronize table and field comments. include-comments.enabled: true # (Optional) Prioritize dispatching unbounded splits to avoid potential TaskManager OutOfMemory issues. scan.incremental.snapshot.unbounded-chunk-first.enabled: true # (Optional) Enable parsing filters to speed up reads. scan.only.deserialize.captured.tables.changelog.enabled: true route: - source-table: order_dw.\.* sink-table: topic_<> replace-symbol: <> sink: type: upsert-kafka name: upsert-kafka Sink properties.bootstrap.servers: xxxx.alikafka.aliyuncs.com:9092 # The following parameters are required for ApsaraMQ for Kafka. aliyun.kafka.accessKeyId: #{ak} aliyun.kafka.accessKeySecret: #{sk} aliyun.kafka.instanceId: #{instanceId} aliyun.kafka.endpoint: #{endpoint} aliyun.kafka.regionId: #{regionId} Dans le coin supérieur droit, cliquez sur Deploy pour déployer le job.
Dans la barre de navigation de gauche, cliquez sur . Cliquez sur Start dans la colonne Actions du job cible, sélectionnez Initial Mode, puis cliquez sur Start.
-
-
Consommez les données Kafka en temps réel.
Le job d'ingestion écrit les données de la base de données MySQL en amont vers Kafka au format JSON. Plusieurs jobs en aval peuvent alors consommer les données d'un seul topic pour récupérer l'état le plus récent des tables de la base de données. Vous pouvez consommer les données des tables synchronisées vers Kafka selon l'une des méthodes suivantes :
Via un catalog
Lisez les données d'un topic Kafka en l'utilisant comme table source.
-
Sur la page , créez un job SQL en streaming et copiez le code suivant dans l'éditeur SQL.
CREATE TEMPORARY TABLE print_user_proudct( order_id BIGINT, product STRING, user_name STRING ) WITH ( 'connector'='print', 'logger'='true' ); CREATE TEMPORARY TABLE print_user_feedback( feedback_id BIGINT, `comment` STRING, user_name STRING ) WITH ( 'connector'='print', 'logger'='true' ); BEGIN STATEMENT SET; -- Required when writing to multiple sinks. -- Join order information with the user table in the Kafka JSON Catalog to display the username and product name for each order. INSERT INTO print_user_proudct SELECT `order`.key_id as order_id, value_product as product, `user`.value_name as user_name FROM `kafka-catalog`.`kafka`.`order`/*+OPTIONS('properties.group.id'='<yourGroupName>', 'scan.startup.mode'='earliest-offset')*/ as `order` -- Specify the group and startup mode. LEFT JOIN `kafka-catalog`.`kafka`.`user`/*+OPTIONS('properties.group.id'='<yourGroupName>', 'scan.startup.mode'='earliest-offset')*/ as `user` -- Specify the group and startup mode. ON `order`.value_user_id = `user`.key_id; -- Join reviews with the user table to display the content of each review and the corresponding username. INSERT INTO print_user_feedback SELECT feedback.key_id as feedback_id, value_comment as `comment`, `user`.value_name as user_name FROM `kafka-catalog`.`kafka`.feedback/*+OPTIONS('properties.group.id'='<yourGroupName>', 'scan.startup.mode'='earliest-offset')*/ as feedback -- Specify the group and startup mode. LEFT JOIN `kafka-catalog`.`kafka`.`user`/*+OPTIONS('properties.group.id'='<yourGroupName>', 'scan.startup.mode'='earliest-offset')*/ as `user` -- Specify the group and startup mode. ON feedback.value_user_id = `user`.key_id; END; -- Required when writing to multiple sinks.Cet exemple utilise le connecteur Print pour imprimer directement les résultats. Vous pouvez également diriger les résultats vers une table de résultats utilisant un autre connecteur pour une analyse plus poussée. Pour plus d'informations sur la syntaxe d'écriture vers plusieurs sinks, consultez Instruction INSERT INTO.
RemarqueLorsque vous utilisez directement cette méthode, des modifications de schéma peuvent survenir. Par conséquent, le schéma analysé par le catalog JSON Kafka peut différer du schéma de la table MySQL correspondante. Par exemple, des champs supprimés peuvent encore apparaître et certains champs peuvent présenter des valeurs nulles.
Le schéma lu depuis le catalog se compose des champs issus des données consommées. Si un champ est supprimé mais que ses messages n'ont pas expiré, le champ peut encore apparaître avec une valeur nulle. Aucun traitement particulier n'est nécessaire dans ce cas.
Dans le coin supérieur droit, cliquez sur Deploy pour déployer le job.
Dans la barre de navigation de gauche, cliquez sur , cliquez sur Start dans la colonne Actions du job cible, sélectionnez Initial Mode, puis cliquez sur Start.
Via une table temporaire
Définissez un schéma personnalisé et lisez les données depuis une table temporaire.
-
Sur la page , créez un job SQL en streaming et copiez le code suivant dans l'éditeur SQL.
CREATE TEMPORARY TABLE user_source ( key_id BIGINT, value_name STRING ) WITH ( 'connector' = 'kafka', 'topic' = 'user', 'properties.bootstrap.servers' = '<yourKafkaBrokers>', 'scan.startup.mode' = 'earliest-offset', 'key.format' = 'json', 'value.format' = 'json', 'key.fields' = 'key_id', 'key.fields-prefix' = 'key_', 'value.fields-prefix' = 'value_', 'value.fields-include' = 'EXCEPT_KEY', 'value.json.infer-schema.flatten-nested-columns.enable' = 'false', 'value.json.infer-schema.primitive-as-string' = 'false' ); CREATE TEMPORARY TABLE order_source ( key_id BIGINT, value_product STRING, value_user_id BIGINT ) WITH ( 'connector' = 'kafka', 'topic' = 'order', 'properties.bootstrap.servers' = '<yourKafkaBrokers>', 'scan.startup.mode' = 'earliest-offset', 'key.format' = 'json', 'value.format' = 'json', 'key.fields' = 'key_id', 'key.fields-prefix' = 'key_', 'value.fields-prefix' = 'value_', 'value.fields-include' = 'EXCEPT_KEY', 'value.json.infer-schema.flatten-nested-columns.enable' = 'false', 'value.json.infer-schema.primitive-as-string' = 'false' ); CREATE TEMPORARY TABLE feedback_source ( key_id BIGINT, value_user_id BIGINT, value_comment STRING ) WITH ( 'connector' = 'kafka', 'topic' = 'feedback', 'properties.bootstrap.servers' = '<yourKafkaBrokers>', 'scan.startup.mode' = 'earliest-offset', 'key.format' = 'json', 'value.format' = 'json', 'key.fields' = 'key_id', 'key.fields-prefix' = 'key_', 'value.fields-prefix' = 'value_', 'value.fields-include' = 'EXCEPT_KEY', 'value.json.infer-schema.flatten-nested-columns.enable' = 'false', 'value.json.infer-schema.primitive-as-string' = 'false' ); CREATE TEMPORARY TABLE print_user_proudct( order_id BIGINT, product STRING, user_name STRING ) WITH ( 'connector'='print', 'logger'='true' ); CREATE TEMPORARY TABLE print_user_feedback( feedback_id BIGINT, `comment` STRING, user_name STRING ) WITH ( 'connector'='print', 'logger'='true' ); BEGIN STATEMENT SET; -- Required when writing to multiple sinks. -- Join order information with the user table from the Kafka JSON Catalog to display the username and product name for each order. INSERT INTO print_user_proudct SELECT order_source.key_id as order_id, value_product as product, user_source.value_name as user_name FROM order_source LEFT JOIN user_source ON order_source.value_user_id = user_source.key_id; -- Join reviews with the user table to display the content of each review and the corresponding username. INSERT INTO print_user_feedback SELECT feedback_source.key_id as feedback_id, value_comment as `comment`, user_source.value_name as user_name FROM feedback_source LEFT JOIN user_source ON feedback_source.value_user_id = user_source.key_id; END; -- Required when writing to multiple sinks.Cet exemple utilise le connecteur Print pour imprimer directement les résultats. Vous pouvez également diriger les résultats vers une table de résultats d'un autre connecteur pour une analyse plus poussée. Pour plus d'informations sur la syntaxe d'écriture vers plusieurs sinks, consultez Instruction INSERT INTO.
Le tableau suivant décrit les paramètres de configuration de la table temporaire.
Paramètre
Description
Notes
connectorType de connecteur.
Définissez la valeur sur
kafka.topicNom du topic correspondant.
Doit correspondre à la description du catalog JSON Kafka.
properties.bootstrap.serversAdresses des brokers Kafka.
Le format est
host:port,host:port,host:port, séparé par des virgules (,).scan.startup.modePosition de départ pour la lecture des données depuis Kafka.
Valeurs valides :
-
earliest-offset: commence la lecture au premier offset disponible. -
latest-offset: commence la lecture au dernier offset. -
group-offsets(par défaut) : lit à partir de l'offset validé pour le groupe spécifié par properties.group.id. -
timestamp : lit à partir de l'horodatage spécifié par scan.startup.timestamp-millis.
-
specific-offsets : commence la lecture aux offsets spécifiés dans scan.startup.specific-offsets.
Note
Ce paramètre prend effet lorsque le job démarre sans état sauvegardé. Lorsque le job redémarre ou récupère à partir d'un checkpoint, il privilégie la lecture depuis l'état sauvegardé.
key.formatFormat utilisé par le connecteur Flink Kafka pour sérialiser ou désérialiser la clé du message Kafka.
Définissez la valeur sur
json.key.fieldsChamps de la table source ou de résultat correspondant à la clé du message Kafka.
Utilisez des points-virgules (;) pour séparer plusieurs noms de champs. Par exemple,
field1;field2.key.fields-prefixPréfixe personnalisé pour tous les champs de clé de message Kafka afin d'éviter les conflits de noms avec les champs de valeur de message ou les champs de métadonnées.
Cette valeur doit correspondre à celle du paramètre key.fields-prefix du catalog JSON Kafka.
value.formatFormat utilisé par le connecteur Flink Kafka pour sérialiser ou désérialiser la valeur du message Kafka.
Définissez la valeur sur
json.value.fields-prefixPréfixe personnalisé pour tous les champs de valeur de message Kafka afin d'éviter les conflits de noms avec les champs de clé de message ou les champs de métadonnées.
Doit correspondre à la valeur du paramètre value.fields-prefix du catalog JSON Kafka.
value.fields-includePolitique de gestion des champs de clé de message dans la valeur du message.
Définissez la valeur sur
EXCEPT_KEY. Cela indique que la valeur du message n'inclut pas les champs de la clé du message.value.json.infer-schema.flatten-nested-columns.enableIndique s'il faut développer récursivement les colonnes JSON imbriquées dans la valeur du message Kafka.
Valeur du paramètre infer-schema.flatten-nested-columns.enable du catalog correspondant.
value.json.infer-schema.primitive-as-stringIndique s'il faut déduire tous les types primitifs comme String dans la valeur du message Kafka.
Valeur du paramètre infer-schema.primitive-as-string du catalog correspondant.
-
Dans le coin supérieur droit, cliquez sur Deploy pour déployer le job.
Dans la barre de navigation de gauche, cliquez sur , cliquez sur Start dans la colonne Actions du job cible, sélectionnez Initial Mode, puis cliquez sur Start.
-
-
Consultez les résultats du job.
Dans la barre de navigation de gauche, cliquez sur , puis cliquez sur le job cible.
Sous l'onglet Logs, dans l'onglet Running Task Managers, cliquez sur la tâche dont vous souhaitez consulter le Path, ID.
-
Cliquez sur Logs et recherchez les informations de journal liées à
PrintSinkOutputWriter.Recherchez dans les journaux la sortie de
PrintSinkOutputWriter. La sortie contient quatre enregistrements de données jointes :+I[1, Good., Tom],+I[2, Very good, Jerry],+I[2, Basket, Tom]et+I[1, Football, Jerry]. Cela indique que les jointures entre la table utilisateur et les tables de commandes et de retours ont réussi.