惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

Engineering at Meta
Engineering at Meta
T
Threat Research - Cisco Blogs
V
Vulnerabilities – Threatpost
T
Tor Project blog
T
Troy Hunt's Blog
C
CERT Recently Published Vulnerability Notes
C
Cisco Blogs
W
WeLiveSecurity
Cloudbric
Cloudbric
cs.AI updates on arXiv.org
cs.AI updates on arXiv.org
爱范儿
爱范儿
Google Online Security Blog
Google Online Security Blog
cs.CL updates on arXiv.org
cs.CL updates on arXiv.org
Simon Willison's Weblog
Simon Willison's Weblog
Exploit-DB.com RSS Feed
Exploit-DB.com RSS Feed
Martin Fowler
Martin Fowler
Cisco Talos Blog
Cisco Talos Blog
F
Full Disclosure
MongoDB | Blog
MongoDB | Blog
CTFtime.org: upcoming CTF events
CTFtime.org: upcoming CTF events
I
Intezer
www.infosecurity-magazine.com
www.infosecurity-magazine.com
G
GRAHAM CLULEY
B
Blog RSS Feed
云风的 BLOG
云风的 BLOG
人人都是产品经理
人人都是产品经理
M
MIT News - Artificial intelligence
腾讯CDC
L
LangChain Blog
L
LINUX DO - 热门话题
H
Help Net Security
S
Schneier on Security
N
Netflix TechBlog - Medium
博客园 - Franky
酷 壳 – CoolShell
酷 壳 – CoolShell
Spread Privacy
Spread Privacy
S
Secure Thoughts
T
The Exploit Database - CXSecurity.com
P
Privacy International News Feed
P
Privacy & Cybersecurity Law Blog
Cyberwarzone
Cyberwarzone
A
About on SuperTechFans
NISL@THU
NISL@THU
钛媒体:引领未来商业与生活新知
钛媒体:引领未来商业与生活新知
D
DataBreaches.Net
The GitHub Blog
The GitHub Blog
Recorded Future
Recorded Future
雷峰网
雷峰网
AWS News Blog
AWS News Blog
V2EX - 技术
V2EX - 技术

郑文峰的博客

使用dify对接飞书多维表格 使用n8n对接飞书多维表格 服务启动时出现 OOM 一次服务升级时pg表DDL执行超时失败 Go语言高效IO缓冲技术详解 Go语言延迟初始化(Lazy Initialization)最佳实践 Go语言字符串拼接性能对比与优化指南 Go语言结构体内存对齐完全指南 Go语言空结构体:零内存消耗的高效编程 Go语言堆栈分配与逃逸分析深度解析 Go语言原子操作完全指南 Go语言内存预分配完全指南 Go语言不可变数据共享:无锁并发编程实践 Go语言零拷贝技术完全指南 Go语言遍历性能深度解析:从原理到优化实践 Go语言Interface Boxing原理与性能优化指南 Go协程池深度解析:原理、实现与最佳实践 使用etcd分布式锁导致的协程泄露与死锁问题 基于pre-commit的Python代码规范落地实践 初识 MCP Server pulsar阻塞导致logstash无法接入日志 django-prometheus使用及源码分析 kube-proxy源码分析 kubernetes service如何通过iptables转发 tcp缓存引起的日志丢失 django-apschedule定时任务异常停止 理解calico容器网络通信方案原理 理解flannel的三种容器网络方案原理 理解Linux IPIP隧道 理解VXLAN网络 理解Linux TunTap设备 快速了解iptables kafka中listener和advertised.listeners的作用 django rest_framework 分页 django后端服务、logstash和flink接入VictoriaMetrics指标监控 python中import原理 docker容器单机网络 手动实现docker容器bridge网络模型 mysql之MVCC原理 mysql之日志 使用python实现单例模式的三种方式 redis之缓存 redis之分片集群 redis之哨兵机制 redis之主从库同步 redis之持久化 redis之五种基本数据类型 go中如何处理error pod中将代码与运行环境分离 ddt源码分析 python装饰器的使用方法 读书笔记:如何阅读一本书 使用ddt实现unittest的参数化测试 使用kubeadm安装k8s 优化gin表单的错误提示信息 gin中validator模块的源码分析 go简单使用grpc python简单使用grpc k8s之PV、PVC和StorageClass k8s之StatefulSet k8s之DaemonSet k8s之Job和CronJob k8s之ConfigMap和Secret k8s之Service k8s之Pod k8s之Deployment 容器的本质 docker容器 python迭代器与生成器 python元编程 python垃圾回收机制 python上下文管理器 django rest_framework使用jwt django rest_framework异常处理 django rest_framework 自定义文档 django压缩文件下载 django rest_framework使用pytest单元测试 django restframework choice 自定义输出数据 django Filtering 使用 django viewset 和 Router 配合使用时报的错 django model的序列化 django中使用AbStractUser django.core.exceptions.ImproperlyConfigured Application labels aren't unique, duplicates users django 中 media配置 django 外键引用自身和on_delete参数 django 警告 while time zone support is active Flask使用flask_socketio实现websocket flask结合mongo tornado 文件上传 tornado 使用jwt完成用户异步认证 tornado 用户密码 bcrypt加密 tornado 结合wtforms使用表单操作 tornado finish和write区别 tornado 使用peewee-async 完成异步orm数据库操作 pyspark streaming简介 和 消费 kafka示例 使用hue创建ozzie的pyspark action workflow count的性能优化 django rest_framework Authentication django celery 结合使用 网站
使用java开发logstash的filter插件
zhengwenfeng · 2022-12-20 · via 郑文峰的博客

