Exportez les données d'une rubrique source de votre instance Message Queue for Apache Kafka vers Alibaba Cloud Elasticsearch.
Prérequis
Pour plus d'informations, consultez la section Prérequis.
Étape 1 : Créer les ressources du service cible
Créez une instance et un index dans la console Elasticsearch. Pour en savoir plus, reportez-vous à la section Prise en main.
Ajoutez le bloc CIDR du service Function Compute à la liste d'autorisation de votre instance Elasticsearch. Cette opération est requise uniquement pour les connexions VPC. Pour plus de détails, consultez la page Configurer une liste d'autorisation d'adresses IP publiques ou privées pour une instance.
Étape 2 : Créer un connecteur de réception Elasticsearch
Connectez-vous à la console ApsaraMQ for Kafka. Sur la page Overview, sélectionnez une région dans la section Resource Distribution.
Dans le volet de navigation de gauche, choisissez .
Sur la page Tasks, cliquez sur Create Task.
-
Sur la page Create Task, configurez les paramètres Task Name et Description. Suivez ensuite les instructions à l'écran pour configurer les autres paramètres.
-
Configuration de la tâche
-
Dans l'assistant de configuration Source, définissez le paramètre Data Provider sur ApsaraMQ for Kafka, configurez les paramètres suivants, puis cliquez sur Next.
Parameter
Description
Example
Region
La région où réside l'instance ApsaraMQ for Kafka.
China (Hangzhou)
ApsaraMQ for Kafka Instance
L'ID de l'instance ApsaraMQ for Kafka dans laquelle sont produites les données à acheminer.
alikafka_post-cn-9hdsbdhd****
Topic
La rubrique de l'instance ApsaraMQ for Kafka dans laquelle sont produites les données à acheminer.
guide-sink-topic
Group ID
L'ID du groupe de l'instance ApsaraMQ for Kafka dans laquelle sont produites les données à acheminer.
Quickly Create : Le système crée automatiquement un groupe dont l'ID est au format GID_EVENTBRIDGE_xxx.
Use Existing Group : Sélectionnez l'ID d'un groupe existant inutilisé. Si vous sélectionnez un groupe déjà utilisé, la publication et la consommation des messages existants seront affectées.
Use Existing Group
Consumer Offset
Latest Offset : Consomme les messages à partir du dernier offset.
Earliest Offset : Consomme les messages à partir du premier offset.
Latest Offset
Network Configuration
Si une transmission de données transfrontalière est requise, sélectionnez Self-managed Internet. Dans les autres cas, sélectionnez Basic Network.
Basic Network
Data Format
Le format de données permet d'encoder les données binaires envoyées depuis la source dans un format spécifique. Plusieurs formats sont pris en charge. Si vous n'avez pas d'exigences particulières en matière d'encodage, spécifiez Json comme valeur.
Json : encode les données binaires au format JSON basé sur UTF-8 et les place dans la charge utile.
Text : format par défaut. Encode les données binaires en une chaîne basée sur UTF-8 et les place dans la charge utile.
Binary : encode les données binaires en chaînes Base64, puis les place dans la charge utile.
Text
Messages
Paramètre de Advanced Configuration. Nombre maximal de messages à envoyer par lot. Une requête est envoyée uniquement lorsque le nombre de messages accumulés atteint cette valeur. La valeur doit être un entier compris entre 1 et 10 000.
2000
Interval (Unit: Seconds)
Paramètre de Advanced Configuration. Intervalle d'appel de la fonction. Le système agrège les messages et les envoie à Function Compute à cet intervalle. La valeur doit être un entier compris entre 0 et 15. L'unité est la seconde. Une valeur de 0 signifie que les messages sont livrés immédiatement.
3
À l'étape Filtering, définissez un modèle d'événement pour filtrer les requêtes. Pour plus d'informations, consultez la section Modèles d'événements.
À l'étape Transformation, configurez le nettoyage des données pour un traitement complexe, tel que la division, le mappage, l'enrichissement et le routage dynamique. Pour en savoir plus, reportez-vous à la section Utiliser Function Compute pour nettoyer les données des messages.
-
À l'étape Sink, sélectionnez Alibaba Cloud Elasticsearch acs.elasticSearch pour le paramètre Service Type et configurez les paramètres suivants.
Parameter
Description
Example
Elasticsearch Cluster
L'instance Elasticsearch que vous avez créée.
es-cn-pe336j0gj001e****
Cluster Logon Name
Le nom de connexion de l'instance, qui est
elasticpar défaut.elastic
Instance logon password
Le mot de passe configuré lors de la création de l'instance.
Index Name
Le nom de l'index que vous avez créé. Pour plus d'informations, consultez la section Prise en main. La valeur peut être une constante de chaîne ou une variable JSONPath, telle que
product_infoou$.data.key.product_info
Document Type
Le type du document de données. La valeur peut être une constante de chaîne ou une variable JSONPath.
Exemples :
_docou$.data.key.RemarqueCe paramètre ne peut être configuré que pour les versions d'instance Elasticsearch antérieures à la version 7.0. La valeur par défaut est la constante
_doc._doc
Document
Indiquez si vous souhaitez livrer l'événement complet ou partiel à Elasticsearch. Pour les événements partiels, vous devez configurer une règle d'extraction JSONPath.
Complete Event
Network Configuration
-
VPC : Livrez les messages Kafka à Elasticsearch via un VPC.
-
Public Network : Livrez les messages Kafka à Elasticsearch via Internet.
Public Network
VPC
Le VPC contenant l'instance Elasticsearch. Ce paramètre est requis uniquement lorsque le paramètre Network Configuration est défini sur VPC.
vpc-bp17fapfdj0dwzjkd****
vSwitch
Le vSwitch auquel appartient l'instance Elasticsearch. Ce paramètre est requis uniquement lorsque le paramètre Network Configuration est défini sur VPC.
vsw-bp1gbjhj53hdjdkg****
Security Group
Le groupe de sécurité. Ce paramètre est requis uniquement lorsque le paramètre Network Configuration est défini sur VPC.
test_group
-
-
-
Propriétés de la tâche
Configurez la stratégie de nouvelle tentative pour les livraisons d'événements ayant échoué et la méthode de gestion des erreurs. Pour plus d'informations, consultez la section Nouvelles tentatives et files d'attente de lettres mortes.
-
Une fois la configuration terminée, cliquez sur Save. Sur la page Tasks, localisez la tâche de connecteur de réception Elasticsearch que vous avez créée. La colonne Status affiche Starting. Lorsque le statut passe à Running, le connecteur est créé et prêt à l'emploi.
Étape 3 : Tester le connecteur de réception Elasticsearch
Sur la page Tasks, cliquez sur la rubrique source dans la colonne Event Source de la tâche de connecteur de réception Elasticsearch.
Sur la page des détails de la rubrique, cliquez sur Send Test Message.
-
Dans le panneau Start to Send and Consume Message, configurez le message comme suit, puis cliquez sur OK.
Définissez la méthode d'envoi sur Console. Définissez Message Key sur
es-sink-k1et Message Content sur{"esk1":1,"esk2":"v2"}. Pour Send to Specified Partition, sélectionnez No. Connectez-vous à la console Elasticsearch et accédez à l'instance via Kibana. Pour plus d'informations, consultez la section Prise en main.
-
Sur la console Kibana, exécutez la commande suivante pour afficher le résultat de l'insertion des données.
GET /your-index-name/_searchLa requête renvoie un statut 200 OK avec un document dont
_indexestproduct_infoet_idest1717558528. Le champ_sourcecontient les champstopic,partition,offset,timestamp,headers,keyetvalue, aveckeydéfini sures-sink-k1etvaluedéfini sur{"esk1": 1, "esk2": "v2"}. Cela confirme que les données ont été écrites avec succès dans Elasticsearch.