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

推荐订阅源

F
Fortinet All Blogs
C
Check Point Blog
GbyAI
GbyAI
博客园 - 司徒正美
爱范儿
爱范儿
N
Netflix TechBlog - Medium
H
Hacker News: Front Page
Threat Intelligence Blog | Flashpoint
Threat Intelligence Blog | Flashpoint
Security Latest
Security Latest
C
Cyber Attacks, Cyber Crime and Cyber Security
博客园 - Franky
Recent Announcements
Recent Announcements
P
Privacy International News Feed
T
Tor Project blog
Y
Y Combinator Blog
有赞技术团队
有赞技术团队
cs.CL updates on arXiv.org
cs.CL updates on arXiv.org
G
GRAHAM CLULEY
The Hacker News
The Hacker News
N
News and Events Feed by Topic
I
Intezer
The GitHub Blog
The GitHub Blog
S
SegmentFault 最新的问题
T
The Blog of Author Tim Ferriss
PCI Perspectives
PCI Perspectives
S
Secure Thoughts
P
Proofpoint News Feed
Microsoft Security Blog
Microsoft Security Blog
IT之家
IT之家
T
Threat Research - Cisco Blogs
J
Java Code Geeks
钛媒体:引领未来商业与生活新知
钛媒体:引领未来商业与生活新知
D
DataBreaches.Net
Hacker News - Newest:
Hacker News - Newest: "LLM"
Last Week in AI
Last Week in AI
H
Help Net Security
L
LangChain Blog
大猫的无限游戏
大猫的无限游戏
Help Net Security
Help Net Security
S
Schneier on Security
T
The Exploit Database - CXSecurity.com
Google Online Security Blog
Google Online Security Blog
Cyberwarzone
Cyberwarzone
T
Tailwind CSS Blog
V
Vulnerabilities – Threatpost
Forbes - Security
Forbes - Security
Apple Machine Learning Research
Apple Machine Learning Research
O
OpenAI News
AWS News Blog
AWS News Blog
月光博客
月光博客

郑文峰的博客

使用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. 相关链接