Un connecteur source MySQL synchronise les modifications au niveau des lignes d'une base de données ApsaraDB RDS for MySQL vers les topics de votre instance ApsaraMQ for Kafka. Ce connecteur utilise DataWorks pour capturer les événements CDC (Change Data Capture), notamment les opérations INSERT, UPDATE et DELETE, puis les achemine vers les topics Kafka.
Fonctionnement
Lorsque vous créez un connecteur source MySQL via la console ApsaraMQ for Kafka, le système exécute automatiquement les actions suivantes :
Active DataWorks Basic Edition (gratuitement).
Crée un espace de travail DataWorks ainsi qu'un groupe de ressources exclusif pour l'intégration de données.
Génère les topics de destination dans votre instance Kafka, à raison d'un topic par table source.
Le groupe de ressources exclusif (4 vCPU, 8 Go de mémoire) exécute les tâches de synchronisation des données. Ces groupes de ressources fonctionnent sur abonnement mensuel et se renouvellent automatiquement à expiration.
Les groupes de ressources exclusifs pour l'intégration de données sont payants. Pour plus de détails sur les tarifs, consultez Vue d'ensemble de la facturation.

Nommage des topics et partitions
DataWorks génère un topic de destination par table source selon le modèle de nommage suivant :
<topic-prefix>_<source-table-name>
Le trait de soulignement est ajouté automatiquement. Par exemple, si le préfixe est mysql et que les tables sources sont table_1, table_2 et table_n, les topics générés sont mysql_table_1, mysql_table_2 et mysql_table_n.

