メインコンテンツまでスキップ

Filebeat + Logstash使用ガイド

最終更新 2026/10/03

このセクションでは、データ転送ツール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のパラメータの説明:

パラメータ名タイプ必須デフォルト値説明
urlstringtrueなしAEのデータ受信アドレス
appidstringtrueなしプロジェクトのAPPID
flush_interval_secnumberfalse2flushをトリガーする間隔。単位:秒
flush_batch_sizenumberfalse500flushをトリガーするjsonデータ数。単位:件
compressnumberfalse1データ圧縮。0は圧縮なしで、内部ネットワークで設定できます。1はgzip圧縮で、デフォルトはgzip圧縮です
uuidbooleanfalsefalseUUIDスイッチを有効にするかどうか。短い間隔でのネットワーク変動で発生しうる重複の排除に使用
is_filebeat_status_recordbooleanfalsetrueFilebeatのログ監視状態(offset、ファイル名など)の記録を有効にするかどうか

3.2.3 Logstashの実行設定​

Logstashはデフォルトでconfig/logstash.ymlを実行設定として使用します。

pipeline.workers: 1
queue.type: persisted
queue.drain: true

推奨事項

  1. user_set、つまりユーザープロパティを送信する場合は、pipeline.workersの値を1に変更してください。workersの値が1より大きいと、データの処理順序が変わってしまいます。trackイベントでは1より大きい値を設定できます。
  2. プログラムの予期しない終了でデータが失われないようにするには、queue.type: persistedを設定してください。これはLogstashが使用するバッファキューのタイプを表し、この設定によりLogstashの再起動後もバッファキュー内のデータを引き続き送信できます。

データの永続性について詳しくは公式サイトを参照してください persistent-queues

  1. 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
  1. shutdown_timeout :Filebeatがシャットダウンする前に、パブリッシャーがイベントの送信を完了するのを待つ時間です。
  2. paths:監視するログを指定します。現在はGo言語のglob関数に従って処理され、設定したディレクトリは再帰的に処理されません。
  3. hosts: 送信先として複数のLogstash hostsを指定します。loadbalanceがfalseの場合はプライマリ/バックアップのように動作し、trueの場合は負荷分散を表します
  4. 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
このページは役に立ちましたか?