Realtime Compute for Apache Flink prend en charge les instructions CTAS (CREATE TABLE AS) pour la synchronisation de tables uniques et CDAS (CREATE DATABASE AS SELECT) pour la synchronisation multi-tables ou de bases de données complètes. Ces mécanismes permettent une réplication en temps réel depuis une instance ApsaraDB RDS for MySQL vers un cluster E-MapReduce (EMR) StarRocks.
Contexte
L'instruction CTAS (CREATE TABLE AS) crée automatiquement une table StarRocks dont le schéma est identique à celui de la table source MySQL. Elle synchronise également les données et réplique les modifications de schéma de la source vers la destination en temps réel, ce qui simplifie la création des tables cibles et garantit la cohérence des schémas.
Lors de l'exécution d'une instruction CTAS, Flink effectue les étapes suivantes :
-
Vérifie si la table cible existe.
Si la table n'existe pas, Flink utilise le catalog de destination pour la créer. La nouvelle table cible reprend alors exactement le même schéma que la table source.
Si la table existe déjà, Flink ignore l'étape de création. Une erreur est signalée si le schéma de la table cible diffère de celui de la table source.
Soumet et démarre le déploiement de synchronisation des données. Flink assure alors la réplication des données et des changements de schéma de la table source vers la table cible.
Une instruction CTAS synchronise les données en temps réel et propage les modifications de schéma de la source vers la table de destination.
Les modifications de schéma englobent la création initiale de la table ainsi que ses altérations ultérieures.
-
Modifications de schéma prises en charge :
Ajout d'une colonne nullable : Flink ajoute automatiquement la nouvelle colonne à la fin du schéma de la table cible et synchronise ses données.
Suppression d'une colonne nullable : La colonne n'est pas physiquement supprimée de la table cible. Flink remplit automatiquement ses données avec des valeurs
NULL.-
Renommage d'une colonne : Flink traite cette opération comme la combinaison d'un ajout de nouvelle colonne et de la suppression de l'ancienne. La colonne renommée est ajoutée à la fin de la table cible, tandis que les données de la colonne originale sont remplacées par des valeurs
NULL.Par exemple, si vous renommez col_a en col_b, la colonne col_b est ajoutée à la fin de la table cible et les données de col_a sont automatiquement remplacées par des valeurs NULL.
-
Modifications de schéma non prises en charge :
-
Changements de types de données.
Cela inclut, par exemple, le passage d'un type VARCHAR à BIGINT ou la modification d'une propriété NOT NULL en NULLABLE.
Modifications des contraintes, telles qu'une clé primaire ou un index.
Ajout ou suppression de colonnes non nullables.
Ajustements de la longueur des champs dans une instruction DDL.
-
En cas de modification de schéma non prise en charge, supprimez manuellement la table cible et redémarrez le déploiement CTAS. Cette opération recrée la table cible et resynchronise l'intégralité des données historiques.
CTAS n'identifie pas les types spécifiques de DDL. Il compare plutôt les différences de schéma entre les enregistrements de données avant et après une modification. Si vous supprimez une colonne puis la rajoutez sans aucune modification de données entre les deux opérations DDL, CTAS ne détecte aucun changement de schéma. De même, lors de l'ajout d'une colonne, CTAS ne détecte le changement de schéma et ne le synchronise vers la table cible qu'après qu'une modification de données s'y soit produite.
Pour consulter les types de données pris en charge lors de la création d'une table avec
CTAS, reportez-vous à Data type mapping between Flink and StarRocks.Lorsque vous utilisez une instruction CTAS pour fusionner plusieurs tables MySQL, Flink ajoute automatiquement deux colonnes,
_db_nameet_table_name, au début du schéma de la table générée afin d'identifier la table source. Ce comportement ne peut pas être modifié. Par conséquent, lorsque vous définissez l'ordre des colonnes pour la nouvelle table, commencez à partir de la troisième colonne pour garantir que le schéma résultant corresponde à vos attentes.
Prérequis
Activez Realtime Compute for Apache Flink (Fully Managed) et créez un cluster Flink. Pour plus d'informations, consultez Activate Realtime Compute for Apache Flink (Fully Managed) et Getting started with a Flink SQL deployment.
Créez un cluster StarRocks. Pour plus d'informations, consultez Create a StarRocks cluster.
Créez une instance ApsaraDB RDS for MySQL. Pour plus d'informations, consultez Create an ApsaraDB RDS for MySQL instance.
Les exemples de cette rubrique utilisent MySQL 5.7, un cluster E-MapReduce (EMR) StarRocks (EMR-3.39.1) et Realtime Compute for Apache Flink (version 1.15-vvr-6.0.3).
Limites
Le cluster Flink, le cluster StarRocks et l'instance ApsaraDB RDS for MySQL doivent se trouver dans le même VPC.
L'instance ApsaraDB RDS for MySQL doit être en version 5.7 ou ultérieure.
L'accès Internet doit être activé sur le cluster StarRocks.
La version de Flink utilisée dans votre cluster doit être 1.15-vvr-6.0.3 ou ultérieure.
Étape 1 : Préparer les données de test
-
Créez une base de données et un compte de test. Pour plus d'informations, consultez Create databases and accounts for an ApsaraDB RDS for MySQL instance.
Après avoir créé la base de données et le compte, accordez les autorisations de lecture et d'écriture au compte de test.
RemarqueDans cette rubrique, la base de données se nomme
test_cdcet le comptetest. Connectez-vous à l'instance MySQL à l'aide du compte de test. Pour plus d'informations, consultez Use DMS to log on to an ApsaraDB RDS for MySQL instance.
-
Exécutez les commandes suivantes dans MySQL pour créer une table de données.
use test_cdc; CREATE TABLE IF NOT EXISTS `runoob_tbl`( `runoob_id` INT UNSIGNED AUTO_INCREMENT, `runoob_title` VARCHAR(100) NOT NULL, `runoob_author` VARCHAR(40) NOT NULL, `submission_date` DATE, `add_col` int DEFAULT NULL, PRIMARY KEY ( `runoob_id` ) )ENGINE=InnoDB DEFAULT CHARSET=utf8; INSERT INTO test_cdc.`runoob_tbl` (`runoob_id`,`runoob_title`,`runoob_author`,`submission_date`,`add_col`) values (18,'first','tom','2022-06-22 17:13:44',3) Connectez-vous au cluster StarRocks via SSH. Pour plus d'informations, consultez Log on to a cluster.
-
Exécutez la commande suivante pour vous connecter au cluster StarRocks.
mysql -h127.0.0.1 -P 9030 -uroot -
Exécutez les commandes suivantes pour créer un utilisateur et lui attribuer des autorisations.
CREATE DATABASE test_cdc; CREATE USER 'test' IDENTIFIED by '123456'; GRANT CREATE TABLE ON DATABASE test_cdc TO test;
Étape 2 : Créer des catalogs
Sur la page Draft Editor de la console Realtime Compute for Apache Flink, créez des catalogs pour MySQL et StarRocks. Pour plus d'informations, consultez Getting started with a Flink SQL deployment.
Les paramètres fournis à titre d'exemple sont donnés uniquement pour référence. Configurez-les selon vos besoins spécifiques.
-
Catalog MySQL
-
Exemple
CREATE CATALOG mysql WITH ( 'type' = 'mysql', 'hostname' = 'rm-2zepd6e20u3od****.mysql.rds.aliyuncs.com', 'port' = '3306', 'username' = 'emr-test', 'password' = '123456', 'default-database' = 'test_cdc' ); -
Paramètres
Parameter
Description
type
Type du catalog. Définissez la valeur sur
mysql.hostname
Endpoint interne de l'instance ApsaraDB RDS for MySQL. Vous pouvez copier cet endpoint depuis la page Database Connection de l'instance dans la console ApsaraDB RDS. Exemple :
rm-2zepd6e20u3od****.mysql.rds.aliyuncs.com.port
Numéro de port du service de base de données MySQL. La valeur par défaut est
3306.username
Nom d'utilisateur pour accéder au service de base de données MySQL.
Utilisez le nom d'utilisateur du compte créé lors de l'Étape 1 : Préparer les données de test. Cet exemple utilise test.
password
Mot de passe pour accéder au service de base de données MySQL.
Utilisez le mot de passe du compte créé lors de l'Étape 1 : Préparer les données de test.
default-database
Nom de la base de données MySQL par défaut.
Utilisez le nom de la base de données indiqué dans l'Étape 1 : Préparer les données de test. Cet exemple utilise test_cdc.
-
-
Catalog StarRocks
-
Exemple
CREATE CATALOG sr WITH ( 'type' = 'starrocks', 'endpoint' = '172.16.**.**:9030', 'username' = 'test', 'password' = '123456', 'dbname' = 'test_cdc' ); -
Paramètres
Parameter
Description
type
Type du catalog. Définissez la valeur sur
starrocks.endpoint
Adresse IP et port du Frontend (FE) StarRocks.
username
Nom d'utilisateur pour accéder au cluster StarRocks.
Utilisez le nom d'utilisateur du compte créé lors de l'Étape 1 : Préparer les données de test. Cet exemple utilise test.
password
Mot de passe pour le service de base de données StarRocks.
Utilisez le mot de passe du compte créé lors de l'Étape 1 : Préparer les données de test.
dbname
Nom de la base de données StarRocks.
Utilisez le nom de la base de données indiqué dans l'Étape 1 : Préparer les données de test. Cet exemple utilise test_cdc.
-
Étape 3 : Créer et publier un déploiement
-
Sur la page Draft Editor de la console Realtime Compute for Apache Flink, rédigez une instruction
CTAS.Voici trois exemples d'instructions CTAS.
-
Sémantique At-least-once : Utilisez l'option sink.buffer-flush.interval-ms pour configurer l'intervalle d'écriture des données dans StarRocks. Cette option permet de réduire la latence et l'utilisation de la mémoire.
/* At-least-once semantics */ use CATALOG sr; CREATE TABLE IF NOT EXISTS runoob_tbl_sr with ( 'starrocks.create.table.properties'=' engine = olap primary key(runoob_id) distributed by hash(runoob_id ) buckets 8', 'database-name'='test_cdc', 'jdbc-url'='jdbc:mysql://172.16.**.**:9030', 'load-url'='172.16.**.**:18030', 'table-name'='runoob_tbl_sr', 'username'='test', 'password' = '123456', 'sink.buffer-flush.interval-ms' = '5000', 'sink.properties.row_delimiter' = '\x02', 'sink.properties.column_separator' = '\x01' ) as table mysql.test_cdc.runoob_tbl /*+ OPTIONS ( 'connector' = 'mysql-cdc', 'hostname' = 'rm-2zepd6e20u3od****.mysql.rds.aliyuncs.com', 'port' = '3306', 'username' = 'test', 'password' = '123456', 'database-name' = 'test_cdc', 'table-name' = 'runoob_tbl' )*/; -
Sémantique Exactly-once : Vous devez configurer un intervalle de checkpoint. Cela empêche la perte et la duplication des données en cas de défaillance, mais la visibilité des données dépend de l'intervalle de checkpoint. Pour plus d'informations, consultez Checkpointing.
/* Exactly-once semantics. */ set 'execution.checkpointing.interval' = '1 min'; set 'execution.checkpointing.mode' = 'EXACTLY_ONCE'; set 'execution.checkpointing.timeout' = '10 min'; use CATALOG sr; CREATE TABLE IF NOT EXISTS runoob_tbl with ( 'starrocks.create.table.properties'=' engine = olap primary key(runoob_id) distributed by hash(runoob_id ) buckets 8', 'database-name'='test_cdc', 'jdbc-url'='jdbc:mysql://172.16.**.**:9030', 'load-url'='172.16.**.**:18030', 'table-name'='runoob_tbl', 'username'='test', 'password' = '123456', 'sink.semantic' = 'exactly-once', 'sink.properties.row_delimiter' = '\x02', 'sink.properties.column_separator' = '\x01 ) as table mysql.test_cdc.runoob_tbl /*+ OPTIONS ( 'connector' = 'mysql-cdc', 'hostname' = 'rm-2zepd6e20u3od****.mysql.rds.aliyuncs.com', 'port' = '3306', 'username' = 'test', 'password' = '123456', 'database-name' = 'test_cdc', 'table-name' = 'runoob_tbl' )*/; -
Mode simple : La définition des champs n'est pas requise lors de la création de la table, car le schéma est copié directement depuis MySQL. Toutefois, ce mode ne permet pas de créer de partitions. Pour utiliser des partitions, optez pour le mode normal.
/* The two preceding examples use normal mode. This example demonstrates simple mode. */ use CATALOG sr; CREATE TABLE IF NOT EXISTS runoob_tbl1 with ( 'starrocks.create.table.properties'='buckets 8', 'starrocks.create.table.mode'='simple', 'database-name'='test_cdc', 'jdbc-url'='jdbc:mysql://172.16.**.**:9030', 'load-url'='172.16.**.**:18030', 'table-name'='runoob_tbl_sr', 'username'='test', 'password' = '123456', 'sink.buffer-flush.interval-ms' = '5000', 'sink.properties.row_delimiter' = '\x02', 'sink.properties.column_separator' = '\x01' ) as table mysql.test_cdc.runoob_tbl /*+ OPTIONS ( 'connector' = 'mysql-cdc', 'hostname' = 'rm-2zepd6e20u3od****.mysql.rds.aliyuncs.com', 'port' = '3306', 'username' = 'emr-test', 'password' = '123456', 'database-name' = 'test_cdc', 'table-name' = 'runoob_tbl' )*/;
Tableau 1. Paramètres WITH
Parameter
Required
Description
starrocks.create.table.properties
Yes
Définitions suffixées de l'instruction StarRocks
CREATE TABLE, à l'exclusion des définitions de colonnes. Cela inclut notammentengine,keyetbuckets.database-name
Yes
Nom de la base de données StarRocks.
Cet exemple utilise
test_cdc.jdbc-url
Yes
Permet d'exécuter des opérations de requête dans StarRocks.
Par exemple,
jdbc:mysql://172.16..:9030. La partie172.16..correspond à l'adresse IP interne du cluster StarRocks.load-url
Yes
Adresse IP et port HTTP du frontend (FE) StarRocks au format
Internal IP address of StarRocks cluster:Port. Cette rubrique prend le port 8030 comme exemple. Sélectionnez un port en fonction de la version de votre cluster :18030 : Pour EMR V5.9.0 ou ultérieur et EMR V3.43.0 ou ultérieur.
8030 : Pour EMR V5.8.0 ou antérieur et EMR V3.42.0 ou antérieur.
RemarquePour plus d'informations sur les ports, consultez UI and ports.
sink.semantic
No
Définissez cette valeur sur
exactly-oncepour garantir la sémantique de cohérence des données. La valeur par défaut estat-least-once.starrocks.create.table.mode
No
Valeurs prises en charge :
-
normal(par défaut) : Vous devez fournir des configurations complètes, telles queengine,keyetbuckets, dans l'option starrocks.create.table.properties. -
simple: Le moteur est défini surolapet le type de clé surprimary keypar défaut. La clé primaire est héritée de la table MySQL. Par défaut, la table est distribuée par hachage de toutes les colonnes de clé primaire, sans partitions. Vous devez spécifierbucketsdans l'option starrocks.create.table.properties. D'autres configurations commepropertiessont facultatives.
RemarqueLe paramètre
sink.use.new-apia été supprimé dans les versions Flink 1.15-vvr-6.0.5 et ultérieures. Si vous utilisez une version antérieure à 1.15-vvr-6.0.5, ajoutez'sink.use.new-api'='false',aux paramètres WITH.Pour plus d'informations sur les autres configurations, consultez Continuously load data from Apache Flink.
Tableau 2. Paramètres OPTIONS
Parameter
Description
connector
Type de connecteur. Définissez la valeur sur
mysql-cdc.hostname
Endpoint interne de l'instance ApsaraDB RDS for MySQL.
Vous pouvez copier cet endpoint depuis la page Database Connection de l'instance dans la console ApsaraDB RDS. Par exemple, rm-bp1nu0c46fn9k****.mysql.rds.aliyuncs.com.
port
Numéro de port du service de base de données MySQL. La valeur par défaut est
3306.username
Nom d'utilisateur pour accéder au service de base de données MySQL.
Utilisez le nom d'utilisateur du compte créé lors de l'Étape 1 : Préparer les données de test. Cet exemple utilise test.
password
Mot de passe pour accéder au service de base de données MySQL.
Utilisez le mot de passe du compte créé lors de l'Étape 1 : Préparer les données de test.
table-name
Nom de la table dans StarRocks.
Utilisez le nom de la table indiqué dans l'Étape 1 : Préparer les données de test. Cet exemple utilise runoob_tbl.
database-name
Nom de la base de données MySQL par défaut.
Utilisez le nom de la base de données indiqué dans l'Étape 1 : Préparer les données de test. Cet exemple utilise test_cdc.
-
Dans les Advanced settings de la page Draft Editor, sélectionnez la version Flink 1.15-vvr-6.0.3 ou ultérieure.
Cliquez sur online.
Sur la page Deployments, localisez le déploiement cible et cliquez sur START dans la colonne Actions.
Étape 4 : Démonstration de scénarios
Interrogation des données
Connectez-vous au cluster StarRocks via SSH. Pour plus d'informations, consultez Log on to a cluster.
-
Exécutez la commande suivante pour vous connecter au cluster StarRocks.
mysql -h127.0.0.1 -P 9030 -uroot -
Dans le CLI StarRocks, exécutez les commandes suivantes pour consulter les données de la table.
use test_cdc; select * from runoob_tbl1;Le résultat indique que les données de la table MySQL ont bien été synchronisées vers StarRocks.
+-----------+--------------+---------------+-----------------+---------+ | runoob_id | runoob_title | runoob_author | submission_date | add_col | +-----------+--------------+---------------+-----------------+---------+ | 18 | first | tom | 2022-06-22 | 3 | +-----------+--------------+---------------+-----------------+---------+
Interrogation des données insérées
-
Dans la console SQL de l'instance ApsaraDB RDS for MySQL, exécutez la commande suivante pour insérer des données.
INSERT INTO runoob_tbl(`runoob_id`,`runoob_title`,`runoob_author`,`submission_date`,`add_col`) values(1,'second','tom2','2022-06-23',1) -
Dans le CLI StarRocks, exécutez la commande suivante pour consulter les données de la table.
select * from runoob_tbl1;Le résultat confirme que les données ont été insérées avec succès.
+-----------+--------------+---------------+-----------------+---------+ | runoob_id | runoob_title | runoob_author | submission_date | add_col | +-----------+--------------+---------------+-----------------+---------+ | 1 | second | tom2 | 2022-06-23 | 1 | | 18 | first | tom | 2022-06-22 | 3 | +-----------+--------------+---------------+-----------------+---------+
Synchronisation des mises à jour de données
-
Dans la console SQL de l'instance ApsaraDB RDS for MySQL, exécutez la commande suivante pour mettre à jour des données spécifiques.
update runoob_tbl set runoob_title= 'new' where runoob_id = 18 -
Dans le CLI StarRocks, exécutez la commande suivante pour consulter les données de la table.
select * from runoob_tbl1;Le résultat montre que la mise à jour des données a été synchronisée.
+-----------+--------------+---------------+-----------------+---------+ | runoob_id | runoob_title | runoob_author | submission_date | add_col | +-----------+--------------+---------------+-----------------+---------+ | 1 | second | tom2 | 2022-06-23 | 1 | | 18 | new | tom | 2022-06-22 | 3 | +-----------+--------------+---------------+-----------------+---------+
Synchronisation des suppressions de données
-
Dans la console SQL de l'instance ApsaraDB RDS for MySQL, exécutez la commande suivante pour supprimer des données spécifiques.
DELETE FROM runoob_tbl WHERE runoob_id = 1 -
Dans le CLI StarRocks, exécutez la commande suivante pour consulter les données de la table.
select * from runoob_tbl1;Le résultat indique que la suppression des données a été synchronisée.
+-----------+--------------+---------------+-----------------+---------+ | runoob_id | runoob_title | runoob_author | submission_date | add_col | +-----------+--------------+---------------+-----------------+---------+ | 18 | new | tom | 2022-06-22 | 3 | +-----------+--------------+---------------+-----------------+---------+
Ajout d'une colonne nullable
-
Dans la console SQL de l'instance ApsaraDB RDS for MySQL, exécutez la commande suivante pour ajouter une colonne nullable.
alter table `runoob_tbl` add COLUMN `add_col2` INT; -
Exécutez la commande suivante pour insérer des données.
INSERT INTO runoob_tbl(`runoob_id`,`runoob_title`,`runoob_author`,`submission_date`,`add_col`,`add_col2`) values(1,'second','tom2','2022-06-23',1,2) -
Dans le CLI StarRocks, exécutez la commande suivante pour consulter les données de la table.
select * from runoob_tbl1;Le résultat confirme que la modification du schéma a réussi.
+-----------+--------------+---------------+-----------------+---------+----------+ | runoob_id | runoob_title | runoob_author | submission_date | add_col | add_col2 | +-----------+--------------+---------------+-----------------+---------+----------+ | 18 | new | tom | 2022-06-22 | 3 | NULL | +-----------+--------------+---------------+-----------------+---------+----------+ | 1 | second | tom2 | 2022-06-23 | 1 | 2 | | 18 | first | tom | 2022-06-22 | 3 | NULL | +-----------+--------------+---------------+-----------------+---------+----------+
Présentation de CDAS
Une instruction CDAS constitue un raccourci syntaxique pour CTAS. Elle permet de synchroniser une base de données MySQL entière vers StarRocks via un seul déploiement Flink. Vous pouvez également utiliser la syntaxe « including table » pour ne synchroniser qu'un sous-ensemble de tables.
Comme pour CTAS, vous devez créer les catalogs MySQL et StarRocks correspondants avant d'exécuter une instruction CDAS. L'exemple suivant illustre la syntaxe à utiliser.
CREATE DATABASE IF NOT EXISTS sr_db with (
'starrocks.create.table.properties'=' buckets 8',
'starrocks.create.table.mode'='simple',
'jdbc-url'='jdbc:mysql://172.16.**.**:9030',
'load-url'='172.16.**.**:18030',
'username'='test',
'password' = '123456',
'sink.buffer-flush.interval-ms' = '5000' ,
'sink.properties.row_delimiter' = '\x02',
'sink.properties.column_separator' = '\x01'
)
as DATABASE mysql.test_cdc including table 'tabl1','tbl2','tbl3'
/*+ OPTIONS ( 'connector' = 'mysql-cdc',
'hostname' = 'rm-2zepd6e20u3od****.mysql.rds.aliyuncs.com',
'port' = '3306',
'username' = 'test',
'password' = '123456',
'database-name' = 'test_cdc' )*/;