ラベル kafka の投稿を表示しています。 すべての投稿を表示
ラベル kafka の投稿を表示しています。 すべての投稿を表示

2021年1月9日土曜日

fluentd -> kafka -> logstash -> elasticsearch の順番でログを格納する方法

概要

fluentd -> kafka -> logstash -> elasticsearch という順番でログを確認する方法を紹介します
fluentd, kafka, elasticsearch の構築方法は過去に紹介しているので別記事を参照しています

環境

  • kafka 2.6.0
  • logstash 7.10.1
  • elasticsearch 6.4.0

fluent-kafka-plugin が動作する環境の構築

こちらを参考に構築してください
kafka と fluent-kafka-plugin がインストールされた fluent コンテナが起動している状態になれば OK です

elasticsearch の構築

こちらを参考に構築してください
9200 ポートで elasticsearch にアクセスできれば OK です

logstash のインストール

  • brew install logstash

logstash kafka input plugin のインストールと設定

ここが今回の肝になる部分です
logstash-integration-kafka を使います

  • logstash-plugin install logstash-integration-kafka

インストールが完了したら設定ファイルを作成していきます
kafka からデータを受け取って elasticsearch に流す設定を定義します

  • cp /usr/local/etc/logstash/logstash-sample.conf /usr/local/etc/logstash/logstash.conf
  • vim /usr/local/etc/logstash/logstash.conf
input {
  kafka {
    bootstrap_servers => "192.168.1.2:9092"
    topics => ["test"]
    decorate_events => true
  }
}

filter {
  json {
    source => "message"
  }
}

output {
  elasticsearch {
    hosts => ["http://192.168.1.2:9200"]
    index => "%{[@metadata][kafka][topic]}-%{+YYYY.MM.dd}"
  }
}


decorate_events を true にしないと [@metadata][kafka][topic] が output セクションで使えないので true にしています

また filter で json を使っています
どうやらデフォルトでは kafka から受け取ったログはすべて「message」というフィールドに格納されています
なのでそれぞれのフィールドに分割するために filter を挟んでいます

  • logstash -f /usr/local/etc/logstash/logstash.conf

起動に少し時間がかかりますが ERROR がでなければ OK です

動作確認

まずは適当なログを fluentd コンテナに送ります

  • docker run --rm --log-driver=fluentd --log-opt fluentd-address=192.168.1.2:24224 --log-opt tag="docker.{{.Name}}" alpine /bin/sh -c "while :;do echo \"{\\\"timestamp\\\":\\\"$(date)\\\",\\\"msg\\\":\\\"hello\\\"}\"; sleep 3; done;"


ある程度待った後に elasticsearch にインデックスが作成されているか確認しましょう

  • curl 'http://192.168.1.2:9200/_cat/indices?v'
  • curl 'http://192.168.1.2:9200/test-2021.01.06?pretty'

最後に

必ずしも kafka は必要ではないですがスケールによっては必要になります

参考サイト

2021年1月8日金曜日

ruby-kafka を使って kafka との証明書認証をサクっと確認する方法

概要

fluent-kafka-plugin で証明書を使って kafka に接続する場合にうまく接続できない場合があると思います
そんな場合は ruby-kafka を使って確認しましょう

概要

  • macOS 11.1
  • Ruby 3.0.0
  • ruby-kafka 1.3.0

準備

  • bundle init
  • vim Gemfile
gem "ruby-kafka"
  • bundle install

テストコード

  • vim test.rb
require 'kafka'

kafka = Kafka.new(
  ["kafka:9092"],
  ssl_ca_cert: File.read('./ca.pem'),
  ssl_client_cert: File.read('./cert.pem'),
  ssl_client_cert_key: File.read('./key.pem')
)
puts kafka

ret = kafka.deliver_message("Hello, World!", topic: "test")
puts ret


  • bundle exec ruby test.rb

これでエラーが出なければ OK です
証明書が間違っている場合などは OpenSSL のエラーなどが出ると思います

参考サイト

2020年12月14日月曜日

fluent-plugin-kafka を使ってみた

概要

fluentd から kafka にメッセージを送信することができるプラグインがあるので使ってみました
kafka の構築に関してはこちらを参考にしてください
また今回 fluentd はコンテナで動作させます

環境

  • macOS 10.15.7
  • fluentd
  • kafka 2.6.0

事前準備

kafka と zookeeper を起動させておきましょう

  • brew services start zookeeper
  • brew services start kafka

fluent-plugin-kafka がインストールされたイメージの作成

fluent/fluentd には fluent-plugin-kafka がインストールされていないのでインストールされているイメージを作成します

  • vim Dockerfile
FROM fluent/fluentd

