Blog

Wie man Daten über Kafka in Elasticsearch einliest

Eine Schritt-für-Schritt-Anleitung zur Integration von Apache Kafka mit Elasticsearch für effiziente Datenerfassung, -indizierung und -visualisierung mit Python, Docker Compose und Kafka Connect.

In diesem Artikel zeigen wir, wie man Apache Kafka mit Elasticsearch zur Datenerfassung und -indizierung integriert. Wir geben einen Überblick über Kafka, sein Konzept von Produzenten und Konsumenten und erstellen einen Log-Index, in dem Nachrichten über Apache Kafka empfangen und indexiert werden. Das Projekt wurde in Python implementiert, und der Code ist auf GitHub verfügbar.

Voraussetzungen

  • Docker und Docker Compose: Stellen Sie sicher, dass Docker und Docker Compose auf Ihrem Rechner installiert sind.

  • Python 3.x: Zum Ausführen der Producer- und Consumer-Skripte.

Einführung in Apache Kafka

Apache Kafka ist eine verteilte Streaming-Plattform, die hohe Skalierbarkeit und Verfügbarkeit sowie Fehlertoleranz ermöglicht. In Kafka erfolgt die Datenverwaltung über die Hauptkomponenten:

  • Broker: zuständig für die Speicherung und Verteilung von Nachrichten zwischen Produzenten und Konsumenten.

  • Zookeeper: verwaltet und koordiniert die Kafka-Broker und kontrolliert den Zustand des Clusters, die Partitionsleiter und die Verbraucherinformationen.

  • Themen: Kanäle, auf denen Daten veröffentlicht und zur Nutzung gespeichert werden.

  • Konsumenten und Produzenten: Während Produzenten Daten an die Themen senden, rufen Konsumenten diese Daten ab.

Diagramm Apache Kafka

Diese Komponenten arbeiten zusammen und bilden das Kafka-Ökosystem, das ein robustes Framework für das Datenstreaming bietet.

Projektstruktur

Um den Datenerfassungsprozess zu verstehen, haben wir ihn in Phasen unterteilt:

  • Infrastrukturbereitstellung: Einrichtung der Docker-Umgebung zur Unterstützung von Kafka, Elasticsearch und Kibana.

  • Producer-Erstellung: Implementierung des Kafka-Producers, der Daten an das Logs-Topic sendet.

  • Consumer Creation: Entwicklung des Kafka-Consumers zum Lesen und Indizieren von Nachrichten in Elasticsearch.

  • Aufnahmevalidierung: Überprüfung und Validierung der gesendeten und empfangenen Daten.

Infrastrukturkonfiguration mit Docker Compose

Wir haben Docker Compose verwendet, um die notwendigen Dienste zu konfigurieren und zu verwalten. Nachfolgend finden Sie den Docker Compose-Code, der die für die Integration von Apache Kafka, Elasticsearch und Kibana erforderlichen Dienste einrichtet und so einen Datenaufnahmeprozess sicherstellt.

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"]'

Sie können direkt über das Elasticsearch Labs GitHub- Repository auf die Datei zugreifen.

Datenübertragung mit dem Kafka Producer

Der Produzent ist für das Senden von Nachrichten an das Log-Thema verantwortlich. Durch das Senden von Nachrichten in Stapeln wird die Netzwerknutzungseffizienz gesteigert. Optimierungen sind mit den Einstellungen batch_size und linger_ms möglich, die die Anzahl bzw. die Latenz der Stapel steuern. Die Konfiguration acks='all' gewährleistet, dass Nachrichten dauerhaft gespeichert werden, was für wichtige Protokolldaten unerlässlich ist.

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()

Beim Starten des Producers werden die Nachrichten in Batches an das Topic gesendet, wie unten dargestellt:

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

Konsum und Indizierung von Daten mit dem Kafka Consumer

Der Consumer ist so konzipiert, dass er Nachrichten effizient verarbeitet, indem er Batches aus dem Logs-Topic empfängt und diese in Elasticsearch indexiert. Mit auto_offset_reset='latest' wird sichergestellt, dass der Consumer mit der Verarbeitung der neuesten Nachrichten beginnt und die älteren ignoriert, und max_poll_records=10 begrenzt den Batch auf 10 Nachrichten. Bei fetch_max_wait_ms=2000 wartet der Konsument bis zu 2 Sekunden, um genügend Nachrichten zu sammeln, bevor er den Batch verarbeitet.

Im Hauptzyklus verarbeitet der Consumer die Logmeldungen, speichert sie in Elasticsearch und indexiert sie, um eine kontinuierliche Datenaufnahme zu gewährleisten.

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")

Datenvisualisierung in Kibana

Mit Kibana können wir die von Kafka aufgenommenen und in Elasticsearch indexierten Daten untersuchen und validieren. Durch den Zugriff auf die Entwicklertools in Kibana können Sie die indizierten Nachrichten anzeigen und überprüfen, ob die Daten den Erwartungen entsprechen. Wenn beispielsweise unser Kafka-Producer 5 Batches mit jeweils 10 Nachrichten sendet, sollten wir insgesamt 50 Datensätze im Index sehen.

