如何在自定义Logstash Java过滤插件中超时后生成并推送新事件
自定义Logstash Java过滤器异步发送事件实现方案
你需要的在filter方法外生成并发送事件的需求是可行的,不需要依赖额外外部插件,以下是具体实现方案:
核心实现逻辑
Logstash的Filter接口没有提供异步主动推送事件的公开API,所有事件都需要通过filter方法的返回集合传递到下一个管道节点,你可以通过「待发送事件队列中转」的方式实现你的需求。
1. 调整插件结构适配异步逻辑
首先在插件类中新增线程安全的存储结构,同时在构造器中启动你的后台处理线程,所有共享变量都要做并发安全处理:
// 新增成员变量 private final Queue<Event> pendingEmitEvents = new ConcurrentLinkedQueue<>(); private Thread backgroundProcessThread; private volatile boolean isRunning = true; // 构造器调整 public JavaFilterExample(String id, Configuration config, Context context) { this.id = id; this.sourceField = config.get(SOURCE_CONFIG); // 原buffer改为线程安全实现 this.buffer = Collections.synchronizedList(new ArrayList<>()); // 启动后台处理线程 this.backgroundProcessThread = new Thread(() -> { while (isRunning) { try { // 这里写你的超时检测、多消息关联逻辑 boolean triggerCondition = /* 满足超时条件 || buffer包含finish标识 */; if (triggerCondition) { // 构建新的关联事件 Event aggregatedEvent = new org.logstash.bail.BasicEvent(); aggregatedEvent.setField("aggregated_content", String.join(",", buffer)); aggregatedEvent.setField("event_count", buffer.size()); // 事件放到待发送队列 pendingEmitEvents.add(aggregatedEvent); buffer.clear(); } Thread.sleep(1000); // 按需调整检测间隔 } catch (InterruptedException e) { Thread.currentThread().interrupt(); break; } } }); backgroundProcessThread.start(); }
2. 改造filter方法返回新增事件
每次filter方法被调用时,除了处理当前输入的事件,还要把待发送队列里的所有事件一起返回给管道:
@Override public Collection<Event> filter(Collection<Event> events, FilterMatchListener matchListener) { List<Event> outputEvents = new ArrayList<>(events); // 处理当前输入的事件 for (Event e : events) { Object f = e.getField(sourceField); if (f instanceof String) { String content = (String) f; buffer.add(content); matchListener.filterMatched(e); // 收到finish标识立刻触发关联 if ("finish".equals(content)) { Event aggregatedEvent = new org.logstash.bail.BasicEvent(); aggregatedEvent.setField("aggregated_content", String.join(",", buffer)); aggregatedEvent.setField("trigger_type", "finish_flag"); pendingEmitEvents.add(aggregatedEvent); buffer.clear(); } } } // 把待发送队列里的所有事件加到返回集合,传递到下一个过滤器 Event event; while ((event = pendingEmitEvents.poll()) != null) { outputEvents.add(event); } return outputEvents; }
3. 补充资源释放逻辑
重写close方法,停止后台线程避免内存泄漏,同时把插件关闭前剩余的缓存数据也生成事件输出:
@Override public void close() { isRunning = false; if (backgroundProcessThread != null) { backgroundProcessThread.interrupt(); try { backgroundProcessThread.join(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } // 关闭前处理剩余缓存 if (!buffer.isEmpty()) { Event lastEvent = new org.logstash.bail.BasicEvent(); lastEvent.setField("aggregated_content", String.join(",", buffer)); lastEvent.setField("trigger_type", "plugin_shutdown"); pendingEmitEvents.add(lastEvent); } }
注意事项
- 你提到的直接在外部类调用
Logstash.sendEvent的方式没有公开API支持,上述队列中转的方案是官方自定义插件场景下的标准实现方式 - 如果不需要保留原始输入事件,可以在处理完成后把原始事件从
outputEvents中移除,只返回你生成的关联事件 - 新创建的事件要使用
org.logstash.bail.BasicEvent实现类,不要自行实现Event接口
内容的提问来源于stack exchange,提问作者oliver
相关产品推荐
相关产品推荐

