Filebeat + Logstash user guide
This section describes how to use the data transfer tool Filebeat + Logstash:
Before you start the integration, read Data rules. After you're familiar with the data format and data rules of AE, read this guide to complete the integration.
Data uploaded through Filebeat + Logstash must follow the AE data format
Note: Logstash has low throughput. To import a large amount of historical data, we recommend that you use the DataX engine or the Logbus tool
1. Filebeat + Logstash overview
The Filebeat + Logstash tool is mainly used to import log data into the AE backend in real time. It monitors the file streams in the log directory of the server, and when new data is written to any log file in the directory, it sends the data to the AE backend in real time.
Logstash is an open-source, server-side data processing pipeline that ingests data from multiple sources simultaneously, transforms it, and then sends it to your favorite "stash". Logstash official introduction
Filebeat is a log data shipper for local files. It monitors log directories or specific log files (tail file) and gives you a lightweight way to forward and centralize logs and files. Filebeat official introduction
The following figure shows the data collection flow based on Filebeat + Logstash:
2. Download and install Filebeat + Logstash
Note: Logstash-6.x or later is required, and the server must have a JDK environment
2.1 Download and install Logstash
See the Logstash official installation documentation and choose a download method
2.2 logstash-output-thinkingdata plugin
Latest version: 1.2.1
Update time: 2023-07-26
2.2.1 Install and uninstall the logstash-output-thinkingdata plugin
The plugin checks whether the data is JSON data, and then packages the data and sends it to AE
Run the following command in the logstash directory:
bin/logstash-plugin install logstash-output-thinkingdata
The installation takes a while. After the installation succeeds, run:
bin/logstash-plugin list
If logstash-output-thinkingdata appears in the list, the installation is successful.
Other commands are as follows:
To upgrade the plugin, run:
bin/logstash-plugin update logstash-output-thinkingdata
To uninstall the plugin, run:
bin/logstash-plugin uninstall logstash-output-thinkingdata
2.2.2 Change Log
v1.2.1 2023/07/26
- Added format validation for the data in message
v1.2.0 2023/04/25
- Added support for passing multiple data records in message
v1.1.0 2021/01/27
- Added support for the #app_id format in data
v1.0.0 2020/06/09
- Receives the message passed by an event in Logstash and sends it to AE
2.3 Download and install Filebeat
See the Filebeat official installation documentation and choose a download method
3. Filebeat + Logstash usage
3.1 Data preparation
1. First, use ETL to convert the data to be transferred into the AE data format, and write it to local files or send it to a Kafka cluster. If you use a consumer of a server-side SDK (such as Java) that writes to Kafka or local files, the data is already in the correct format and doesn't need to be converted.
2. Determine the directory where the files of the data to be uploaded are stored, or the Kafka address and topic, and configure Filebeat + Logstash accordingly. Filebeat + Logstash monitors file changes in the file directory (new files are monitored, and existing files are tailed), or subscribes to the data in Kafka.
3. Don't directly rename data logs that are stored in the monitored directory and have already been uploaded. Renaming a log is equivalent to creating a new file, and Filebeat may upload these files again, which causes duplicate data.
3.2 Logstash configuration
3.2.1 Logstash pipeline configuration
Logstash can run multiple pipelines at the same time. Pipelines don't affect each other, and each has its own input and output configuration. The pipeline configuration file is located at config/pipelines.yml. If you're already using Logstash for other log collection work, you can add a new pipeline to your existing Logstash that's dedicated to collecting AE log data and sending it to AE
The configuration of the new pipeline is as follows:
# Pipeline configuration for thinkingdata
- pipeline.id: thinkingdata-output
# If you upload user properties, set the number of cores to 1. If you only upload event properties, you can set it to a value less than or equal to the number of CPUs on this machine
pipeline.workers: 1
# Type of buffer queue to use
queue.type: persisted
# Use a separate input and output configuration
path.config: "/home/elk/conf/ta_output.conf"
For more pipeline configuration options, see the official website: Pipeline
3.2.2 Logstash input and output configuration
ta_output.conf reference example:
This scenario shows the Logstash input and output configuration for data files generated by the SDK
# Use beats as the input
input {
beats {
port => "5044"
}
}
# If the data isn't generated by a server-side SDK or doesn't follow the AE format, use a filter to process the raw data. We recommend a ruby plugin as the filter; see section 4 for examples
#filter {
# if "log4j" in [tags] {
# ruby {
# path => "/home/elk/conf/rb/log4j_filter.rb"
# }
# }
#}
# Use thinkingdata as the output
output{
thinkingdata {
url => "http://YOUR_RECEIVER_URL/logbus"
appid =>"YOUR_APPID"
}
}
thinkingdata parameters:
| Parameter name | Type | Required | Default value | Description |
|---|---|---|---|---|
| url | string | true | None | AE data receiving URL |
| appid | string | true | None | APPID of the project |
| flush_interval_sec | number | false | 2 | Interval for triggering a flush, in seconds |
| flush_batch_size | number | false | 500 | Number of JSON records that triggers a flush |
| compress | number | false | 1 | Data compression. 0: no compression, which you can use on an intranet; 1: gzip compression. Defaults to gzip compression |
| uuid | boolean | false | false | Whether to enable UUID, used for deduplication when network fluctuations occur within short intervals |
| is_filebeat_status_record | boolean | false | true | Whether to enable Filebeat log status monitoring, such as offset and file name |
3.2.3 Logstash runtime configuration
By default, Logstash uses config/logstash.yml as the runtime configuration.
pipeline.workers: 1
queue.type: persisted
queue.drain: true
Recommendations
- When you report user properties with user_set, change pipeline.workers to 1. If workers is greater than 1, the order in which data is processed changes. For track events, you can set it to a value greater than 1.
- To make sure data isn't lost when the program terminates unexpectedly, set queue.type: persisted, which specifies the type of buffer queue that Logstash uses. With this setting, Logstash continues sending the data in the buffer queue after it restarts.
For more about data persistence, see the official website persistent-queues
- Set queue.drain to true. This setting makes Logstash send all the data in the buffer queue before it exits normally.
For more details, see the official website logstash.yml
3.2.4 Start Logstash
In the Logstash installation directory: 1. Start Logstash directly, which uses config/pipelines.yml as the pipeline configuration and runtime configuration
bin/logstash
2. Start Logstash with ta_output.conf specified as the input and output configuration file, which uses config/logstash.yml as the runtime configuration
bin/logstash -f /youpath/ta_output.conf
3. Start Logstash in the background
nohup bin/logstash -f /youpath/ta_output.conf > /dev/null 2>&1 &
For more startup options, see the Logstash official startup documentation
3.3 Filebeat configuration
3.3.1 Filebeat runtime configuration
Filebeat reads the log files of the server-side SDK. The default Filebeat configuration file is filebeat.yml. A reference configuration for config/filebeat.yml is as follows:
#======================= Filebeat inputs =========================
filebeat.shutdown_timeout: 5s
filebeat.inputs:
- type: log
enabled: true
paths:
- /var/log/log.*
#- c:\programdata\elasticsearch\logs\*
#------------------------- Logstash output ----------------------------
output.logstash:
# You can enter the Logstash process on one server, or Logstash processes on multiple servers
hosts: ["ip1:5044","ip2:5044"]
loadbalance: false
- shutdown_timeout: how long Filebeat waits on shutdown for the publisher to finish sending events before Filebeat shuts down.
- paths: specifies the logs to monitor. Paths are currently processed with the Go glob function, and configured directories aren't processed recursively.
- hosts: the addresses of multiple Logstash hosts to send data to. When loadbalance is false, it works like an active/standby setup; true enables load balancing
- loadbalance: if you report user properties with user_set, don't set loadbalance : true. With this setting, data is sent to all Logstash instances in round-robin fashion, which is likely to disrupt the order of the data.
If you only import track events, you can set up multiple Logstash processes and set loadbalance : true. The default Filebeat configuration is loadbalance : false
For more about Filebeat, see the Filebeat official documentation
3.3.2 Start Filebeat
After you start filebeat, view the related output:
./filebeat -e -c filebeat.yml
Start Filebeat in the background
nohup ./filebeat -c filebeat.yml > /dev/null 2>&1 &
4. Configuration examples
4.1 Configure filebeat to monitor data in different log formats
For a single Filebeat, you can configure filebeat.yml as follows. It uses about 10 MB of memory at runtime. You can also start multiple filebeat processes to monitor different log formats. For details, see the official website
#=========================== Filebeat inputs =============================
filebeat.inputs:
- type: log
enabled: true
#Directory to monitor
paths:
- /home/elk/sdk/*/log.*
#Tag the data so that Logstash can match it; no filter processing is needed
tags: ["sdklog"]
#Data split by special characters
- type: log
enabled: true
paths: /home/txt/split.*
#Tag the data so that Logstash can match it; filter processing is needed
tags: ["split"]
# Data received from log4j
- type: log
enabled: true
paths: /home/web/logs/*.log
# Filter processing is needed
tags: ["log4j"]
# nginx log data
- type: log
enabled: true
paths: /home/web/logs/*.log
tags: ["nginx"]
4.2 Logstash configuration
Check whether the Logstash configuration file contains errors
bin/logstash -f /home/elk/ta_output.conf --config.test_and_exit
Note: All ruby plugins in the following scripts are written in ruby syntax. If you're familiar with java, you can use the custom parser in the logbus directory
4.2.1 Server-side SDK logs
You can configure ta_output.conf as follows
input {
beats {
port => "5044"
}
}
# Use thinkingdata as the output
output{
thinkingdata {
url => "url"
appid =>"appid"
# compress => 0
# uuid => true
}
}
4.2.2 log4j logs
You can set the log4j format as follows, based on your business logs:
//Log format
//[%d{yyyy-MM-dd HH:mm:ss.SSS}] Log time that matches the AE input time format; it can also be yyyy-MM-dd HH:mm:ss
//[%level{length=5}] Log level: debug, info, warn, error
//[%thread-%tid] Current thread information
//[%logger] Fully qualified name of the class that the current log message belongs to
//[%X{hostName}] Host name of the current node. Must be customized through MDC.
//[%X{ip}] IP of the current node. Must be customized through MDC.
//[%X{userId}] Unique ID of the user login. You can set it to account_id or another value. AE requires that account_id and distinct_id can't both be empty. You can set other properties based on your business. Must be customized through MDC.
//[%X{applicationName}] Name of the current application. Must be customized through MDC.
//[%F,%L,%C,%M] %F: name of the file (class) that the current log message belongs to; %L: line number of the log message in its file; %C: fully qualified class name of the file that the current log belongs to; %M: name of the method that the current log belongs to
//[%m] Log details
//%ex Exception information
//%n Line break
<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 configuration
input {
beats {
port => "5044"
}
}
filter {
if "log4j" in [tags]{
#You can also perform other filter data processing
ruby {
path => "/home/conf/log4j.rb"
}
}
}
# Use thinkingdata as the output
output{
thinkingdata {
url => "url"
appid =>"appid"
}
}
The /home/conf/log4j.rb script is as follows
# The event parameter gives you access to all the properties in input
def filter(event)
_message = event.get('message') #message is each log line that you upload
begin
#A regex extracts the data in the correct format here. We recommend putting error logs in a separate file: error logs are too long, have no analysis scenario, and span multiple lines
#The data here is in a format like this: _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, #Can be used when type is track; if there's no #event_name, don't upload it
'#type' =>'track', #Can be obtained from the file if the file has it; determine whether the reported data is user properties or event properties
'properties' => properties
}
event.set('message',data.to_json)
return [event]
rescue
# puts _message
puts "Data does not match the regex format"
return [] #Not reported
end
end
4.2.3 Nginx logs
First, define the Nginx log format. If you set it to JSON format:
input {
beats {
port => "5044"
}
}
filter {
#If the data is all in the same format, you don't need to check tags
if "nginx" in [tags]{
ruby {
path => "/home/conf/nginx.rb"
}
}
}
# Use thinkingdata as the output
output{
thinkingdata {
url => "url"
appid =>"appid"
}
}
The /home/conf/nginx.rb script is as follows:
require 'date'
def filter(event)
#Extract log information like the following
# {"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')
#Parse the log level and request time, and save them to the event object
#Parse JSON-format logs
#Get the log content
#Convert it to a JSON object
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
#Skip reporting when #account_id and #distinct_id are both 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'],#Can be null
'#time' => time, #Can't be null
'#account_id'=>account_id, # account_id and distinct_id can't both be null
'#event_name'=>url, #Can be used when type is track; if there's no #event_name, don't upload it
'#type' =>'track', #Can be obtained from the file if the file has it; determine whether the reported data is user properties or event properties
'properties' => properties
}
event.set('message',data.to_json)
return [event]
end
4.2.4 Other logs
Configure ta_output.conf as follows
input {
beats {
port => "5044"
}
}
filter {
if "other" in [tags]{
#You can also perform other filter data processing
ruby {
path => "/home/conf/other.rb"
}
}
}
# Use thinkingdata as the output
output{
thinkingdata {
url => "url"
appid =>"appid"
}
}
The /home/conf/other.rb script is as follows:
# def register(params)
# # The parameters obtained through params here are passed in through script_params in the logstash file
# #@message = params["xx"]
# end
def filter(event)
#Process business data here. If no grok or similar processing has been done, get the raw data directly from message and process it
begin
_message = event.get('message') #message is each log line that you upload
#The log processed here is data like the following
#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(/,/)
#Or match the data with a regex; for JSON data, parse it with ruby syntax. You can process the data here and then report it
#mess = /regex/.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 is the int type, to_f is the float type, to_s is the string type (default)
}
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
#If you don't want a record to go through, or an error occurs, return empty for that record. For example, if account_id and distinct_id are both null, return [] directly
return []
end
end

