Tous les produits
Search
Centre de documentation

ApsaraMQ for Kafka:Create an AnalyticDB sink connector

Dernière mise à jour :Aug 11, 2026

Un connecteur de réception AnalyticDB lit les messages d'une rubrique ApsaraMQ for Kafka et les écrit dans une base de données AnalyticDB for MySQL ou AnalyticDB for PostgreSQL. Le connecteur utilise Function Compute pour transférer les données entre les services au sein de la même région.

Fonctionnement

Les données transitent par trois composants :

  1. ApsaraMQ for Kafka lit les messages depuis la rubrique source.

  2. Function Compute reçoit les messages et les écrit dans la base de données de destination.

  3. AnalyticDB stocke les données dans la table spécifiée.

En interne, le connecteur utilise cinq rubriques (pour les offsets, la configuration, l'état, la file d'attente de lettres mortes et les données d'erreur) ainsi qu'un groupe de consommateurs. Vous pouvez créer ces ressources automatiquement ou manuellement.

Prérequis

Avant de commencer, assurez-vous de disposer des éléments suivants :

Rassemblez les informations suivantes pour configurer le connecteur :

Informations Description Exemple
Nom de la rubrique source La rubrique Kafka à partir de laquelle exporter les données adb-test-input
ID d'instance AnalyticDB L'ID de l'instance de base de données de destination am-bp139yqk8u1ik****
Nom de la base de données La base de données de destination adb_demo
Nom de la table La table de destination user
Nom d'utilisateur de la base de données Le nom d'utilisateur de connexion à la base de données adbmysql
Mot de passe de la base de données Le mot de passe de connexion à la base de données ********

Notes d'utilisation

  • L'exportation des données est limitée à la même région. L'exportation interrégionale n'est pas prise en charge. Pour plus d'informations, consultez les limites.

  • Function Compute offre un quota de ressources gratuit. L'utilisation au-delà du quota gratuit est facturée selon la tarification de Function Compute.

  • ApsaraMQ for Kafka sérialise les messages sous forme de chaînes encodées en UTF-8. Les données binaires ne sont pas prises en charge.

  • Si la base de données de destination utilise un endpoint privé, configurez le même VPC et le même vSwitch pour le service Function Compute. Sinon, Function Compute ne pourra pas atteindre la base de données. Pour plus d'informations, consultez la section Mettre à jour un service.

  • ApsaraMQ for Kafka crée automatiquement un rôle lié au service lors de la création d'un connecteur, s'il n'existe pas déjà.

  • Pour résoudre les problèmes d'exécution des fonctions, utilisez la journalisation de Function Compute. Pour plus d'informations, consultez la section Configurer la journalisation.

Configurer le connecteur

Pour configurer le connecteur :

  1. (Facultatif) Créer les rubriques requises et le groupe de consommateurs

  2. Créer et déployer le connecteur

  3. Configurer la mise en réseau de Function Compute

  4. Ajouter le bloc CIDR du VPC à la liste d'autorisation d'AnalyticDB

  5. Vérifier l'exportation des données

Étape 1 : (Facultatif) Créer les rubriques requises et le groupe de consommateurs

Pour qu'ApsaraMQ for Kafka crée automatiquement ces ressources, ignorez cette étape et définissez Resource Creation Method sur Auto à l' étape 2 .
Important

Si votre instance ApsaraMQ for Kafka exécute la version majeure 0.10.2, les rubriques nécessitant le moteur de stockage Local Storage doivent être créées automatiquement. La création manuelle de rubriques Local Storage n'est pas prise en charge dans cette version.

Le connecteur nécessite cinq rubriques internes et un groupe de consommateurs. Créez-les manuellement uniquement si vous avez besoin de configurations spécifiques.

Rubriques requises

Rubrique Préfixe de nom Partitions Moteur de stockage Politique de nettoyage des journaux
Offset de tâche connect-offset Supérieur à 1 Local Storage Compact
Configuration de tâche connect-config 1 Local Storage Compact
État de tâche connect-status 6 (recommandé) Local Storage Compact
File d'attente de lettres mortes connect-error 6 (recommandé) Local Storage ou Cloud Storage Tout
Données d'erreur connect-error 6 (recommandé) Local Storage ou Cloud Storage Tout
Conseil : La rubrique de file d'attente de lettres mortes et la rubrique de données d'erreur peuvent partager la même rubrique afin d'économiser des ressources.

Créer une rubrique

  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.

Important

Vous devez créer les rubriques dans la région où votre instance Elastic Compute Service (ECS) est déployée. Une rubrique ne peut pas être utilisée entre différentes régions. Par exemple, si les producteurs et les consommateurs de messages s'exécutent sur une instance ECS déployée dans la région Chine (Pékin), la rubrique doit également être créée dans la région Chine (Pékin).

  1. Sur la page Instances, cliquez sur le nom de l'instance.

  2. Dans le volet de navigation de gauche, cliquez sur Topics.

  3. Sur la page Topics, cliquez sur Create Topic.

  4. Dans le panneau Create Topic, configurez les paramètres suivants et cliquez sur OK.

ParamètreDescription
NameLe nom de la rubrique. Utilisez le préfixe de nommage indiqué dans le tableau ci-dessus (par exemple, connect-offset-kafka-adb-sink).
PartitionsLe nombre de partitions. Consultez les valeurs requises dans le tableau ci-dessus.
Storage EngineLe type de moteur de stockage. Disponible uniquement sur les instances Professional Edition. Les instances Standard Edition utilisent Cloud Storage par défaut. Options : Cloud Storage (disques Alibaba Cloud, stockage distribué à 3 réplicas avec faible latence et haute fiabilité) ou Local Storage (algorithme ISR Apache Kafka, stockage distribué à 3 réplicas). Les instances Standard (High Write) ne prennent en charge que Cloud Storage.
Message TypeLa garantie d'ordre des messages. Normal Message (par défaut pour Cloud Storage) : l'ordre des partitions peut ne pas être préservé en cas de défaillance du broker. Partitionally Ordered Message (par défaut pour Local Storage) : l'ordre est préservé en cas de défaillance du broker, mais les partitions concernées ne sont pas disponibles jusqu'à leur restauration.
Log Cleanup PolicyRequis lorsque Storage Engine est défini sur Local Storage. Delete (par défaut) : conserve les messages selon la période de rétention ; supprime les messages les plus anciens lorsque le stockage dépasse 85 %. Compact : conserve uniquement la dernière valeur par clé. Requis pour les rubriques internes du connecteur qui utilisent Kafka Connect.
Important

Les rubriques compactées par journalisation ne peuvent être utilisées que dans des composants cloud-native spécifiques, tels que Kafka Connect et Confluent Schema Registry. Pour plus d'informations, consultez aliware-kafka-demos.

DescriptionUne description facultative.
TagDes tags facultatifs.
  1. Répétez l'opération pour créer les cinq rubriques requises.

Créer le groupe de consommateurs

Le connecteur nécessite un groupe de consommateurs nommé connect-<nom-du-connecteur> (par exemple, connect-kafka-adb-sink).

  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. Sur la page Instances, cliquez sur le nom de l'instance.

  4. Dans le volet de navigation de gauche, cliquez sur Groups.

  5. Sur la page Groups, cliquez sur Create Group.

  6. Dans le panneau Create Group, saisissez le nom du groupe dans le champ Group ID, ajoutez une description et des tags facultatifs, puis cliquez sur OK.

Étape 2 : 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. Sur la page Instances, cliquez sur le nom de l'instance.

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

  5. Sur la page Connectors, cliquez sur Create Connector.

  6. Suivez l'assistant Create Connector :

Configurer les informations de base

Paramètre Description Exemple
Name Le nom du connecteur. 1 à 48 caractères : chiffres, lettres minuscules et traits d'union (-). Ne peut pas commencer par un trait d'union. Doit être unique au sein de l'instance. kafka-adb-sink
Instance L'instance ApsaraMQ for Kafka. Affiche le nom et l'ID de l'instance. demo alikafka_post-cn-st21p8vj****

Cliquez sur Next.

Configurer le service source

Sélectionnez Message Queue for Apache Kafka comme service source et configurez les paramètres suivants.

Paramètre Description Exemple
Data Source Topic La rubrique à partir de laquelle exporter les données. adb-test-input
Consumer Thread Concurrency Le nombre de threads de consommation simultanés. Par défaut : 6. Valeurs valides : 1, 2, 3, 6, 12. 6
Consumer Offset Point de départ de la consommation. Earliest Offset : commencez depuis le début. Latest Offset : commencez depuis le message le plus récent. Earliest Offset

Cliquez sur Configure Runtime Environment pour développer les paramètres avancés.

ParamètreDescriptionExemple
VPC IDLe VPC pour la tâche d'exportation. Par défaut, il s'agit du VPC de l'instance ApsaraMQ for Kafka.vpc-bp1xpdnd3l***
vSwitch IDLe vSwitch pour la tâche d'exportation. Doit se trouver dans le même VPC que l'instance.vsw-bp1d2jgg81***
Failure Handling PolicyAction en cas d'échec d'envoi de message. Continue Subscription : continuer la consommation et enregistrer l'erreur. Stop Subscription : arrêter la consommation et enregistrer l'erreur. Pour plus de détails sur les journaux, consultez Gérer un connecteur. Pour les codes d'erreur, consultez Codes d'erreur.Continue Subscription
Resource Creation MethodMéthode de création des rubriques internes requises et du groupe de consommateurs. Auto : ApsaraMQ for Kafka les crée automatiquement. Manual : utiliser les ressources créées à l'étape 1.Auto
Connector Consumer GroupLe groupe de consommateurs pour le connecteur. Format : connect-<nom-du-connecteur>.connect-kafka-adb-sink
Task Offset TopicStocke les offsets des consommateurs. Préfixe de nom : connect-offset. Partitions : supérieur à 1. Moteur de stockage : Local Storage. Politique de nettoyage : Compact.connect-offset-kafka-adb-sink
Task Configuration TopicStocke les configurations des tâches. Préfixe de nom : connect-config. Partitions : 1. Moteur de stockage : Local Storage. Politique de nettoyage : Compact.connect-config-kafka-adb-sink
Task Status TopicStocke l'état des tâches. Préfixe de nom : connect-status. Partitions : 6 (recommandé). Moteur de stockage : Local Storage. Politique de nettoyage : Compact.connect-status-kafka-adb-sink
Dead-letter Queue TopicStocke les données d'erreur du framework Kafka Connect. Préfixe de nom : connect-error. Partitions : 6 (recommandé). Moteur de stockage : Local Storage ou Cloud Storage. Peut partager une rubrique avec la rubrique de données d'erreur.connect-error-kafka-adb-sink
Error Data TopicStocke les données d'erreur du connecteur de réception. Préfixe de nom : connect-error. Partitions : 6 (recommandé). Moteur de stockage : Local Storage ou Cloud Storage. Peut partager une rubrique avec la rubrique de file d'attente de lettres mortes.connect-error-kafka-adb-sink

Cliquez sur Next.

Configurer le service de destination

Sélectionnez AnalyticDB comme service de destination et configurez les paramètres suivants.

Paramètre Description Exemple
Instance Type Le type de base de données de destination : AnalyticDB for MySQL ou AnalyticDB for PostgreSQL. AnalyticDB for MySQL
AnalyticDB Instance ID L'ID de l'instance de destination. am-bp139yqk8u1ik****
Database Name La base de données de destination. adb_demo
Table Name La table de destination pour les données exportées. user
Database Username Le nom d'utilisateur de connexion à la base de données. adbmysql
Database Password Le mot de passe de connexion à la base de données. Pour réinitialiser un mot de passe oublié : pour AnalyticDB for MySQL, réinitialisez-le dans la console AnalyticDB for MySQL. Pour AnalyticDB for PostgreSQL, accédez à Account Management et cliquez sur Reset Password. ********
Le nom d'utilisateur et le mot de passe de la base de données sont transmis à Function Compute en tant que variables d'environnement lorsqu'ApsaraMQ for Kafka crée la tâche d'exportation. ApsaraMQ for Kafka ne stocke pas ces identifiants après la création de la tâche.

Cliquez sur Create.

  1. Sur la page Connectors, recherchez le connecteur et cliquez sur Deploy dans la colonne Actions.

Étape 3 : Configurer la mise en réseau de Function Compute

Après le déploiement, Function Compute crée automatiquement un service (nommé kafka-service-<nom_du_connecteur>-<chaîne_aléatoire>) et une fonction (nommée fc-adb-<chaîne_aléatoire>) pour le connecteur.

  1. Sur la page Connectors, recherchez le connecteur et cliquez sur Configure Function dans la colonne Actions. Cela ouvre la console Function Compute.

  2. Dans la console Function Compute, localisez le service créé automatiquement et configurez le VPC et le vSwitch pour qu'ils correspondent à la base de données de destination. Pour plus d'informations, consultez la section Mettre à jour un service.

Étape 4 : Ajouter le bloc CIDR du VPC à la liste d'autorisation d'AnalyticDB

Ajoutez le bloc CIDR du VPC configuré dans Function Compute à la liste d'autorisation d'AnalyticDB. Trouvez le bloc CIDR sur la page vSwitch de la console VPC — il se trouve dans la ligne correspondant au VPC et au vSwitch du service Function Compute.

Étape 5 : Vérifier l'exportation des données

Envoyer des messages de test

Important

Le contenu du message doit être un JSON valide. Chaque clé JSON correspond à un nom de colonne dans la table de destination, et chaque valeur correspond aux données de la colonne. Vérifiez que les clés JSON correspondent aux noms de colonnes de votre table de destination avant d'envoyer des messages de test.

  1. Sur la page Connectors, recherchez le connecteur et cliquez sur Test dans la colonne Actions.

  2. Dans le panneau Send Message, sélectionnez une Sending Method :

    • Console :

      1. Dans le champ Message Key, saisissez une clé (par exemple, demo).

      2. Dans le champ Message Content, saisissez du contenu JSON (par exemple, {"key": "test"}).

      3. Pour Send to Specified Partition, sélectionnez Yes et saisissez un Partition ID (par exemple, 0) pour cibler une partition spécifique, ou sélectionnez No pour laisser ApsaraMQ for Kafka attribuer la partition. Pour plus d'informations sur les ID de partition, consultez la section Afficher l'état des partitions.

    • Docker : Exécutez la commande Docker affichée dans la section Run the Docker container to produce a sample message.

    • SDK : Sélectionnez un SDK pour votre langage de programmation et une méthode d'accès pour envoyer le message de test.

Vérifier le résultat

Après avoir envoyé des messages, vérifiez la table de destination pour confirmer que les données ont été exportées :

  1. Connectez-vous à la console AnalyticDB for MySQL ou à la console AnalyticDB for PostgreSQL.

  2. Connectez-vous à la base de données de destination.

  3. Dans la SQLConsole de la console Data Management Service 5.0, ouvrez la table de destination et confirmez que les données exportées sont présentes.

Étapes suivantes

  • Gérer un connecteur : affichez l'état du connecteur, mettez en pause, reprenez ou supprimez des connecteurs.

  • Configurer la journalisation : configurez les journaux de Function Compute pour résoudre les problèmes d'exportation de données.