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 :
ApsaraMQ for Kafka lit les messages depuis la rubrique source.
Function Compute reçoit les messages et les écrit dans la base de données de destination.
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 :
Une instance ApsaraMQ for Kafka avec la fonctionnalité de connecteur activée
Une rubrique source dans l'instance ApsaraMQ for Kafka
-
Une base de données de destination AnalyticDB :
AnalyticDB for MySQL : un cluster, un compte de base de données, une connexion client et une base de données
AnalyticDB for PostgreSQL : une instance, un compte de base de données et une connexion client
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 :
(Facultatif) Créer les rubriques requises et le groupe de consommateurs
Ajouter le bloc CIDR du VPC à la liste d'autorisation d'AnalyticDB
É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 .
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
Connectez-vous à la console ApsaraMQ for Kafka.
Dans la section Resource Distribution de la page Overview, sélectionnez la région de votre instance.
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).
Sur la page Instances, cliquez sur le nom de l'instance.
Dans le volet de navigation de gauche, cliquez sur Topics.
Sur la page Topics, cliquez sur Create Topic.
Dans le panneau Create Topic, configurez les paramètres suivants et cliquez sur OK.
| Paramètre | Description |
|---|---|
| Name | Le nom de la rubrique. Utilisez le préfixe de nommage indiqué dans le tableau ci-dessus (par exemple, connect-offset-kafka-adb-sink). |
| Partitions | Le nombre de partitions. Consultez les valeurs requises dans le tableau ci-dessus. |
| Storage Engine | Le 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 Type | La 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 Policy | Requis 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. |
| Description | Une description facultative. |
| Tag | Des tags facultatifs. |
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).
Connectez-vous à la console ApsaraMQ for Kafka.
Dans la section Resource Distribution de la page Overview, sélectionnez la région de votre instance.
Sur la page Instances, cliquez sur le nom de l'instance.
Dans le volet de navigation de gauche, cliquez sur Groups.
Sur la page Groups, cliquez sur Create Group.
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
Connectez-vous à la console ApsaraMQ for Kafka.
Dans la section Resource Distribution de la page Overview, sélectionnez la région de votre instance.
Sur la page Instances, cliquez sur le nom de l'instance.
Dans le volet de navigation de gauche, cliquez sur Connectors.
Sur la page Connectors, cliquez sur Create Connector.
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ètre | Description | Exemple |
|---|---|---|
| VPC ID | Le VPC pour la tâche d'exportation. Par défaut, il s'agit du VPC de l'instance ApsaraMQ for Kafka. | vpc-bp1xpdnd3l*** |
| vSwitch ID | Le vSwitch pour la tâche d'exportation. Doit se trouver dans le même VPC que l'instance. | vsw-bp1d2jgg81*** |
| Failure Handling Policy | Action 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 Method | Mé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 Group | Le groupe de consommateurs pour le connecteur. Format : connect-<nom-du-connecteur>. | connect-kafka-adb-sink |
| Task Offset Topic | Stocke 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 Topic | Stocke 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 Topic | Stocke 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 Topic | Stocke 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 Topic | Stocke 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.
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.
Sur la page Connectors, recherchez le connecteur et cliquez sur Configure Function dans la colonne Actions. Cela ouvre la console Function Compute.
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.
AnalyticDB for MySQL : configurez la liste d'autorisation dans la console AnalyticDB for MySQL. Pour plus d'informations, consultez la section Configurer une liste d'autorisation d'adresses IP.
AnalyticDB for PostgreSQL : configurez la liste d'autorisation dans la console AnalyticDB for PostgreSQL. Pour plus d'informations, consultez la section Configurer une liste d'autorisation d'adresses IP.
Étape 5 : Vérifier l'exportation des données
Envoyer des messages de test
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.
Sur la page Connectors, recherchez le connecteur et cliquez sur Test dans la colonne Actions.
-
Dans le panneau Send Message, sélectionnez une Sending Method :
-
Console :
Dans le champ Message Key, saisissez une clé (par exemple,
demo).Dans le champ Message Content, saisissez du contenu JSON (par exemple,
{"key": "test"}).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 :
Connectez-vous à la console AnalyticDB for MySQL ou à la console AnalyticDB for PostgreSQL.
Connectez-vous à la base de données de destination.
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.