Tous les produits
Search
Centre de documentation

ApsaraMQ for Kafka:Créer un connecteur source MySQL

Dernière mise à jour :Aug 12, 2026

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 :

  1. Active DataWorks Basic Edition (gratuitement).

  2. Crée un espace de travail DataWorks ainsi qu'un groupe de ressources exclusif pour l'intégration de données.

  3. 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.

Important

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.

mysql_connector

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.

table_topic_match

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 :

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

  1. Connectez-vous à la console ApsaraMQ for Kafka.

  2. Dans la section Resource Distribution de la page Overview, sélectionnez la région de votre instance.

  3. Dans le volet de navigation de gauche, cliquez sur Connectors.

  4. Sur la page Connectors, sélectionnez votre instance dans la liste déroulante Select Instance, puis cliquez sur Create Connector.

  5. 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-mysql
    Instance 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-kafka
    Database Account Nom d'utilisateur pour la connexion à la base de données. mysql_to_kafka
    Password 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_tbl
    Tables 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
  6. 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.

Remarque

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

  1. 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());
  2. 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