RUN apk add --update --virtual .build-deps \
        sudo build-base ruby-dev \
 && sudo gem install \
        fluent-plugin-kafka zookeeper \
 && sudo gem sources --clear-all \
 && apk del .build-deps \
 && rm -rf /var/cache/apk/* \
           /home/fluent/.gem/ruby/*/cache/*.gem


  • docker build -t my_fluentd .

fluent.conf の作成

次に作成した fluentd イメージ上で動作させる設定ファイルを作成します
コンテナを作成する場合にホストマシン上のファイルをマウントして動作させます

今回はわかりやすいように copy を使って kafka にログを流すのと同時に fluentd コンテナの標準出力にもログを出しています

  • vim fluent.conf
<source>
  @type forward
  port 24224
  bind 0.0.0.0
</source>

<match docker.**>
  @type copy
  <store>
    @type stdout
  </store>
  <store>
    @type kafka2
    brokers 192.168.1.2:9092
    zookeeper 192.168.1.2:2181
    default_topic test
    <format>
      @type json
    </format>
  </store>
</match>

fluentd コンテナの起動

作成したイメージと設定ファイルを使ってコンテナ起動します
問題なくコンテナが起動していることを確認しましょう

  • docker run -d -p 24224:24224 -p 24224:24224/udp -v $(pwd):/fluentd/etc -e FLUENTD_CONF=fluent.conf my_fluentd

動作確認用コンテナの作成

何でも OK です
今回は JSON の情報を echo で 3 秒おきに出力するコンテナにしています
ロギングドライバだけ fluentd を指定しましょう

  • docker run --rm --log-driver=fluentd --log-opt fluentd-address=192.168.1.2:24224 --log-opt tag="docker.{{.Name}}" alpine /bin/sh -c "while :;do echo \"{\"timestamp\":\"$(date)\"}\"; sleep 3; done;"

動作確認

kafka-console-consumer を使って fluentd からログが飛んできているか確認しましょう
fluentd がデフォルトで 10s バッファするのでログが飛んでくるのは 10s おきになっているのが確認できると思います

  • kafka-console-consumer --bootstrap-server 192.168.1.2:9092 --topic test --from-beginning


また fluentd コンテナで logs を確認しても良いと思います
そもそも fluentd コンテナに出力されていない場合は kafka にも当然ログは飛んできません

トラブルシューティング

zookeeper が localhost でしか LISTEN していない場合には設定ファイルを編集しましょう

  • vim /usr/local/etc/kafka/zookeeper.properties
clientPort=2181
clientPortAddress=192.168.1.2

最後に

fluent-plugin-kafka を使ってみました
今回は Output プラグインを使いましたが Input プラグインもあり kafka からの入力を受け取ることもできます

kafka がすでにあればかなり簡単に使える印象です

参考サイト

2020年12月12日土曜日

MacOS 上で kafka をインストールして使ってみる

概要

kafka はストリームでメッセージのやり取りをするための基盤です
今回は MacOS 上で簡単に動かす方法を紹介します

環境

  • macOS 10.15.7
  • kafka 2.6.0
  • zookeeper 3.6.2

kafka/zookeeper インストール

  • brew install kafka

一緒に zookeeper もインストールされます

kafka/zookeeper 起動

  • brew services start zookeeper
  • brew services start kafka

zookeeper は localhost:2181, kafka は localhost:9092 で起動します

トピックの作成

  • kafka-topics --create --zookeeper localhost:2181 --replication-factor 1 --partitions 1 --topic test

    すべて必須のパラメータになります

  • --create でトピックの作成

  • --zookeeper localhost:2181 で zookeeper の LISTEN アドレスの指定
  • --replication-factor 1 でレプリカを 1 つのみ生成
  • --partitions 1 でパーティションを 1 つのみ生成
  • --topic でトピック名を指定

replication-factor と partitions についてはこのあたりがイメージしやすいかなと思います

トピック確認

  • kafka-topics --list --zookeeper localhost:2181

    => test

  • kafka-topics --describe --zookeeper localhost:2181 --topic test

Topic: test PartitionCount: 1 ReplicationFactor: 1 Configs: Topic: test Partition: 0 Leader: 0 Replicas: 0 Isr: 0

メッセージの受信準備

  • kafka-console-consumer --bootstrap-server localhost:9092 --topic test --from-beginning

    メッセージ待受状態になります

メッセージの送信

  • kafka-console-producer --broker-list localhost:9092 --topic test

    インタラクティブモードになった適当に文字列を入力しエンターします
    そして受信状態のコンソールにメッセージが表示されれば OK です

最後に

次回は fluentd から kafka にメッセージを送信してみたいと思います

参考サイト