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을 추가하여 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로 서로 다른 로그 형식의 데이터를 모니터링하는 설정
단일 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

