Tous les produits
Search
Centre de documentation

ApsaraMQ for Kafka:Create an Elasticsearch sink connector

Dernière mise à jour :Aug 11, 2026

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

Function Compute

Elasticsearch

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

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

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

  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 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****
Important

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ètreDescriptionExemple
VPC IDVPC de l'instance source. Rempli automatiquement ; aucune modification n'est nécessaire.vpc-bp1xpdnd3l***
vSwitch IDvSwitch de l'instance source. Doit se trouver dans le même VPC.vsw-bp1d2jgg81***
Failure Handling PolicyAction à 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 MethodL'option Auto crée automatiquement les rubriques internes requises. L'option Manual vous permet de les créer vous-même.Auto
Connector Consumer GroupGroupe 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 --
Remarque

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
Remarque
  • 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.

  1. Sur la page Connectors, recherchez le connecteur. Dans la colonne Actions, choisissez More > Configure Function.

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

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

  2. Dans le panneau Send Message, définissez l'option Method of Sending sur Console.

  3. Dans le champ Message Key, saisissez une clé, par exemple demo.

  4. Dans le champ Message Content, saisissez un corps JSON, par exemple :

       {"key": "test"}
  5. 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

  1. Connectez-vous à la console Kibana.

  2. Exécutez la requête suivante pour rechercher dans l'index cible :

       GET /<index_name>/_search
  3. 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.