# 0. 前言

在工作中遇到,logstash 中的 filter 中写了大量的解析逻辑,解析性能遇到瓶颈,所以希望将该部分的逻辑转换成 java 开发的插件,以提高解析速度。

本文主要记录我开发插件的过程。

# 1. 准备开发环境

下载 logstash 源码

直接可以去 logstash github (opens new window) 中选择自己使用的版本进行下载即可。

构建 logstash

将下载的 logstash 压缩包解压出来,进入 logstash 根目录下,当前路径下有 gradlew 和 gradlew.bat 两个脚本文件,前者是在 linux 下执行的,后者是在 windows 执行的脚本。

假设当前环境是 windows,执行 gradlew.bat assemble 命令可以对当前模块进行构建。在这个过程中会去下载所有的依赖包到本地。等待构建完成,直至输出 BUILD SUCCESSFUL 代表构建成功。

gradlew.bat 脚本是对 gradle 的封装,在执行该命令时,会主动根据 gradle/wrapper/ 下的配置去下载 gradle 工具,然后再调用 gradle 进行构建模块

# 2. 编写 logstash java filter 插件

# 2.1 准备官方 demo

下载 java 插件官方模板

logstash-filter-java_filter_example (opens new window) 下载到本地使用,自定义开发的插件是基于该 example 进行修改的。

构建插件

在该项目的根目录下,创建 gradle.properties 文件,需要添加变量指定 logstash 下的 logstash-core 目录路径,使用绝对路径即可。

LOGSTASH_CORE_PATH=<target_folder>/logstash-core

1

该变量是给 build.gradle 文件中使用的。

# 2.2 开发 Filter 代码

首先来看官方提供的 demo Filter 代码,代码路径在:src\main\java\org\logstashplugins\JavaFilterExample.java,我们开发的插件基本是按照这个例子进行修改实现的。

  • 设置 pipeline 中的插件名称

首先可以看到有一个注解 @LogstashPlugin(name = "java_filter_example") name 的值是指我们在 pipeline 中填写的插件名称。

  • 在 pipeline 中传参到插件中

通过 PluginConfigSpec.stringSetting 定义变量

public static final PluginConfigSpec<String> SOURCE_CONFIG = PluginConfigSpec.stringSetting("source", "message");

1

再通过在构造方法中调用 get 方法即可获取到传入的值