Le nombre de partitions dépend de la présence d'une clé primaire dans la table source :
| Type de table source | Partitions par topic |
|---|---|
| Avec clé primaire | 6 |
| Sans clé primaire | 1 |
Vérifiez que votre instance Kafka dispose d'un nombre suffisant de topics et de partitions disponibles. Si l'une de ces ressources vient à manquer, la création du topic échoue et le déploiement du connecteur devient impossible.
Région et réseau
Même région : lorsque l'instance RDS MySQL et l'instance Kafka se trouvent dans la même région, le système crée automatiquement une interface réseau élastique (ENI) dans le Virtual Private Cloud (VPC) correspondant et l'associe à l'instance Elastic Compute Service (ECS) du groupe de ressources exclusif. Aucune configuration réseau manuelle n'est requise.
Régions différentes : si les instances se situent dans des régions distinctes, vérifiez les points suivants :
Une instance Cloud Enterprise Network (CEN) existe sous le même compte Alibaba Cloud.
Les VPC de l'instance RDS MySQL et de l'instance Kafka sont tous deux rattachés à l'instance CEN.
La bande passante interrégionale est configurée pour l'instance CEN.
En l'absence d'une configuration CEN adéquate, le système risque de créer automatiquement une instance CEN avec une bande passante minimale. Cela peut entraîner des erreurs de connectivité lors de la création du connecteur ou pendant son exécution.
Limites du groupe de ressources exclusif
| Limite | Valeur |
|---|---|
| Nombre maximal de connecteurs par groupe de ressources | 3 |
| Nombre maximal d'associations ENI VPC par groupe de ressources | 2 |
Si un groupe de ressources existant héberge moins de trois connecteurs, DataWorks le réutilise pour les nouveaux connecteurs. Toutefois, en cas de chevauchement de blocs CIDR ou d'autres contraintes techniques empêchant cette réutilisation, DataWorks crée un nouveau groupe de ressources.
Prérequis
Avant de commencer, assurez-vous de disposer des éléments suivants :
Une instance ApsaraMQ for Kafka avec la fonctionnalité connecteur activée, déployée dans l'une des régions suivantes : Chine (Shenzhen), Chine (Chengdu), Chine (Pékin), Chine (Zhangjiakou), Chine (Hangzhou), Chine (Shanghai) ou Singapour.
Une instance ApsaraDB RDS for MySQL comprenant une base de données, un compte de base de données et au moins une table créée. Pour les références SQL, consultez Instructions SQL courantes pour MySQL.
Des blocs CIDR ne se chevauchant pas entre le VPC de votre instance RDS MySQL et celui de votre instance Kafka.
DataWorks autorisé à accéder à vos ENI via la page Cloud Resource Access Authorization.
La source de données (instance ApsaraDB RDS for MySQL) et la destination de données (instance ApsaraMQ for Kafka) créées sous le même compte Alibaba Cloud.
Permissions du compte de base de données
Le compte de base de données MySQL doit posséder au minimum les permissions suivantes :
| Permission | Objectif |
|---|---|
SELECT |
Lire les lignes des tables sources |
REPLICATION SLAVE |
Se connecter au journal binaire MySQL et le lire |
REPLICATION CLIENT |
Interroger l'état du journal binaire |
Accordez ces permissions à l'aide de l'instruction suivante :
GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO '<your-username>'@'%';
Remplacez <your-username> par le nom d'utilisateur réel du compte de base de données.
Permissions de l'utilisateur RAM
Si vous utilisez un utilisateur Resource Access Management (RAM) plutôt qu'un compte Alibaba Cloud, attachez les politiques suivantes à cet utilisateur RAM :
| Politique | Description |
|---|---|
AliyunDataWorksFullAccess |
Gérer toutes les ressources DataWorks du compte Alibaba Cloud |
AliyunBSSOrderAccess |
Acheter des services Alibaba Cloud |
Pour obtenir les instructions détaillées, consultez Accorder des permissions aux utilisateurs RAM.
Créer et déployer le connecteur
Connectez-vous à la console ApsaraMQ for Kafka.
Dans la section Resource Distribution de la page Overview, sélectionnez la région de votre instance.
Dans le volet de navigation de gauche, cliquez sur Connectors.
Sur la page Connectors, sélectionnez votre instance dans la liste déroulante Select Instance, puis cliquez sur Create Connector.
-
Suivez l'assistant en trois étapes : cliquez sur Next. Step 2: Configure source service Sélectionnez ApsaraDB RDS for MySQL comme service source et configurez les paramètres suivants : cliquez sur Next. Step 3: Configure destination service Confirmez l'instance Kafka de destination et cliquez sur Create.
Step 1: Configure basic information
Paramètre Description Exemple Name Nom unique du connecteur au sein de l'instance. De 1 à 48 caractères ; chiffres, lettres minuscules et traits d'union ( -) autorisés. Ne peut pas commencer par un trait d'union. Le système crée automatiquement un groupe de consommateurs nomméconnect-<connector-name>.kafka-source-mysqlInstance Affiche le nom et l'ID de l'instance Kafka sélectionnée. demo alikafka_post-cn-st21p8vj****Paramètre Description Exemple Region of ApsaraDB RDS for MySQL Instance Région de l'instance RDS MySQL source. China (Shenzhen) ApsaraDB RDS for MySQL Instance ID ID de l'instance de la base de données source. rm-wz91w3vk6owmz****Database Name Nom de la base de données à synchroniser. mysql-to-kafkaDatabase Account Nom d'utilisateur pour la connexion à la base de données. mysql_to_kafkaPassword of Database Account Mot de passe pour la connexion à la base de données. -- Database Table Un ou plusieurs noms de tables, séparés par des virgules ( ,). Chaque table correspond à un topic.mysql_tblTables to Automatically Add Expression régulière permettant de détecter et synchroniser automatiquement les nouvelles tables. Utilisez .*pour inclure toutes les tables..*Topic Prefix Préfixe pour les noms de topics générés automatiquement. Doit être globalement unique. mysql Sur la page Connectors, localisez le connecteur et cliquez sur Deploy dans la colonne Actions. Lorsque la colonne Status affiche Running, le connecteur est actif et synchronise les données.
Si le déploiement du connecteur échoue, vérifiez que tous les prérequis sont remplis. Pour modifier les paramètres du connecteur après le déploiement, cliquez sur
Task Configurations
dans la colonne
Actions
pour ouvrir la console DataWorks.
Vérifier la synchronisation des données
-
Insérez une ligne de test dans une table source : pour plus d'exemples SQL, consultez Instructions SQL courantes pour MySQL.
INSERT INTO mysql_tbl (mysql_title, mysql_author, submission_date) VALUES ("mysql2kafka", "tester", NOW()); -
Utilisez la fonctionnalité de requête de messages d'ApsaraMQ for Kafka pour confirmer l'arrivée des données dans le topic correspondant. Pour les instructions, consultez Requêter des messages. Un événement INSERT synchronisé avec succès se présente comme suit :
{ "schema": { "dataColumn": [ { "name": "mysql_id", "type": "LONG" }, { "name": "mysql_title", "type": "STRING" }, { "name": "mysql_author", "type": "STRING" }, { "name": "submission_date", "type": "DATE" } ], "primaryKey": [ "mysql_id" ], "source": { "dbType": "MySQL", "dbName": "mysql_to_kafka", "tableName": "mysql_tbl" } }, "payload": { "before": null, "after": { "dataColumn": { "mysql_title": "mysql2kafka", "mysql_author": "tester", "submission_date": 1614700800000 } }, "sequenceId": "1614748790461000000", "timestamp": { "eventTime": 1614748870000, "systemTime": 1614748870925, "checkpointTime": 1614748870000 }, "op": "INSERT", "ddl": null }, "version": "0.0.1" }
Champs du message
| Champ | Description |
|---|---|
schema.dataColumn |
Noms des colonnes et types de données de la table source |
schema.primaryKey |
Colonnes de clé primaire de la table source |
schema.source |
Type de base de données source, nom de la base de données et nom de la table |
payload.before |
État de la ligne avant la modification (null pour les opérations INSERT) |
payload.after |
État de la ligne après la modification (null pour les opérations DELETE) |
payload.op |
Type d'opération : INSERT, UPDATE ou DELETE |
payload.timestamp |
Horodatages de l'événement, du système et du point de contrôle |
Pour consulter la spécification complète du format de message, reportez-vous à Formats de message.
Étapes suivantes
Requêter des messages dans vos topics Kafka
Activer la fonctionnalité connecteur pour d'autres instances