Un connecteur de sortie Elasticsearch lit les messages d'une rubrique de votre instance ApsaraMQ for Kafka et les écrit dans un index Elasticsearch. Le connecteur utilise Function Compute comme intermédiaire : il consomme les messages depuis la rubrique source, les transmet à une fonction Function Compute, puis cette fonction les écrit dans Elasticsearch via l'API Bulk. Chaque message devient un document dans l'index cible, avec des métadonnées telles que la rubrique, la partition, le décalage (offset) et l'horodatage.
Avant de commencer
Effectuez les configurations suivantes avant de créer le connecteur.
ApsaraMQ for Kafka
Activer la fonctionnalité de connecteur pour votre instance.
Créer une rubrique à utiliser comme source de données.
Function Compute
Elasticsearch
Créer un cluster et un index Elasticsearch dans la console Elasticsearch. Utilisez la version 7.0 ou ultérieure pour assurer la compatibilité avec le client Function Compute (version 7.7.0).
Ajouter le bloc CIDR du endpoint Function Compute à la liste d'autorisation Elasticsearch. Pour les tests initiaux, spécifiez
0.0.0.0/0afin d'autoriser toutes les adresses IP du VPC, puis restreignez la plage après avoir vérifié la connectivité.
Informations à recueillir
Rassemblez les informations suivantes avant de démarrer l'assistant.
| Information | Où la trouver | Exemple |
|---|---|---|
| ID de l'instance Elasticsearch | Console Elasticsearch | es-cn-oew1o67x0000**** |
| Endpoint Elasticsearch (public ou privé) | Informations de base du cluster | es-cn-oew1o67x0000****.elasticsearch.aliyuncs.com |
| Port Elasticsearch | Informations de base du cluster | 9200 (HTTP/HTTPS) ou 9300 (TCP) |
| Nom d'utilisateur et mot de passe Elasticsearch | Définis lors de la création du cluster ; réinitialiser si nécessaire | elastic / **** |
| Nom de l'index Elasticsearch | Console Elasticsearch | elastic_test |
| Nom de la rubrique source | Console ApsaraMQ for Kafka | elasticsearch-test-input |
Limites
L'instance ApsaraMQ for Kafka et le cluster Elasticsearch doivent se trouver dans la même région.
ApsaraMQ for Kafka sérialise les messages sous forme de chaînes UTF-8. Les données binaires ne sont pas prises en charge.
Si vous spécifiez l'endpoint privé du cluster Elasticsearch, Function Compute ne peut pas y accéder par défaut. Pour activer la connectivité, configurez le service Function Compute pour utiliser le même VPC et le même vSwitch que le cluster Elasticsearch. Consultez la section Configurer le service Function Compute.
Pour connaître les autres limites relatives aux connecteurs, consultez la section Limites.
Facturation
Le connecteur utilise Function Compute pour exporter les données. Function Compute propose une offre gratuite. L'utilisation au-delà de cette offre gratuite est facturée selon les modalités de facturation de Function Compute.
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 votre région.
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 et cliquez sur Create Connector.
Étape 1 : Configurer les informations de base
À l'étape Configure Basic Information, définissez le nom du connecteur et examinez les détails de l'instance.
| Paramètre | Description | Exemple |
|---|---|---|
| Name | Un nom unique au sein de l'instance ApsaraMQ for Kafka. Utilisez 1 à 48 caractères : chiffres, lettres minuscules et traits d'union (-). Ne peut pas commencer par un trait d'union. Le connecteur crée automatiquement un groupe de consommateurs nommé connect-<connector-name>. |
kafka-elasticsearch-sink |
| Instance | Affiche le nom et l'ID de l'instance actuelle. | demo alikafka_post-cn-st21p8vj**** |
Par défaut, l'option Authorize to Create Service Linked Role est sélectionnée. ApsaraMQ for Kafka crée un rôle lié au service s'il n'existe pas déjà.
Cliquez sur Next.
Étape 2 : Configurer le service source
À l'étape Configure Source Service, 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 depuis laquelle les données sont consommées. | elasticsearch-test-input |
| Consumer Thread Concurrency | Nombre de threads de consommation simultanés. Valeurs valides : 1, 2, 3, 6, 12. Par défaut : 6. | 6 |
| Consumer Offset | Point de départ de la consommation. L'option Earliest Offset lit depuis le début. L'option Latest Offset ne lit que les nouveaux messages. | Earliest Offset |
Cliquez sur Configure Runtime Environment pour afficher les paramètres supplémentaires.
| Paramètre | Description | Exemple |
|---|---|---|
| VPC ID | VPC de l'instance source. Rempli automatiquement ; aucune modification n'est nécessaire. | vpc-bp1xpdnd3l*** |
| vSwitch ID | vSwitch de l'instance source. Doit se trouver dans le même VPC. | vsw-bp1d2jgg81*** |
| Failure Handling Policy | Action à entreprendre en cas d'échec d'un message. L'option Continue Subscription consigne l'erreur et poursuit la consommation. L'option Stop Subscription consigne l'erreur et arrête la partition. Consultez la section Gérer un connecteur pour plus de détails sur les journaux et la section Codes d'erreur pour le dépannage. Remarque
| Continue Subscription |
| Resource Creation Method | L'option Auto crée automatiquement les rubriques internes requises. L'option Manual vous permet de les créer vous-même. | Auto |
| Connector Consumer Group | Groupe de consommateurs pour la tâche du connecteur. Format : connect-<connector-name>. | connect-kafka-elasticsearch-sink |
Rubriques internes (création manuelle uniquement)
Si vous définissez l'option Resource Creation Method sur Manual, créez les rubriques suivantes. Toutes les rubriques nécessitant l'option Local Storage sont disponibles uniquement sur les instances de l'édition Professionnelle.
| Paramètre | Convention de nommage | Partitions | Moteur de stockage | **cleanup.policy** |
|---|---|---|---|---|
| Task Offset Topic | connect-offset-* | Plus de 1 | Local Storage | Compact |
| Task Configuration Topic | connect-config-* | Exactement 1 | Local Storage | Compact |
| Task Status Topic | connect-status-* | 6 (recommandé) | Local Storage | Compact |
| Dead-letter Queue Topic | connect-error-* | 6 (recommandé) | Local Storage or Cloud Storage | -- |
| Error Data Topic | connect-error-* | 6 (recommandé) | Local Storage or Cloud Storage | -- |
Pour économiser les ressources de rubriques, utilisez la même rubrique pour la file d'attente des messages morts et pour la rubrique des données d'erreur.
Cliquez sur Next.
Étape 3 : Configurer le service de destination
À l'étape Configure Destination Service, sélectionnez Elasticsearch comme service de destination et configurez les paramètres suivants.
| Paramètre | Description | Exemple |
|---|---|---|
| Elasticsearch Instance ID | L'ID du cluster Elasticsearch. | es-cn-oew1o67x0000**** |
| Endpoint | L'endpoint public ou privé du cluster. Consultez la section Afficher les informations de base du cluster. | es-cn-oew1o67x0000****.elasticsearch.aliyuncs.com |
| Port | 9200 pour HTTP/HTTPS, ou 9300 pour TCP. | 9300 |
| Username | Le nom d'utilisateur Elasticsearch. Par défaut : elastic. Personnalisez via X-Pack RBAC si nécessaire. Le compte doit disposer d'autorisations d'écriture sur l'index cible. |
elastic |
| Password | Le mot de passe défini lors de la création du cluster. Réinitialisez le mot de passe si vous l'avez oublié. | **** |
| Index | Le nom de l'index Elasticsearch cible. | elastic_test |
Le nom d'utilisateur et le mot de passe sont transmis à Function Compute en tant que variables d'environnement lors de la création de la tâche du connecteur. ApsaraMQ for Kafka ne stocke pas ces identifiants après la création de la tâche.
Le compte doit disposer d'autorisations d'écriture sur l'index, car les messages sont envoyés via l'API Bulk d'Elasticsearch.
Cliquez sur Create.
Déployer le connecteur
Après la création, le connecteur apparaît sur la page Connectors. Cliquez sur Deploy dans la colonne Actions pour démarrer le connecteur.
Configurer le service Function Compute
Après le déploiement du connecteur, Function Compute crée automatiquement un service nommé kafka-service-<connector-name>-<random-string>. Si le cluster Elasticsearch utilise un endpoint privé, configurez le service Function Compute pour utiliser le même VPC et le même vSwitch que le cluster Elasticsearch.
Sur la page Connectors, recherchez le connecteur. Dans la colonne Actions, choisissez More > Configure Function.
Dans la console Function Compute, localisez le service créé automatiquement et mettez à jour les paramètres VPC et vSwitch pour qu'ils correspondent au cluster Elasticsearch.
Vérifier le flux de données
Envoyez un message de test pour confirmer que les données circulent bien d'ApsaraMQ for Kafka vers Elasticsearch.
Envoyer un message de test
Sur la page Connectors, recherchez le connecteur et cliquez sur Test dans la colonne Actions.
Dans le panneau Send Message, définissez l'option Method of Sending sur Console.
Dans le champ Message Key, saisissez une clé, par exemple
demo.-
Dans le champ Message Content, saisissez un corps JSON, par exemple :
{"key": "test"} Pour l'option Send to Specified Partition, cliquez sur Yes et saisissez un Partition ID (par exemple,
0) pour cibler une partition spécifique, ou cliquez sur No pour laisser le système en attribuer une. Pour rechercher les ID de partition, consultez la section Afficher l'état des partitions.
Vous pouvez également envoyer des messages de test via Docker ou un SDK. Sélectionnez l'option correspondante dans le champ Method of Sending et suivez les instructions affichées à l'écran.
Vérifier l'index Elasticsearch
-
Exécutez la requête suivante pour rechercher dans l'index cible :
GET /<index_name>/_search -
Confirmez que la réponse contient le message que vous avez envoyé. Une réponse réussie ressemble à ceci :
{ "took": 8, "timed_out": false, "_shards": { "total": 5, "successful": 5, "skipped": 0, "failed": 0 }, "hits": { "total": { "value": 1, "relation": "eq" }, "max_score": 1.0, "hits": [ { "_index": "product_****", "_type": "_doc", "_id": "TX3TZHgBfHNEDGoZ****", "_score": 1.0, "_source": { "msg_body": { "key": "test", "offset": 2, "overflowFlag": false, "partition": 2, "timestamp": 1616599282417, "topic": "dv****", "value": "test1", "valueSize": 8 }, "doc_as_upsert": true } } ] } }
Dépannage
| Symptôme | Cause possible | Solution |
|---|---|---|
| Le connecteur échoue à écrire dans Elasticsearch | Function Compute ne peut pas atteindre le cluster Elasticsearch | Vérifiez que le service Function Compute utilise le même VPC et le même vSwitch que le cluster Elasticsearch. Consultez la section Configurer le service Function Compute. |
| Les messages ne sont pas consommés | Paramètre de décalage du consommateur incorrect | Vérifiez le paramètre Consumer Offset. Utilisez l'option Earliest Offset pour lire tous les messages existants, ou l'option Latest Offset uniquement pour les nouveaux messages. |
| Erreurs d'authentification | Identifiants invalides ou autorisations insuffisantes | Confirmez que le nom d'utilisateur et le mot de passe sont corrects et que le compte dispose d'autorisations d'écriture sur l'index cible. |
Pour consulter les journaux d'appels de fonctions Function Compute, reportez-vous à la section Configurer la fonctionnalité de journalisation.