MaxCompute propose un connecteur Flink Change Data Capture (CDC) qui synchronise en temps réel les modifications de données provenant de sources telles que MySQL vers des tables standard MaxCompute ou des tables Delta. Utilisez le connecteur de pipeline Flink CDC pour configurer un pipeline d'extraction, de transformation et de chargement (ETL) en continu de MySQL vers MaxCompute.
Le connecteur prend en charge :
La création automatique de tables basée sur la structure de la table source
L'évolution du schéma : les modifications apportées au schéma de la table source sont propagées à MaxCompute
La synchronisation complète de la base de données sur plusieurs tables
Le routage des tables : mappez les tables sources vers différents noms de tables cibles, y compris la fusion de bases de données fragmentées
Prérequis
Avant de commencer, assurez-vous de disposer des éléments suivants :
Un projet MaxCompute. Pour en créer un, consultez la rubrique Créer un projet MaxCompute.
Docker et Docker Compose installés sur votre machine
Un environnement d'exécution Java (JRE) compatible avec Flink 1.18.0
Notes d'utilisation
Sélection du type de table : si la table source possède une clé primaire, une table Delta est créée automatiquement. En l'absence de clé primaire, une table standard MaxCompute est créée.
Comportement d'écriture pour les tables standard : les opérations
DELETEsont ignorées. Les opérationsUPDATEsont traitées comme des opérationsINSERT.Sémantique de livraison : seule la sémantique « au moins une fois » (at-least-once) est prise en charge. Pour mettre en œuvre des écritures idempotentes, utilisez une table Delta ; la clé primaire de la table Delta gère la déduplication.
Évolution du schéma : les modifications de schéma de la table source sont propagées à MaxCompute. Les nouvelles colonnes sont toujours ajoutées en dernière position. Les types de données des colonnes ne peuvent être modifiés que vers un type compatible. Pour plus de détails, consultez la rubrique Modifier le type de données d'une colonne.
Fusion de bases de données fragmentées : vous pouvez fusionner plusieurs tables sources en une seule table de destination à l'aide d'expressions régulières. La fusion de tables présentant des valeurs de clé primaire en double entre les sources n'est pas encore prise en charge.
Premiers pas
Ce tutoriel configure un travail ETL en continu complet qui synchronise les modifications de données de MySQL vers MaxCompute, couvrant la synchronisation complète de la base de données, l'évolution du schéma et le routage des tables.
Configurer les environnements
Configurer un cluster Flink en mode autonome
-
Téléchargez le fichier flink-1.18.0-bin-scala_2.12.tgz et décompressez-le pour obtenir le répertoire
flink-1.18.0. Dans le répertoireflink-1.18.0, exécutez la commande suivante pour définir la variableFLINK_HOME:export FLINK_HOME=$(pwd) -
Ouvrez le fichier
$FLINK_HOME/conf/flink-conf.yamlavec un éditeur de texte et ajoutez les paramètres suivants :# Enable checkpointing. Run a checkpoint every 3 seconds. # For production, set the checkpoint interval to at least 30 seconds. execution.checkpointing.interval: 3000 # The Flink CDC pipeline connector relies on Flink's communication mechanism # for data synchronization. Increase the timeout to prevent connection drops. pekko.ask.timeout: 60s -
Démarrez le cluster Flink :
./bin/start-cluster.shAprès le démarrage, ouvrez
http://localhost:8081/dans un navigateur pour accéder à l'interface web Flink. Exécutezstart-cluster.shplusieurs fois pour ajouter des TaskManagers.
Configurer un environnement MySQL
Ce tutoriel utilise Docker Compose pour exécuter MySQL.
-
Créez un fichier nommé
docker-compose.yamlavec le contenu suivant :Paramètre
description
versionVersion de Docker Compose
imageVersion de l'image Docker. Définissez-la sur
debezium/example-mysql:1.1portsMappage du port MySQL
environmentMot de passe root MySQL et identifiants utilisateur
version: '2.1' services: mysql: image: debezium/example-mysql:1.1 ports: - "3306:3306" environment: - MYSQL_ROOT_PASSWORD=123456 - MYSQL_USER=mysqluser - MYSQL_PASSWORD=mysqlpw -
Dans le répertoire contenant le fichier
docker-compose.yaml, démarrez les conteneurs :docker-compose up -dExécutez
docker pspour confirmer que tous les conteneurs sont en cours d'exécution.
Préparer les données dans MySQL
-
Connectez-vous au conteneur MySQL :
docker-compose exec mysql mysql -uroot -p123456 -
Créez la base de données
app_dbet renseignez trois tables :-
Créez la base de données :
CREATE DATABASE app_db; USE app_db; -
Créez et renseignez la table
orders:CREATE TABLE `orders` ( `id` INT NOT NULL, `price` DECIMAL(10,2) NOT NULL, PRIMARY KEY (`id`) ); INSERT INTO `orders` (`id`, `price`) VALUES (1, 4.00); INSERT INTO `orders` (`id`, `price`) VALUES (2, 100.00); -
Créez et renseignez la table
shipments:CREATE TABLE `shipments` ( `id` INT NOT NULL, `city` VARCHAR(255) NOT NULL, PRIMARY KEY (`id`) ); INSERT INTO `shipments` (`id`, `city`) VALUES (1, 'beijing'); INSERT INTO `shipments` (`id`, `city`) VALUES (2, 'xian'); -
Créez et renseignez la table
products:CREATE TABLE `products` ( `id` INT NOT NULL, `product` VARCHAR(255) NOT NULL, PRIMARY KEY (`id`) ); INSERT INTO `products` (`id`, `product`) VALUES (1, 'Beer'); INSERT INTO `products` (`id`, `product`) VALUES (2, 'Cap'); INSERT INTO `products` (`id`, `product`) VALUES (3, 'Peanut');
-
Soumettre un travail de pipeline YAML
-
Téléchargez les packages JAR requis :
Package Flink CDC : téléchargez le fichier flink-cdc-3.1.1-bin.tar.gz et décompressez-le pour obtenir le répertoire
flink-cdc-3.1.1. Déplacez le contenu des répertoiresbin,lib,logetconfvers les répertoires correspondants sousflink-1.18.0.-
Packages de connecteurs : téléchargez les fichiers JAR suivants et placez-les dans le répertoire
flink-1.18.0/lib: > Remarque : les liens de téléchargement ci-dessus pointent uniquement vers des versions publiées. Pour utiliser une version SNAPSHOT, compilez le code source à partir de la branche master ou release localement. Pilote JDBC : téléchargez le fichier mysql-connector-java-8.0.27.jar. Transmettez-le à l'interface CLI Flink CDC avec l'option
--jar, ou placez-le dans$FLINK_HOME/libet redémarrez le cluster. Les connecteurs CDC n'incluent plus les pilotes JDBC.
-
Créez le fichier YAML du pipeline. L'exemple suivant (
mysql-to-maxcompute.yaml) synchronise toutes les tables de la base de donnéesapp_db:Espace réservé
description
${your_accessId}Votre ID AccessKey. Obtenez-le depuis la page Paire de clés d'accès.
${your_accessKey}Le secret AccessKey correspondant à votre ID AccessKey.
${your_maxcompute_endpoint}L'endpoint MaxCompute pour votre région et votre réseau. Consultez la rubrique Endpoints.
${your_project}Le nom de votre projet MaxCompute. Retrouvez-le dans la console MaxCompute sous Workspace > Projects.
################################################################################ # Description: Sync MySQL all tables to MaxCompute ################################################################################ source: type: mysql hostname: localhost port: 3306 username: root password: 123456 tables: app_db\.* server-id: 5400-5404 server-time-zone: UTC # Configure the accessId, accessKey, endpoint, and project parameters. sink: type: maxcompute name: MaxComputeSink accessId: ${your_accessId} accessKey: ${your_accessKey} endpoint: ${your_maxcompute_endpoint} project: ${your_project} bucketsNum: 8 pipeline: name: Sync MySQL Database to MaxCompute parallelism: 1Remplacez les espaces réservés suivants : Pour connaître tous les paramètres source disponibles, consultez la rubrique Connecteur MySQL. Pour connaître tous les paramètres de destination disponibles, consultez la section Configuration du connecteur.
-
Soumettez le travail de pipeline au cluster Flink autonome :
./bin/flink-cdc.sh mysql-to-maxcompute.yamlUne soumission réussie renvoie une sortie similaire à :
Pipeline has been submitted to cluster. Job ID: f9f9689866946e25bf151ecc179ef46f Job Description: Sync MySQL Database to MaxComputeL'interface web Flink affiche un travail en cours d'exécution nommé
Sync MySQL Database to MaxCompute. -
Dans MaxCompute, vérifiez que les tables sont créées et que les données sont écrites :
-- Query the orders table. read orders; -- Expected result: +------------+------------+ | id | price | +------------+------------+ | 1 | 4 | | 2 | 100 | +------------+------------+ -- Query the shipments table. read shipments; -- Expected result: +------------+------------+ | id | city | +------------+------------+ | 1 | beijing | | 2 | xian | +------------+------------+ -- Query the products table. read products; -- Expected result: +------------+------------+ | id | product | +------------+------------+ | 3 | Peanut | | 1 | Beer | | 2 | Cap | +------------+------------+
Synchroniser les modifications en temps réel
Lorsque le pipeline est en cours d'exécution, les modifications effectuées dans MySQL sont reflétées dans MaxCompute en temps réel. Les exemples suivants utilisent la table orders pour illustrer les opérations d'insertion, de modification du schéma, de mise à jour et de suppression.
-
Connectez-vous au conteneur MySQL :
docker-compose exec mysql mysql -uroot -p123456 -
Insérez une ligne dans la table
orders:INSERT INTO app_db.orders (id, price) VALUES (3, 100.00);Exécutez
read orders;dans MaxCompute. La nouvelle ligne apparaît :+------------+------------+ | id | price | +------------+------------+ | 3 | 100 | | 1 | 4 | | 2 | 100 | +------------+------------+ -
Ajoutez une colonne à la table
ordersdans MySQL :ALTER TABLE app_db.orders ADD amount varchar(100) NULL;Exécutez
read orders;dans MaxCompute. La nouvelle colonneamountest propagée automatiquement :+------------+------------+------------+ | id | price | amount | +------------+------------+------------+ | 3 | 100 | NULL | | 1 | 4 | NULL | | 2 | 100 | NULL | +------------+------------+------------+ -
Mettez à jour une ligne dans la table
orders:UPDATE app_db.orders SET price=100.00, amount=100.00 WHERE id=1;Exécutez
read orders;dans MaxCompute :+------------+------------+------------+ | id | price | amount | +------------+------------+------------+ | 3 | 100 | NULL | | 1 | 100 | 100.00 | | 2 | 100 | NULL | +------------+------------+------------+ -
Supprimez une ligne de la table
orders:DELETE FROM app_db.orders WHERE id=2;Exécutez
read orders;dans MaxCompute :+------------+------------+------------+ | id | price | amount | +------------+------------+------------+ | 3 | 100 | NULL | | 1 | 100 | 100.00 | +------------+------------+------------+
Router les tables
Flink CDC permet de router les schémas et les données des tables sources vers différents noms de tables de destination. Utilisez cette fonctionnalité pour les migrations de base de données, les renommages de tables ou les fusions de bases de données fragmentées.
L'exemple YAML suivant route chaque table source vers une table cible portant un nom différent :
################################################################################
# Description: Sync MySQL all tables to MaxCompute
################################################################################
source:
type: mysql
hostname: localhost
port: 3306
username: root
password: 123456
tables: app_db\.*
server-id: 5400-5404
server-time-zone: UTC
# Configure the accessId, accessKey, endpoint, and project parameters.
sink:
type: maxcompute
name: MaxComputeSink
accessId: ${your_accessId}
accessKey: ${your_accessKey}
endpoint: ${your_maxcompute_endpoint}
project: ${your_project}
bucketsNum: 8
route:
- source-table: app_db.orders
sink-table: ods_db.ods_orders
- source-table: app_db.shipments
sink-table: ods_db.ods_shipments
- source-table: app_db.products
sink-table: ods_db.ods_products
pipeline:
name: Sync MySQL Database to MaxCompute
parallelism: 1
Le champ source-table prend en charge les expressions régulières. Pour fusionner plusieurs tables fragmentées en une seule table cible :
route:
- source-table: app_db.order\.*
sink-table: ods_db.ods_orders
Cette configuration agrège les données de app_db.order01, app_db.order02, app_db.order03 et des tables similaires dans ods_db.ods_orders.
Les scénarios dans lesquels des valeurs de clé primaire en double existent entre plusieurs tables sources ne sont pas pris en charge et le seront dans les prochaines versions du connecteur.
Pour connaître tous les paramètres de routage, consultez la rubrique Route.
Nettoyer l'environnement
Une fois le tutoriel terminé, arrêtez tous les composants en cours d'exécution.
-
Arrêtez les conteneurs Docker. Exécutez la commande suivante dans le répertoire contenant le fichier
docker-compose.yaml:docker-compose down -
Arrêtez le cluster Flink. Exécutez la commande suivante dans le répertoire
flink-1.18.0:./bin/stop-cluster.sh
Annexe
Configuration du connecteur
Le connecteur de destination MaxCompute prend en charge les paramètres de configuration suivants.
Les données sont vidées vers MaxCompute lorsque le tampon d'une partition ou d'un compartiment dépasse le seuil de taille configuré, ou lorsqu'un point de contrôle (checkpoint) est déclenché.
|
Paramètre |
Obligatoire |
Valeur par défaut |
Type |
description |
|
|
Oui |
— |
String |
Le type de connecteur. Définissez-le sur |
|
|
Non |
— |
String |
Le nom d'affichage de la destination. |
|
|
Oui |
— |
String |
Votre ID AccessKey. Obtenez-le depuis la page Paire de clés d'accès. |
|
|
Oui |
— |
String |
Le secret AccessKey correspondant à votre ID AccessKey. |
|
|
Oui |
— |
String |
L'endpoint MaxCompute. Configurez-le en fonction de votre région et de votre méthode de connexion réseau. Consultez la rubrique Endpoints. |
|
|
Oui |
— |
String |
Le nom du projet MaxCompute. Retrouvez-le dans la console MaxCompute sous Workspace > Projects. |
|
|
Non |
— |
String |
L'endpoint MaxCompute Tunnel. Le routage automatique est pris en charge dans la plupart des cas. Configurez ce paramètre uniquement lors de l'utilisation d'un environnement réseau spécial, tel qu'un proxy. |
|
|
Non |
— |
String |
Le nom du groupe de ressources exclusif pour MaxCompute Tunnel. Si ce paramètre n'est pas défini, un groupe de ressources partagé est utilisé. |
|
|
Non |
— |
String |
Requis lors de l'utilisation d'un jeton STS émis par un rôle RAM pour l'authentification. |
|
|
Non |
16 |
Integer |
Le nombre de compartiments pour la création automatique de table Delta. Consultez la rubrique Présentation de l'entrepôt de données quasi temps réel. |
|
|
Non |
|
String |
L'algorithme de compression pour les écritures. Valeurs valides : |
|
|
Non |
|
String |
La taille maximale du tampon par partition ou table non partitionnée. Lorsque ce seuil est dépassé, les données sont vidées vers MaxCompute. Chaque partition ou table non partitionnée conserve un tampon indépendant. |
|
|
Non |
|
String |
La taille maximale du tampon par compartiment. S'applique uniquement aux écritures de table Delta. Lorsque ce seuil est dépassé, les données sont vidées vers MaxCompute. Chaque compartiment conserve un tampon indépendant. |
|
|
Non |
16 |
Integer |
Le nombre maximal de partitions ou de tables traitées simultanément lors d'un point de contrôle. |
|
|
Non |
4 |
Integer |
Le nombre maximal de compartiments écrits simultanément dans MaxCompute. S'applique uniquement aux écritures de table Delta. |
|
|
Non |
3 |
Integer |
Le nombre maximal de tentatives en cas d'erreurs réseau. |
|
|
Non |
— |
Long |
La durée d'attente entre les tentatives en cas d'erreurs réseau. Unité : millisecondes. |
Mappages d'emplacement de table
Lorsque le connecteur crée automatiquement une table, il mappe les informations d'emplacement de la source vers MaxCompute comme suit.
Si le projet MaxCompute ne prend pas en charge les modèles de schéma, chaque tâche de synchronisation ne peut synchroniser qu'une seule base de données MySQL. Le connecteur ignore les informations TableId.namespace pour les autres sources de données.
|
Objet dans Flink CDC |
Emplacement dans MaxCompute |
Emplacement dans MySQL |
|
Projet (issu du fichier de configuration) |
Projet |
— |
|
|
Schéma (nécessite la prise en charge du modèle de schéma ; ignoré si le projet ne prend pas en charge les modèles de schéma) |
Base de données |
|
|
Table |
Table |
Mappages de types de données
Le tableau suivant montre comment les types de données Flink sont mappés vers les types de données MaxCompute.
|
Type de données Flink |
Type de données MaxCompute |
notes |
|
CHAR/VARCHAR |
STRING |
MaxCompute stocke tous les types de chaîne sous forme de STRING. |
|
BOOLEAN |
BOOLEAN |
— |
|
BINARY/VARBINARY |
BINARY |
— |
|
DECIMAL |
DECIMAL |
— |
|
TINYINT |
TINYINT |
— |
|
SMALLINT |
SMALLINT |
— |
|
INTEGER |
INTEGER |
— |
|
BIGINT |
BIGINT |
— |
|
FLOAT |
FLOAT |
— |
|
DOUBLE |
DOUBLE |
— |
|
TIME_WITHOUT_TIME_ZONE |
STRING |
MaxCompute ne dispose pas de type TIME natif. Les valeurs temporelles sont stockées sous forme de chaînes. |
|
DATE |
DATE |
— |
|
TIMESTAMP_WITHOUT_TIME_ZONE |
TIMESTAMP_NTZ |
Stocké sans information de fuseau horaire. |
|
TIMESTAMP_WITH_LOCAL_TIME_ZONE |
TIMESTAMP |
— |
|
TIMESTAMP_WITH_TIME_ZONE |
TIMESTAMP |
— |
|
ARRAY |
ARRAY |
— |
|
MAP |
MAP |
— |
|
ROW |
STRUCT |
— |