Filebeat + Logstash使用ガイド
このセクションでは、データ転送ツールFilebeat + Logstashの使用方法を説明します:
接続を始める前に、まずデータルールを読んでください。AEのデータ形式とデータルールを理解したうえで、このガイドを読んで接続を行ってください。
Filebeat + LogstashでアップロードするデータはAEのデータ形式に従う必要があります
注意:Logstashはスループットが低いため、大量の履歴データをインポートする場合は、次の使用をお勧めします: DataXエンジン またはLogbusツール
1. Filebeat + Logstashの概要
Filebeat + Logstashツールは主に、ログデータをリアルタイムでAE管理画面にインポートするために使用します。サーバーのログディレクトリ内のファイルストリームを監視し、ディレクトリ内のいずれかのログファイルに新しいデータが生成されると、リアルタイムでAE管理画面に送信します。
Logstashはオープンソースのサーバーサイドデータ処理パイプラインで、複数のソースから同時にデータを収集・変換し、お好みの「ストレージ」に送信できます。Logstash公式サイトの紹介
Filebeatはローカルファイルのログデータ収集ツールで、ログディレクトリや特定のログファイル(tail file)を監視できます。Filebeatは、ログやファイルを転送・集約するための軽量な方法を提供します。Filebeat公式サイトの紹介
Filebeat + Logstashによる収集の流れは次の図のとおりです:
2. Filebeat + Logstashのダウンロードとインストール
注意:Logstash-6.x以上のバージョンで、サーバーにJDK環境が必要です
2.1 Logstashのダウンロードとインストール
Logstash公式インストールドキュメントを参照し、方法を選択してダウンロードしてください
2.2 logstash-output-thinkingdataプラグイン
最新バージョン:1.2.1
更新日:2023-07-26
2.2.1 logstash-output-thinkingdataプラグインのインストールとアンインストール
このプラグインは、データがjsonデータかどうかをチェックし、データをまとめてAEに送信します
logstashディレクトリで次を実行します:
bin/logstash-plugin install logstash-output-thinkingdata
インストールには少し時間がかかります。インストールが成功するまでお待ちください。その後、次を実行します:
bin/logstash-plugin list
一覧にlogstash-output-thinkingdataがあれば、インストールは成功しています。
その他のコマンドは次のとおりです:
プラグインをアップグレードするには、次を実行します:
bin/logstash-plugin update logstash-output-thinkingdata
プラグインをアンインストールするには、次を実行します:
bin/logstash-plugin uninstall logstash-output-thinkingdata
2.2.2 Change Log
v1.2.1 2023/07/26
- messageのデータに対する形式チェックを追加
v1.2.0 2023/04/25
- messageで複数のデータの受け渡しに対応
v1.1.0 2021/01/27
- データ内の#app_idの形式に対応
v1.0.0 2020/06/09
- Logstashのeventで渡されたmessageを受け取り、AEに送信
2.3 Filebeatのダウンロードとインストール
Filebeat公式インストールドキュメントを参照し、方法を選択してダウンロードしてください
3. Filebeat + Logstashの使用方法
3.1 データの準備
1.まず、転送するデータをETLでAEのデータ形式に変換し、ローカルに書き出すかKafkaクラスターに転送します。JavaなどのサーバーSDKの、Kafkaまたはローカルファイルに書き込むconsumerを使用している場合、データはすでに正しい形式になっているため、変換は不要です。
2.アップロードするデータのファイルを格納するディレクトリ、またはKafkaのアドレスとtopicを決め、Filebeat + Logstashの関連設定を行います。Filebeat + Logstashは、ファイルディレクトリ内のファイルの変更(新規ファイルの作成や既存ファイルのtail)を監視するか、Kafkaのデータをサブスクライブします。
3.監視ディレクトリに格納され、すでにアップロード済みのデータログの名前を直接変更しないでください。ログの名前を変更すると新しいファイルを作成したのと同じ扱いになり、Filebeatがこれらのファイルを再アップロードしてデータが重複する可能性があります。
3.2 Logstashの設定
3.2.1 Logstash Pipelineの設定
Logstashは複数のPipelineの同時実行に対応しています。各Pipelineは互いに影響せず、それぞれ独立した入出力設定を持ちます。Pipelineの設定ファイルはconfig/pipelines.ymlにあります。現在Logstashで他のログ収集作業を行っている場合は、既存のLogstashにAEのログデータ収集専用のPipelineを1つ追加し、AEに送信できます
追加するPipelineの設定は次のとおりです:
# thinkingdataのパイプライン設定
- pipeline.id: thinkingdata-output
# ユーザープロパティをアップロードする場合はコア数を1に設定し、イベントプロパティのみをアップロードする場合はマシンのCPU数以下に設定できます
pipeline.workers: 1
# 使用するバッファキューのタイプ
queue.type: persisted
# 異なる入出力設定を使用
path.config: "/home/elk/conf/ta_output.conf"
Pipelineの詳しい設定は公式サイトを参照してください: Pipeline
3.2.2 Logstashの入出力設定
ta_output.confの設定例:
ここでは、SDKで生成されたデータファイルを対象としたLogstashの入出力設定を示します
# 入力としてbeatsを使用
input {
beats {
port => "5044"
}
}
# データがサーバーSDKで生成されたものでない場合、またはAEの形式に準拠していない場合は、filterで元データをフィルタリングする必要があります。ここではrubyプラグインでのフィルタリングをお勧めします。第4章に例があります
#filter {
# if "log4j" in [tags] {
# ruby {
# path => "/home/elk/conf/rb/log4j_filter.rb"
# }
# }
#}
# 出力としてthinkingdataを使用
output{
thinkingdata {
url => "http://データ収集アドレス/logbus"
appid =>"お客様のAppID"
}
}
thinkingdataのパラメータの説明:
| パラメータ名 | タイプ | 必須 | デフォルト値 | 説明 |
|---|---|---|---|---|
| url | string | true | なし | AEのデータ受信アドレス |
| appid | string | true | なし | プロジェクトのAPPID |
| flush_interval_sec | number | false | 2 | flushをトリガーする間隔。単位:秒 |
| flush_batch_size | number | false | 500 | flushをトリガーするjsonデータ数。単位:件 |
| compress | number | false | 1 | データ圧縮。0は圧縮なしで、内部ネットワークで設定できます。1はgzip圧縮で、デフォルトはgzip圧縮です |
| uuid | boolean | false | false | UUIDスイッチを有効にするかどうか。短い間隔でのネットワーク変動で発生しうる重複の排除に使用 |
| is_filebeat_status_record | boolean | false | true | Filebeatのログ監視状態(offset、ファイル名など)の記録を有効にするかどうか |
3.2.3 Logstashの実行設定
Logstashはデフォルトでconfig/logstash.ymlを実行設定として使用します。
pipeline.workers: 1
queue.type: persisted
queue.drain: true
推奨事項
- user_set、つまりユーザープロパティを送信する場合は、pipeline.workersの値を1に変更してください。workersの値が1より大きいと、データの処理順序が変わってしまいます。trackイベントでは1より大きい値を設定できます。
- プログラムの予期しない終了でデータが失われないようにするには、queue.type: persistedを設定してください。これはLogstashが使用するバッファキューのタイプを表し、この設定によりLogstashの再起動後もバッファキュー内のデータを引き続き送信できます。
データの永続性について詳しくは公式サイトを参照してください persistent-queues
- queue.drainの値をtrueに設定します。この設定により、Logstashは正常終了する前にバッファキュー内のすべてのデータを送信し終えます。
詳しくは公式サイトを参照してください logstash.yml
3.2.4 Logstashの起動
Logstashのインストールディレクトリで、1.直接起動する場合は、config/pipelines.ymlをPipeline設定と実行設定として使用します
bin/logstash
2.ta_output.confを入出力設定ファイルとして指定して起動する場合は、config/logstash.ymlを実行設定として使用します
bin/logstash -f /youpath/ta_output.conf
3.バックグラウンドで起動
nohup bin/logstash -f /youpath/ta_output.conf > /dev/null 2>&1 &
その他の起動方法については、Logstash公式起動ドキュメントを参照してください
3.3 Filebeatの設定
3.3.1 Filebeatの実行設定
FilebeatはバックエンドSDKのログファイルを読み取ります。Filebeatのデフォルトの設定ファイルはfilebeat.ymlです。config/filebeat.ymlの設定例は次のとおりです:
#======================= Filebeat inputs =========================
filebeat.shutdown_timeout: 5s
filebeat.inputs:
- type: log
enabled: true
paths:
- /var/log/log.*
#- c:\programdata\elasticsearch\logs\*
#------------------------- Logstash output ----------------------------
output.logstash:
# 1台のサーバーのプロセスを指定することも、複数のサーバー上のLogstashプロセスを指定することもできます
hosts: ["ip1:5044","ip2:5044"]
loadbalance: false
- shutdown_timeout :Filebeatがシャットダウンする前に、パブリッシャーがイベントの送信を完了するのを待つ時間です。
- paths:監視するログを指定します。現在はGo言語のglob関数に従って処理され、設定したディレクトリは再帰的に処理されません。
- hosts: 送信先として複数のLogstash hostsを指定します。loadbalanceがfalseの場合はプライマリ/バックアップのように動作し、trueの場合は負荷分散を表します
- loadbalance: user_setでユーザープロパティを送信する場合は、loadbalance : trueに設定しないでください。設定するとラウンドロビン方式ですべてのLogstashにデータが送信されるため、データの順序が乱れる可能性が高くなります。
trackイベントのみをインポートする場合は、複数のLogstashプロセスを設定でき、loadbalance : trueに設定することもできます。Filebeatのデフォルト設定はloadbalance : falseです
Filebeatの資料は次を参照してください:Filebeat公式ドキュメント
3.3.2 Filebeatの起動
filebeatの起動後、関連する出力情報を確認します:
./filebeat -e -c filebeat.yml
バックグラウンドで起動
nohup ./filebeat -c filebeat.yml > /dev/null 2>&1 &
4. 設定例
4.1 filebeatで異なるログ形式のデータを監視する設定
1つのFilebeatのfilebeat.ymlは次のように設定できます。実行メモリは約10 MBです。複数のfilebeatプロセスを起動して異なるログ形式を監視することもできます。詳しくは公式サイトを参照してください
#=========================== Filebeat inputs =============================
filebeat.inputs:
- type: log
enabled: true
#監視ディレクトリ
paths:
- /home/elk/sdk/*/log.*
#データにtagsを付けてLogstashでマッチングします。filter処理は不要です
tags: ["sdklog"]
#特殊文字で区切られたデータ
- type: log
enabled: true
paths: /home/txt/split.*
#データにtagsを付けてLogstashでマッチングします。filter処理が必要です
tags: ["split"]
# log4jで受信したデータ
- type: log
enabled: true
paths: /home/web/logs/*.log
# filter処理が必要です
tags: ["log4j"]
# nginxのログデータ
- type: log
enabled: true
paths: /home/web/logs/*.log
tags: ["nginx"]
4.2 Logstashの設定
Logstashの設定ファイルにエラーがないかを検査します
bin/logstash -f /home/elk/ta_output.conf --config.test_and_exit
注意:以下のスクリプトのrubyプラグインはすべてruby構文で記述されています。javaに慣れている場合は、logbusディレクトリにあるカスタムパーサーを使用できます
4.2.1 サーバーSDKのログ
ta_output.confは次のように設定できます
input {
beats {
port => "5044"
}
}
# 出力としてthinkingdataを使用
output{
thinkingdata {
url => "url"
appid =>"appid"
# compress => 0
# uuid => true
}
}
4.2.2 log4jのログ
log4jの形式は次のように設定できます。業務ログの状況に応じて設定してください:
//ログ形式
//[%d{yyyy-MM-dd HH:mm:ss.SSS}] AEの入力時間形式に準拠したログ。yyyy-MM-dd HH:mm:ssも可
//[%level{length=5}] ログレベル。debug、info、warn、error
//[%thread-%tid] 現在のスレッド情報
//[%logger] 現在のログ情報が属するクラスのフルパス
//[%X{hostName}] 現在のノードのホスト名。MDCでカスタマイズする必要があります。
//[%X{ip}] 現在のノードのip。MDCでカスタマイズする必要があります。
//[%X{userId}] ユーザーログインの一意のID。account_idを設定することも、他の値を設定することもできます。AEではaccount_idとdistinct_idを同時に空にすることはできません。業務に応じて他のプロパティを設定することもできます。MDCでカスタマイズする必要があります。
//[%X{applicationName}] 現在のアプリケーション名。MDCでカスタマイズする必要があります。
//[%F,%L,%C,%M] %F:現在のログ情報が属するファイル(クラス)名、%L:ログ情報の所属ファイル内での行番号、%C:現在のログが属するファイルの完全なクラス名、%M:現在のログが属するメソッド名
//[%m] ログの詳細
//%ex 例外情報
//%n 改行
<property name="patternLayout">[%d{yyyy-MM-dd HH:mm:ss.SSS}] [%level{length=5}] [%thread-%tid] [%logger] [%X{hostName}] [%X{ip}] [%X{userId}] [%X{applicationName}] [%F,%L,%C,%M] [%m] ## '%ex'%n
</property>
ta_output.confの設定
input {
beats {
port => "5044"
}
}
filter {
if "log4j" in [tags]{
#他のfilterでデータ処理を行うこともできます
ruby {
path => "/home/conf/log4j.rb"
}
}
}
# 出力としてthinkingdataを使用
output{
thinkingdata {
url => "url"
appid =>"appid"
}
}
/home/conf/log4j.rbのスクリプトは次のとおりです
# ここでは引数eventで、input内のすべてのプロパティを取得できます
def filter(event)
_message = event.get('message') #messageはアップロードする各ログです
begin
#ここでは正規表現で正しい形式のデータを抽出します。errorログは別ファイルにまとめることをお勧めします。エラーログは長すぎて分析シーンがなく、複数行にまたがるためです
#ここでのデータは次のような形式です _message ="[2020-06-08 23:19:56.003] [INFO] [main-1] [cn.thinkingdata] [x] [123.123.123.123] [x] [x] [StartupInfoLogger.java,50,o)] ## ''"
mess = /\[(.*?)\] \[(.*?)\] \[(.*?)\] \[(.*?)\] \[(.*?)\] \[(.*?)\] \[(.*?)\] \[(.*?)\] \[(.*?)\] ## '(.*?)'/.match(_message)
time = mess[1]
level = mess[2]
thread = mess[3]
class_name_all = mess[4]
event_name = mess[5]
ip = mess[6]
account_id = mess[7]
application_name = mess[8]
other_mess = mess[9]
exp = mess[10]
if event_name.empty? || account_id.empty?
return []
end
properties = {
'level' => level,
'thread' => thread,
'application_name' => application_name,
'class_name_all' => class_name_all,
'other_mess' => other_mess,
'exp' => exp
}
data = {
'#ip' => ip,
'#time' => time,
'#account_id'=>account_id,
'#event_name'=>event_name, #typeがtrackの場合に使用できます。#event_nameがない場合はアップロード不要です
'#type' =>'track', #ファイルから取得できます。ファイルにある場合は、送信データがユーザープロパティかイベントプロパティかを確認してください
'properties' => properties
}
event.set('message',data.to_json)
return [event]
rescue
# puts _message
puts "データが正規表現の形式に合致しません"
return [] #送信しません
end
end
4.2.3 Nginxのログ
まずNginxのログ形式を定義します。json形式に設定した場合は次のとおりです
input {
beats {
port => "5044"
}
}
filter {
#同じ形式のデータであればtagsを判定する必要はありません
if "nginx" in [tags]{
ruby {
path => "/home/conf/nginx.rb"
}
}
}
# 出力としてthinkingdataを使用
output{
thinkingdata {
url => "url"
appid =>"appid"
}
}
/home/conf/nginx.rbのスクリプトは次のとおりです:
require 'date'
def filter(event)
#次のようなログ情報を取り出します
# {"accessip_list":"124.207.82.22","client_ip":"123.207.82.22","http_host":"203.195.163.239","@timestamp":"2020-06-03T19:47:42+08:00","method":"GET","url":"/urlpath","status":"304","http_referer":"-","body_bytes_sent":"0","request_time":"0.000","user_agent":"s","total_bytes_sent":"180","server_ip":"10.104.137.230"}
logdata= event.get('message')
#ログレベルとリクエスト時間を解析し、eventオブジェクトに保存
#json形式のログを解析
#ログの内容を取得
#jsonオブジェクトに変換
logInfoJson=JSON.parse logdata
time = DateTime.parse(logInfoJson['@timestamp']).to_time.localtime.strftime('%Y-%m-%d %H:%M:%S')
url = logInfoJson['url']
account_id = logInfoJson['user_agent']
#event_name
#account_idと#distinct_idが両方ともnullの場合は送信をスキップ
if url.empty? || url == "/" || account_id.empty?
return []
end
properties = {
"accessip_list" => logInfoJson['accessip_list'],
"http_host"=>logInfoJson['http_host'],
"method" => logInfoJson['method'],
"url"=>logInfoJson['url'],
"status" => logInfoJson['status'].to_i,
"http_referer" => logInfoJson['http_referer'],
"body_bytes_sent" =>logInfoJson['body_bytes_sent'],
"request_time" => logInfoJson['request_time'],
"total_bytes_sent" => logInfoJson['total_bytes_sent'],
"server_ip" => logInfoJson['server_ip'],
}
data = {
'#ip' => logInfoJson['client_ip'],#nullでも可
'#time' => time, #nullは不可
'#account_id'=>account_id, # account_idとdistinct_idを同時にnullにすることはできません
'#event_name'=>url, #typeがtrackの場合に使用できます。#event_nameがない場合はアップロード不要です
'#type' =>'track', #ファイルから取得できます。ファイルにある場合は、送信データがユーザープロパティかイベントプロパティかを確認してください
'properties' => properties
}
event.set('message',data.to_json)
return [event]
end
4.2.4 その他のログ
ta_output.confの設定は次のとおりです
input {
beats {
port => "5044"
}
}
filter {
if "other" in [tags]{
#他のfilterでデータ処理を行うこともできます
ruby {
path => "/home/conf/other.rb"
}
}
}
# 出力としてthinkingdataを使用
output{
thinkingdata {
url => "url"
appid =>"appid"
}
}
/home/conf/other.rbのスクリプトは次のとおりです:
# def register(params)
# # ここでparamsから取得するパラメータは、logstashファイルでscript_paramsを通じて渡されたものです
# #@message = params["xx"]
# end
def filter(event)
#ここで業務データを処理します。grokなどの一連の処理を行っていない場合は、messageから直接元データを取得して処理します
begin
_message = event.get('message') #messageはアップロードする各ログです
#ここで処理するlogは次のようなデータです
#2020-06-01 22:20:11.222,123455,123.123.123.123,login,nihaoaaaaa,400
time,account_id,ip,event_name,msg,intdata=_message.split(/,/)
#または正規表現でデータをマッチングします。jsonデータの場合はruby構文でjsonを解析するだけです。ここでデータを処理して送信できます
#mess = /正規表現/.match(_message)
#time = mess[1]
#level = mess[2]
#thread = mess[3]
#class_name_all = mess[4]
#event_name = mess[5]
#account_id = mess[6]
#api = mess[7]
properties = {
'msg' => msg,
'int_data' => intdata.to_i #to_iはint型、to_fはfloat型、to_sはstring型(デフォルト)
}
data = {
'#time' => time,
'#account_id'=>account_id,
'#event_name'=>event_name,
'#ip' => ip,
'#type' =>'track',
'properties' => properties
}
event.set('message',data.to_json)
return [event]
rescue
#あるデータを送信したくない場合やエラーが発生した場合は、そのデータに空を返すことができます。たとえばaccount_idとdistinct_idがどちらもnullの場合は、直接[]を返します
return []
end
end