this.sourceField = config.get(SOURCE_CONFIG);

1

并且需要将新增的字段添加到 configSchema 方法中并返回出去。

@Override
public Collection<PluginConfigSpec<?>> configSchema() {
	// should return a list of all configuration options for this plugin
	return Collections.singletonList(SOURCE_CONFIG);
}

1
2
3
4
5

  • filter 主体编码

该插件的主体是 filter 方法,也就是数据的过滤走的 filter 方法,我们将想要做的解析规则实现在该方法中即可。

可以看到该方法中有一个对 events 遍历的处理,每一个 Event 都是进来的每一条数据,然后对该条数据进行处理转换,最后再将转换好的 events 传出去。

可以看到官方的案例是将传入的 message 字符串翻转。

@Override
public Collection<Event> filter(Collection<Event> events, FilterMatchListener matchListener) {
	for (Event e : events) {
		Object f = e.getField(sourceField);
		if (f instanceof String) {
			e.setField(sourceField, StringUtils.reverse((String)f));
			matchListener.filterMatched(e);
		}
	}

	return events;
}

1
2
3
4
5
6
7
8
9
10
11
12

# 3. 单元测试

单测对插件来说至关重要,插件的规则转换流程、判断逻辑都非常多,各种类型的数据都可能导致插件出错,而插件验证需要编译、打包、安装再测试,流程较长,所以我们可以通过单测来减少以上流程的进行,在单测中就把所有的可能性都验证到,节省大量的时间。并且在后续迭代修改中,可以减少改动引发。

建议可以使用 junit 的参数化单测方式,可以提高单测的效率和数量。这个需要在 build.gradle 文件中的 dependencies 添加支持参数化的库来支持。

# 4. 打包部署 Filter 插件

# 4.1 元数据信息

我们需要在 build.gradle 文件中修改部分的插件元数据信息,像 description、authors 和 email 等字段都可以随意填写,以下字段需要注意:

  • group,需要和包名相同
  • pluginClass,需要和插件 Filter 的类名相同
  • pluginName,需要和 @LogstashPlugin 中的 name 相同

# 4.2 打包任务

通过执行 gradlew.bat gem 进行插件打包任务,最后会在插件根目下生成 .gem 的插件安装包文件。

# 4.3 安装

安装有在线安装和离线安装两种方式。

注意:我们需要去官网下载可以直接使用的 logstash,而不能使用上面自己下载的 logstash 源码。

在线安装

在线安装会去访问 Elastic 的官网,所以需要是在线的环境。

通过执行 logstash/bin 路径下的 logstash-plugin 命令进行安装,等待片刻即可安装成功。

logstash-plugin install /path/javaPlugin.gem

1

离线安装

在某些场景下,环境是不能连接外网的,所以需要使用离线安装的方式。

将生成的 gem 插件压缩到 zip 包中,然后再使用 logstash-plugin 命令进行安装。

logstash-plugin install file:///tmp/plugin.zip

1

# 5. 验证

官方的插件 example 的功能是翻转字符串的功能,所以我们只需要验证该功能即可。

  1. 创建一个 pipeline.conf
input {
    # 输入一个字符串
    generator { message => "Hello world!" count => 1 }
}

filter {
	# 在插件中@LogstashPlugin配置的插件名称
    java_filter_example {}
}

output {
    # 直接打印到控制台
    stdout { }
}

1
2
3
4
5
6
7
8
9
10
11
12
13
14

  1. 启动 logstash 加载上面的 pipeline.conf
logstash -f pipeline.conf

1

输出如下,可以看到 message 字段中的 Hello world!被翻转了。

{
	"host" => {
		"name" => "4-sip0060"
	},
	"event" => {
		"original" => "Hello world!",
		"sequence" => 0
	},
	"@timestamp" => 2022-12-20T07:27:46.634166300Z,
	"@version" => "1",
	"message" => "!dlrow olleH"
}

1
2
3
4
5
6
7
8
9
10
11
12

# 6. 相关链接