Um die Daten zu überprüfen, können Sie die folgende Abfrage im Abschnitt „Entwicklertools“ verwenden:

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

Abwehr:

Antwort zur Datenprüfung – Kafka & Elasticsearch​

Darüber hinaus bietet Kibana die Möglichkeit, Visualisierungen und Dashboards zu erstellen, die die Analyse intuitiver und interaktiver gestalten können. Nachfolgend sehen Sie einige Beispiele der von uns erstellten Dashboards und Visualisierungen, die die Daten in verschiedenen Formaten veranschaulichen und so unser Verständnis der verarbeiteten Informationen verbessern.

Kibana-Visualisierung – Kafka & Elasticsearch​

Datenaufnahme mit Kafka Connect

Kafka Connect ist ein Dienst, der die Integration zwischen Datenquellen und Zielen (Senken) wie Datenbanken oder Dateisystemen erleichtert. Es arbeitet mit vordefinierten Konnektoren, die den Datentransfer automatisch abwickeln. In unserem Fall fungiert Elasticsearch als Datensenke.

Datenerfassung mit Kafka Connect

Mit Kafka Connect können wir den Datenaufnahmeprozess vereinfachen und die Notwendigkeit eliminieren, den Datenaufnahme-Workflow manuell in Elasticsearch zu implementieren. Mit dem passenden Konnektor ermöglicht Kafka Connect die direkte Indizierung von Daten, die an ein Kafka-Topic gesendet werden, in Elasticsearch – mit minimalem Einrichtungsaufwand und ohne zusätzliche Programmierung.

Arbeiten mit Kafka Connect

Um Kafka Connect zu implementieren, fügen wir den kafka-connect-Dienst zu unserem Docker Compose-Setup hinzu. Ein wichtiger Bestandteil dieser Konfiguration ist die Installation des Elasticsearch-Connectors, der für die Datenindizierung zuständig ist.

Nach der Konfiguration des Dienstes und der Erstellung des Kafka Connect-Containers wird eine Konfigurationsdatei für den Elasticsearch-Connector benötigt. Diese Datei definiert wichtige Parameter wie zum Beispiel:

  • connection.url: Verbindungs-URL für Elasticsearch.

  • topics: Das Kafka-Thema, das der Connector überwachen wird (in diesem Fall "logs").

  • type.name: Dokumenttyp in Elasticsearch (typischerweise _doc).

  • value.converter: Konvertiert Kafka-Nachrichten in das JSON-Format.

  • value.converter.schemas.enable: Gibt an, ob das Schema einbezogen werden soll.

  • schema.ignore und key.ignore: Einstellungen zum Ignorieren von Kafka-Schemas und -Schlüsseln während der Indizierung.

Nachfolgend der Befehl curl zum Erstellen des Elasticsearch-Connectors in 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"
    }
}'

Mit dieser Konfiguration beginnt Kafka Connect automatisch mit der Erfassung der an das Topic "logs" gesendeten Daten und deren Indizierung in Elasticsearch. Dieser Ansatz ermöglicht die vollautomatische Datenerfassung und -indizierung ohne zusätzlichen Programmieraufwand und vereinfacht so den gesamten Integrationsprozess.

Fazit

Durch die Integration von Kafka und Elasticsearch entsteht eine leistungsstarke Pipeline für die Datenerfassung und -analyse in Echtzeit. Dieser Leitfaden bietet einen grundlegenden Ansatz für den Aufbau einer robusten Datenerfassungsarchitektur mit nahtloser Visualisierung und Analyse in Kibana, die sich zukünftig an komplexere Anforderungen anpassen lässt.

Darüber hinaus vereinfacht die Verwendung von Kafka Connect die Integration zwischen Kafka und Elasticsearch zusätzlich, wodurch die Notwendigkeit von zusätzlichem Code zur Verarbeitung und Indizierung von Daten entfällt. Kafka Connect ermöglicht es, Daten, die an ein bestimmtes Thema gesendet werden, mit minimalem Konfigurationsaufwand automatisch in Elasticsearch zu indizieren.

Zugehörige Inhalte

Anzeigen von Feldern in einem Elasticsearch-Index

Kofi Bartlett

Ruby-Skripting in Logstash

Dai Sugimori

Indexvorlagen in Elasticsearch: Wie man zusammensetzbare Vorlagen verwendet

Kofi Bartlett

Wie man Felder eines Elasticsearch-Index anzeigt

JD Armada

Wie man Daten über LlamaIndex in Elasticsearch einliest

Andre Luiz

Sind Sie bereit, hochmoderne Sucherlebnisse zu schaffen?

Eine ausreichend fortgeschrittene Suche kann nicht durch die Bemühungen einer einzelnen Person erreicht werden. Elasticsearch wird von Datenwissenschaftlern, ML-Ops-Experten, Ingenieuren und vielen anderen unterstützt, die genauso leidenschaftlich an der Suche interessiert sind wie Sie. Lasst uns in Kontakt treten und zusammenarbeiten, um das magische Sucherlebnis zu schaffen, das Ihnen die gewünschten Ergebnisse liefert.

Probieren Sie es selbst aus