Articles de blog

Comment ingérer des données dans Elasticsearch via Kafka ?

Un guide pas à pas pour intégrer Apache Kafka avec Elasticsearch pour une ingestion, une indexation et une visualisation efficaces des données en utilisant Python, Docker Compose et Kafka Connect.

Dans cet article, nous montrons comment intégrer Apache Kafka avec Elasticsearch pour l'ingestion et l'indexation des données. Nous donnerons un aperçu de Kafka, de son concept de producteurs et de consommateurs, et nous créerons un index de logs où les messages seront reçus et indexés par Apache Kafka. Le projet est mis en œuvre en Python, et le code est disponible sur GitHub.

Produits requis

  • Docker et Docker Compose : Assurez-vous que Docker et Docker Compose sont installés sur votre machine.

  • Python 3.x : Pour exécuter les scripts du producteur et du consommateur.

Introduction à Apache Kafka

Apache Kafka est une plateforme de diffusion en continu distribuée qui permet une évolutivité et une disponibilité élevées, ainsi qu'une tolérance aux pannes. Dans Kafka, la gestion des données s'effectue par le biais des principaux composants :

  • Courtier: responsable du stockage et de la distribution des messages entre les producteurs et les consommateurs.

  • Zookeeper: gère et coordonne les courtiers Kafka, en contrôlant l'état de la grappe, les chefs de partition et les informations sur les consommateurs.

  • Sujets: canaux où les données sont publiées et stockées pour être consommées.

  • Consommateurs et producteurs: tandis que les producteurs envoient des données aux thèmes, les consommateurs récupèrent ces données.

Diagramme Apache Kafka

Ces composants fonctionnent ensemble pour former l'écosystème Kafka, qui fournit un cadre robuste pour la diffusion de données en continu.

Structure du projet

Pour comprendre le processus d'ingestion des données, nous l'avons divisé en plusieurs étapes :

  • Provisionnement de l'infrastructure: mise en place de l'environnement Docker pour prendre en charge Kafka, Elasticsearch et Kibana.

  • Création du producteur: mise en œuvre du producteur Kafka, qui envoie des données au sujet des journaux.

  • Création du consommateur: développement du consommateur Kafka pour lire et indexer les messages dans Elasticsearch.

  • Validation de l'ingestion: vérification et validation des données envoyées et consommées.

Configuration de l'infrastructure avec Docker Compose

Nous avons utilisé Docker Compose pour configurer et gérer les services nécessaires. Vous trouverez ci-dessous le code Docker Compose qui met en place chaque service nécessaire à l'intégration d'Apache Kafka, Elasticsearch et Kibana, en assurant un processus d'ingestion des données.

version: "3"

services:

  zookeeper:
    image: confluentinc/cp-zookeeper:latest
    container_name: zookeeper
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181

  kafka:
    image: confluentinc/cp-kafka:latest
    container_name: kafka
    depends_on:
      - zookeeper
    ports:
      - "9092:9092"
      - "9094:9094"
    environment:
      KAFKA_BROKER_ID: 1
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:29092,PLAINTEXT_HOST:${HOST_IP}:9092
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
      KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1

  elasticsearch:
    image: docker.elastic.co/elasticsearch/elasticsearch:8.15.1
    container_name: elasticsearch-8.15.1
    environment:
      - node.name=elasticsearch
      - xpack.security.enabled=false
      - discovery.type=single-node
      - "ES_JAVA_OPTS=-Xms512m -Xmx512m"
    volumes:
      - ./elasticsearch:/usr/share/elasticsearch/data
    ports:
      - 9200:9200

  kibana:
    image: docker.elastic.co/kibana/kibana:8.15.1
    container_name: kibana-8.15.1
    ports:
      - 5601:5601
    environment:
      ELASTICSEARCH_URL: http://elasticsearch:9200
      ELASTICSEARCH_HOSTS: '["http://elasticsearch:9200"]'

Vous pouvez accéder au fichier directement depuis le repo GitHub d'Elasticsearch Labs.

Envoi de données avec le producteur Kafka

Le producteur est responsable de l'envoi des messages au sujet des journaux. L'envoi de messages par lots augmente l'efficacité de l'utilisation du réseau et permet d'optimiser les paramètres batch_size et linger_ms, qui contrôlent respectivement la quantité et la latence des lots. La configuration acks='all' garantit que les messages sont stockés durablement, ce qui est essentiel pour les données d'enregistrement importantes.

