Kafka 経由で Elasticsearch にデータを取り込む方法
Python、Docker Compose、Kafka Connect を使用して、Apache Kafka を Elasticsearch と統合し、効率的なデータの取り込み、インデックス作成、視覚化を行うためのステップバイステップ ガイドです。
この記事では、データの取り込みとインデックス作成のために Apache Kafka と Elasticsearch を統合する方法を説明します。Kafka の概要、プロデューサーとコンシューマーの概念について説明し、Apache Kafka を介してメッセージを受信してインデックスを作成するログ インデックスを作成します。このプロジェクトは Python で実装されており、コードはGitHubで入手できます。
要件
Docker と Docker Compose: マシンに Docker と Docker Compose がインストールされていることを確認します。
Python 3.x: Producer スクリプトと Consumer スクリプトを実行します。
Apache Kafka の紹介
Apache Kafka は、高いスケーラビリティと可用性、およびフォールト トレランスを実現する分散ストリーミング プラットフォームです。Kafka では、データ管理は次の主要コンポーネントを通じて行われます。
ブローカー: プロデューサーとコンシューマー間のメッセージの保存と配信を担当します。
Zookeeper : Kafka ブローカーを管理および調整し、クラスターの状態、パーティション リーダー、およびコンシューマー情報を制御します。
トピック: データが公開され、消費のために保存されるチャネル。
コンシューマーとプロデューサー: プロデューサーはトピックにデータを送信し、コンシューマーはそのデータを取得します。

これらのコンポーネントは連携して Kafka エコシステムを形成し、データ ストリーミングのための堅牢なフレームワークを提供します。
プロジェクト構造
データ取り込みプロセスを理解するために、次の段階に分けました。
インフラストラクチャのプロビジョニング: Kafka、Elasticsearch、Kibana をサポートするための Docker 環境をセットアップします。
プロデューサーの作成: ログ トピックにデータを送信する Kafka プロデューサーを実装します。
コンシューマーの作成: Elasticsearch でメッセージを読み取ってインデックスを作成する Kafka コンシューマーを開発します。
取り込み検証: 送信および消費されたデータを検証および検証します。
Docker Compose によるインフラストラクチャ構成
必要なサービスを構成および管理するために Docker Compose を利用しました。以下に、Apache Kafka、Elasticsearch、Kibana の統合に必要な各サービスを設定し、データ取り込みプロセスを確実に実行する Docker Compose コードを示します。
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"]'Elasticsearch Labs GitHubリポジトリから直接ファイルにアクセスできます。
Kafka Producerによるデータ送信
プロデューサーは、ログ トピックにメッセージを送信する責任を負います。メッセージをバッチで送信することで、ネットワークの使用効率が向上し、バッチの量と待ち時間をそれぞれ制御するbatch_sizeおよびlinger_ms設定による最適化が可能になります。構成acks='all'により、メッセージが永続的に保存されることが保証されます。これは重要なログ データにとって不可欠です。
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()プロデューサーを起動すると、次に示すように、メッセージがトピックにバッチで送信されます。
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/5Kafka Consumerによるデータの消費とインデックス作成
コンシューマーは、ログ トピックからバッチを消費し、Elasticsearch にインデックスを付けることで、メッセージを効率的に処理するように設計されています。auto_offset_reset='latest'使用すると、コンシューマーが最新のメッセージの処理を開始し、古いメッセージを無視することが保証され、 max_poll_records=10バッチを 10 件のメッセージに制限します。fetch_max_wait_ms=2000の場合、コンシューマーはバッチを処理する前に十分なメッセージを蓄積するために最大 2 秒間待機します。
メイン ループでは、コンシューマーはログ メッセージを消費し、各バッチを処理して Elasticsearch にインデックス付けし、継続的なデータ取り込みを保証します。
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")Kibanaでデータを視覚化する
Kibana を使用すると、Kafka から取り込まれ Elasticsearch でインデックス化されたデータを探索および検証できます。Kibana のDev Toolsにアクセスすると、インデックス付けされたメッセージを表示し、データが期待どおりであることを確認できます。たとえば、Kafka プロデューサーが 10 件のメッセージを 5 つのバッチで送信した場合、インデックスには合計 50 件のレコードが表示されます。
データを検証するには、開発ツールセクションで次のクエリを使用できます。
GET /logs/_search
{
"query": {
"match_all": {}
}
}対応:

さらに、Kibana は、分析をより直感的かつインタラクティブにするのに役立つ視覚化とダッシュボードを作成する機能を提供します。以下に、私たちが作成したダッシュボードと視覚化の例をいくつか示します。これらは、さまざまな形式でデータを示し、処理された情報に対する理解を深めます。

Kafka Connectによるデータ取り込み
Kafka Connect は、データベースやファイル システムなどのデータ ソースと宛先 (シンク) 間の統合を容易にするために設計されたサービスです。データの移動を自動的に処理する定義済みのコネクタを使用して動作します。私たちの場合、Elasticsearch はデータ シンクとして機能します。

Kafka Connect を使用すると、データ取り込みプロセスを簡素化できるため、Elasticsearch にデータ取り込みワークフローを手動で実装する必要がなくなります。適切なコネクタを使用すると、Kafka Connect では、追加のコーディングを必要とせず、最小限のセットアップで、Kafka トピックに送信されたデータを Elasticsearch で直接インデックス化できます。
Kafka Connect の操作
Kafka Connect を実装するには、Docker Compose セットアップにkafka-connect サービスを追加します。この構成の重要な部分は、データのインデックス作成を処理する Elasticsearch コネクタをインストールすることです。
サービスを設定し、Kafka Connect コンテナを作成した後、Elasticsearch コネクタの設定ファイルが必要になります。このファイルは次のような重要なパラメータを定義します。
connection.url: Elasticsearch の接続 URL。topics: コネクタが監視する Kafka トピック (この場合は「ログ」)。type.name: Elasticsearch のドキュメント タイプ (通常は _doc)。value.converter: Kafka メッセージを JSON 形式に変換します。value.converter.schemas.enable: スキーマを含めるかどうかを指定します。schema.ignoreおよびkey.ignore: インデックス作成中に Kafka スキーマとキーを無視する設定。
以下は、Kafka Connect で Elasticsearch コネクタを作成するためのcurlコマンドです。
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"
}
}'この構成により、Kafka Connect は「logs」トピックに送信されたデータを自動的に取り込み、Elasticsearch でインデックス作成を開始します。このアプローチにより、追加のコーディングを必要とせずに完全に自動化されたデータの取り込みとインデックス作成が可能になり、統合プロセス全体が簡素化されます。
まとめ
Kafka と Elasticsearch を統合すると、リアルタイムのデータの取り込みと分析のための強力なパイプラインが作成されます。このガイドでは、Kibana でのシームレスな視覚化と分析を備えた堅牢なデータ取り込みアーキテクチャを構築するための基本的なアプローチを提供し、将来のより複雑な要件に適応できるようにします。
さらに、Kafka Connect を使用すると、Kafka と Elasticsearch の統合がさらに効率化され、データの処理やインデックス作成のための追加コードが不要になります。Kafka Connect を使用すると、最小限の構成で特定のトピックに送信されたデータを自動的に Elasticsearch にインデックス付けできます。




