一、Logstash 核心定位1. 是什么Logstash 是 Elastic 技术栈ELK/ELFK的数据采集、过滤、转换工具用 Java 开发管道式处理日志 / 数据流。 完整栈EElasticsearch存储 检索LLogstash数据清洗转换KKibana可视化FFilebeat轻量日志采集常搭配 Logstash2. 核心作用多源数据统一接入日志、数据库、MQ、文件、API数据清洗过滤无用字段、脱敏(敏感数据)、分割、类型转换、GeoIP 解析统一格式化输出到 ES、Redis、Kafka、文件等解耦采集端与存储减轻 Elasticsearch 压力3. 三大核心组件流水线 Pipeline完整处理流程Input → Filter → Output数据源 → Input插件(接收数据) → Filter插件(清洗转换) → Output插件(输出目的地)二、三大模块详细讲解一Input 输入插件数据源入口负责读取各类原始数据支持多输入共存。1.file最常用读取本地日志文件input { file { path /var/log/nginx/*.log # 日志路径支持通配符 start_position beginning # 首次读取从头开始默认end尾部 sincedb_path /var/lib/logstash/sincedb.db # 记录文件读取偏移断点续传 discover_interval 15 # 多久扫描一次新文件 } }关键sincedb记录文件 inode 偏移重启不重复消费。2.beats接收 Filebeat/Metricbeat 数据生产标配Filebeat 轻量采集推送到 Logstash 5044 端口input { beats { port 5044 host 0.0.0.0 } }tcp /udp接收网络日志如 syslogkafka消费 kafka 消息队列高并发解耦jdbc定时读取 MySQL/PostgreSQL 数据库同步业务数据stdin标准输入测试用bin/logstash -e input{stdin{}} ...二Filter 过滤插件核心清洗层最重要中间处理环节对原始日志做解析、切割、删除、转换、匹配无过滤可省略。1.grok日志正则切割最核心原始日志是一整行字符串grok 通过匹配模板拆分字段。 内置大量预设正则模板语法%{模式:字段名}示例 Nginx 日志切割110.24.136.58 - admin [19/Jul/2026:14:30:22 0800] GET /api/user/login HTTP/1.1 200 568filter { grok { match { message %{IP:client_ip} - %{DATA:user} \[%{HTTPDATE:access_time}\] %{DATA:method} %{DATA:uri} HTTP/%{NUMBER:http_version} %{NUMBER:status:int} %{NUMBER:body_size:int} } } }%{IP:client_ip}匹配 IP存入字段 client_ip:int指定字段为数字类型默认字符串 在线调试工具grokdebug2.date时间格式化修正 timestamp日志自带时间需要覆盖系统默认的timestampdate { match [ access_time, dd/MMM/yyyy:HH:mm:ss Z ] target timestamp }3.mutate字段通用操作增删改、类型、替换、合并filter { mutate { add_field { app nginx } # 新增字段 remove_field [ user, message ] # 删除无用字段 convert { status integer } # 类型转换 gsub [ uri, \?, _ ] # 字符替换把?换成_ split { tags , } # 字符串拆数组 } }4.json解析 JSON 格式日志日志是单行 JSON 字符串时使用json { source message # 从message字段解析 target log_data # 解析结果存入log_data不写则平铺到顶层 }5. geoipIP 转地理位置国家、城市、经纬度geoip { source client_ip }6.drop丢弃不需要的日志配合条件过滤垃圾日志if [status] 404 { drop {} }7.conditional 条件判断if [app] nginx { grok { nginx规则 } } else if [app] tomcat { grok { tomcat规则 } }8.multiline多行日志合并Java 异常堆栈必备Java 报错堆栈是多行需要合并成一条事件filter { multiline { pattern ^%{DATE} # 新日志行以日期开头 negate true # 不匹配上面规则的行合并到上一条 what previous } }三Output 输出插件数据目的地清洗完成后输出到存储 / 中间件支持多输出同时写入。1.elasticsearch生产核心输出output { elasticsearch { hosts [http://127.0.0.1:9200] index nginx-access-%{YYYY.MM.dd} # 按天分索引 user elastic password xxx } } stdout { codec rubydebug } # 控制台打印调试用 }%{YYYY.MM.dd}动态日期索引便于按天管理数据。2.stdout控制台输出开发调试专用stdout { codec rubydebug # 格式化打印完整事件 }kafka输出到消息队列分层架构Logstash 输出 Kafka另一台 Logstash 消费 Kafka 做二次清洗。file输出到本地文件redis缓存队列削峰三、核心配套概念1. Codec 编解码器位置Input/Output 内部作用原始数据编解码介于数据源和管道之间。 常用 codecplain纯文本默认json单行 JSON 自动解析multiline多行合并也可在 input 配置rubydebug格式化打印日志示例 input 使用 json codecinput { tcp { port 514 codec json } }2. Pipeline 多管道Logstash 支持同时运行多个独立管道隔离不同业务日志nginx、mysql、java 配置目录config/pipelines.yml- pipeline.id: nginx-pipe path.config: /etc/logstash/conf.d/nginx.conf - pipeline.id: java-pipe path.config: /etc/logstash/conf.d/java.conf3. 工作线程与性能参数logstash.yml 核心性能配置pipeline.workers: 4 # Filter/Ouput工作线程建议等于CPU核心数 pipeline.batch.size: 125 # 批量处理条数攒够一批再过滤输出 pipeline.batch.delay: 5 # 最多等待多少ms凑批 queue.type: persisted # 持久化队列磁盘缓存防止宕机丢数据 queue.max_bytes: 4gb # 磁盘队列上限批量越大吞吐量越高内存占用越高持久队列流量突增时缓冲避免阻塞 Filebeat4. 缓存机制三层内存批队列默认进程重启丢失数据持久化磁盘队列PQ生产必开削峰容错Filebeat 本地缓存最前端兜底防止 Logstash 宕机丢日志5. 数据流转架构对比架构 1Filebeat → Logstash → ES小型业务Filebeat 采集文件Logstash 统一清洗直写 ES。架构 2Filebeat → Kafka → Logstash → ES中大型高并发Kafka 作为缓冲层削峰填谷解耦采集与清洗多台 Logstash 消费 Kafka 水平扩容。四、安装、启停、调试命令1. 启动# 加载conf.d下所有配置 bin/logstash -f config/conf.d/ # 单行简易配置测试 bin/logstash -e input{stdin{}} output{stdout{}} # 指定管道文件 bin/logstash -f test.conf2. 配置校验修改配置后先校验避免启动失败bin/logstash -f xxx.conf -t3. 日志查看运行日志/var/log/logstash/五、常见生产问题与优化1. 日志重复消费原因sincedb 丢失、Logstash 重启未消费完、ES 写入失败重试 优化固定 sincedb 路径开启持久队列ES 开启幂等写入document_id2. Logstash CPU 占用过高grok 正则过于复杂大量回溯简化正则拆分匹配pipeline.workers 设置过大超过 CPU 核心batch.size 过大单次处理数据太多3. 吞吐量不足增加 pipeline.workers中间加 Kafka 队列多实例 Logstash 消费Filebeat 开启批量推送六、完整可运行示例Nginx 日志全流程# input 接收filebeat数据 input { beats { port 5044 } } # filter清洗 filter { # 切割nginx日志 grok { match { message %{IP:client_ip} - %{DATA:user} \[%{HTTPDATE:access_time}\] %{DATA:method} %{DATA:uri} HTTP/%{NUMBER:http_version} %{NUMBER:status:int} %{NUMBER:body_size:int} } } # 转换时间字段 date { match [ access_time, dd/MMM/yyyy:HH:mm:ss Z ] } # 解析IP归属地 geoip { source client_ip } # 新增业务字段删除原始message减少存储 mutate { add_field { service nginx-access } remove_field [message, access_time] } # 过滤静态资源404日志 if [status] 404 and [uri] ~ .js|.css { drop {} } } # output输出ES控制台调试 output { elasticsearch { hosts [127.0.0.1:9200] index nginx-access-log-%{YYYY.MM.dd} } stdout { codec rubydebug } }