producer = KafkaProducer(
   bootstrap_servers=['localhost:9092'],  # Specifies the Kafka server to connect
   value_serializer=lambda x: json.dumps(x).encode('utf-8'),  # Serializes data as JSON and encodes it to UTF-8 before sending
   batch_size=16384,     # Sets the maximum batch size in bytes (here, 16 KB) for buffered messages before sending
   linger_ms=10,         # Sets the maximum delay (in milliseconds) before sending the batch
   acks='all'            # Specifies acknowledgment level; 'all' ensures message durability by waiting for all replicas to acknowledge
)


def generate_log_message():
   levels = ["INFO", "WARNING", "ERROR", "DEBUG"]
   messages = [
       "User login successful",
       "User login failed",
       "Database connection established",
       "Database connection failed",
       "Service started",
       "Service stopped",
       "Payment processed",
       "Payment failed"
   ]
   log_entry = {
       "level": random.choice(levels),
       "message": random.choice(messages),
       "timestamp": time.time()
   }
   return log_entry

def send_log_batches(topic, num_batches=5, batch_size=10):
   for i in range(num_batches):
       logger.info(f"Sending batch {i + 1}/{num_batches}")
       for _ in range(batch_size):
           log_message = generate_log_message()
           producer.send(topic, value=log_message)
       producer.flush()


if __name__ == "__main__":
   topic = "logs"
   send_log_batches(topic)
   producer.close()

Lors du démarrage du producteur, les messages sont envoyés par lots au sujet, comme indiqué ci-dessous :

INFO:kafka.conn:Set configuration …
INFO:log_producer:Sending batch 1/5 
INFO:log_producer:Sending batch 2/5
INFO:log_producer:Sending batch 3/5
INFO:log_producer:Sending batch 4/5

Consommation et indexation des données avec le consommateur Kafka

Le consommateur est conçu pour traiter efficacement les messages, en consommant des lots à partir du sujet des journaux et en les indexant dans Elasticsearch. Avec auto_offset_reset='latest', il s'assure que le consommateur commence à traiter les messages les plus récents, en ignorant les plus anciens, et max_poll_records=10 limite le lot à 10 messages. Avec fetch_max_wait_ms=2000, le consommateur attend jusqu'à 2 secondes pour accumuler suffisamment de messages avant de traiter le lot.

Dans sa boucle principale, le consommateur consomme les messages du journal, traite et indexe chaque lot dans Elasticsearch, assurant ainsi une ingestion continue des données.

consumer = KafkaConsumer(
   'logs',                               
   bootstrap_servers=['localhost:9092'],
   auto_offset_reset='latest',            # Ensures reading from the latest offset if the group has no offset stored
   enable_auto_commit=True,               # Automatically commits the offset after processing
   group_id='log_consumer_group',         # Specifies the consumer group to manage offset tracking
   max_poll_records=10,                   # Maximum number of messages per batch
   fetch_max_wait_ms=2000                 # Maximum wait time to form a batch (in ms)
)

def create_bulk_actions(logs):
   for log in logs:
       yield {
           "_index": "logs",
           "_source": {
               'level': log['level'],
               'message': log['message'],
               'timestamp': log['timestamp']
           }
       }

if __name__ == "__main__":
   try:
       print("Starting message processing…")
       while True:

           messages = consumer.poll(timeout_ms=1000)  # Poll receive messages

           # process each batch messages
           for _, records in messages.items():
               logs = [json.loads(record.value) for record in records]
               bulk_actions = create_bulk_actions(logs)
               response = helpers.bulk(es, bulk_actions)
               print(f"Indexed {response[0]} logs.")
   except Exception as e:
       print(f"Erro: {e}")
   finally:
       consumer.close()
       print(f"Finish")

Visualisation des données dans Kibana

Avec Kibana, nous pouvons explorer et valider les données ingérées depuis Kafka et indexées dans Elasticsearch. En accédant à Dev Tools dans Kibana, vous pouvez visualiser les messages indexés et confirmer que les données sont conformes aux attentes. Par exemple, si notre producteur Kafka a envoyé 5 lots de 10 messages chacun, nous devrions voir un total de 50 enregistrements dans l'index.

Pour vérifier les données, vous pouvez utiliser la requête suivante dans la section Outils de développement :

GET /logs/_search
{
  "query": {
    "match_all": {}
  }
}

Réponse :

Réponse vérifier les données - Kafka & Elasticsearch​

En outre, Kibana permet de créer des visualisations et des tableaux de bord qui peuvent rendre l'analyse plus intuitive et interactive. Vous trouverez ci-dessous quelques exemples de tableaux de bord et de visualisations que nous avons créés, qui illustrent les données sous différents formats, améliorant ainsi notre compréhension des informations traitées.

Visualisation Kibana - Kafka & Elasticsearch​

Ingestion de données avec Kafka Connect

Kafka Connect est un service conçu pour faciliter l'intégration entre les sources de données et les destinations (puits), telles que les bases de données ou les systèmes de fichiers. Il fonctionne avec des connecteurs prédéfinis qui gèrent automatiquement les mouvements de données. Dans notre cas, Elasticsearch fait office de puits de données.

Ingestion de données avec Kafka Connect

En utilisant Kafka Connect, nous pouvons simplifier le processus d'ingestion de données, en éliminant la nécessité d'implémenter manuellement le workflow d'ingestion de données dans Elasticsearch. Avec le connecteur approprié, Kafka Connect permet aux données envoyées à un sujet Kafka d'être directement indexées dans Elasticsearch avec une configuration minimale et sans codage supplémentaire.

Travailler avec Kafka Connect

Pour mettre en œuvre Kafka Connect, nous allons ajouter le service kafka-connect à notre installation Docker Compose. Un élément clé de cette configuration est l'installation du connecteur Elasticsearch, qui se chargera de l'indexation des données.

Après avoir configuré le service et créé le conteneur Kafka Connect, un fichier de configuration pour le connecteur Elasticsearch sera nécessaire. Ce fichier définit des paramètres essentiels tels que

  • connection.url: URL de connexion pour Elasticsearch.

  • topics: Le sujet Kafka que le connecteur surveillera (dans ce cas, "logs").

  • type.name: Type de document dans Elasticsearch (typiquement _doc).

  • value.converter: Convertit les messages Kafka au format JSON.

  • value.converter.schemas.enable: Spécifie si le schéma doit être inclus.

  • schema.ignore et key.ignore: Paramètres permettant d'ignorer les schémas et les clés Kafka lors de l'indexation.

Voici la commande curl pour créer le connecteur Elasticsearch dans Kafka Connect :

curl --location '{{url}}/connectors' \
--header 'Content-Type: application/json' \
--data '{
    "name": "elasticsearch-sink-connector",
    "config": {
        "connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector",
        "topics": "logs",
        "connection.url": "http://elasticsearch:9200",
        "type.name": "_doc",
        "value.converter": "org.apache.kafka.connect.json.JsonConverter",
        "value.converter.schemas.enable": "false",
        "schema.ignore": "true",
        "key.ignore": "true"
    }
}'

Avec cette configuration, Kafka Connect commencera automatiquement à ingérer les données envoyées au sujet "logs" et à les indexer dans Elasticsearch. Cette approche permet d'automatiser entièrement l'ingestion et l'indexation des données sans nécessiter de codage supplémentaire, ce qui simplifie l'ensemble du processus d'intégration.

Conclusion

L'intégration de Kafka et d'Elasticsearch crée un pipeline puissant pour l'ingestion et l'analyse de données en temps réel. Ce guide fournit une approche fondamentale pour construire une architecture d'ingestion de données robuste, avec une visualisation et une analyse transparentes dans Kibana, prête à s'adapter à des exigences plus complexes à l'avenir.

En outre, l'utilisation de Kafka Connect rend l'intégration entre Kafka et Elasticsearch encore plus rationnelle, en éliminant le besoin de code supplémentaire pour traiter et indexer les données. Kafka Connect permet aux données envoyées à un sujet spécifique d'être automatiquement indexées dans Elasticsearch avec une configuration minimale.

Pour aller plus loin

Affichage des champs dans un index Elasticsearch

Kofi Bartlett

Scripting Ruby dans Logstash

Dai Sugimori

Modèles d'index dans Elasticsearch : Comment utiliser les modèles composables

Kofi Bartlett

Comment afficher les champs d'un index Elasticsearch ?

JD Armada

Comment ingérer des données dans Elasticsearch via LlamaIndex

Andre Luiz

Prêt à créer des expériences de recherche d'exception ?

Une recherche suffisamment avancée ne se fait pas avec les efforts d'une seule personne. Elasticsearch est alimenté par des data scientists, des ML ops, des ingénieurs et bien d'autres qui sont tout aussi passionnés par la recherche que vous. Mettons-nous en relation et travaillons ensemble pour construire l'expérience de recherche magique qui vous permettra d'obtenir les résultats que vous souhaitez.

Jugez-en par vous-même