加载中

如何编写 Java 输入插件

要开发新的 Logstash Java 输入插件,你需要编写一个符合 Logstash Java Inputs API 的新 Java 类,对其进行打包,并使用 logstash-plugin 实用程序安装它。我们将逐步介绍这些步骤。

首先复制 示例输入插件。插件 API 目前是 Logstash 代码库的一部分,因此你必须拥有该代码库的本地副本。你可以通过以下 git 命令获取 Logstash 代码库的副本

git clone --branch <branch_name> --single-branch https://github.com/elastic/logstash.git <target_folder>
		

branch_name 应对应包含所需 Java 插件 API 版本的 Logstash 版本。

注意

Java 插件 API 的 GA 版本在 Logstash 代码库的 7.2 及更高分支中可用。

为你的 Logstash 代码库本地副本指定 target_folder。如果你不指定 target_folder,它默认为当前文件夹下的一个名为 logstash 的新文件夹。

获取适当版本的 Logstash 代码库副本后,你需要对其进行编译以生成包含 Java 插件 API 的 .jar 文件。从 Logstash 代码库的根目录 ($LS_HOME) 中,你可以使用 ./gradlew assemble(或者如果你在 Windows 上运行,则使用 gradlew.bat assemble)对其进行编译。这应该会生成 $LS_HOME/logstash-core/build/libs/logstash-core-x.y.z.jar,其中 xyz 指的是 Logstash 的版本。

成功编译 Logstash 后,你需要告诉你的 Java 插件到哪里查找 logstash-core-x.y.z.jar 文件。在插件项目的根文件夹中创建一个名为 gradle.properties 的新文件。该文件应包含一行

LOGSTASH_CORE_PATH=<target_folder>/logstash-core
		

其中 target_folder 是你的 Logstash 代码库本地副本的根文件夹。

该示例输入插件在终止前会生成可配置数量的简单事件。让我们看看该示例输入中的主类。

@LogstashPlugin(name="java_input_example")
public class JavaInputExample implements Input {

    public static final PluginConfigSpec<Long> EVENT_COUNT_CONFIG =
            PluginConfigSpec.numSetting("count", 3);

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

    private String id;
    private long count;
    private String prefix;
    private final CountDownLatch done = new CountDownLatch(1);
    private volatile boolean stopped;


    public JavaInputExample(String id, Configuration config, Context context) {
            this.id = id;
        count = config.get(EVENT_COUNT_CONFIG);
        prefix = config.get(PREFIX_CONFIG);
    }

    @Override
    public void start(Consumer<Map<String, Object>> consumer) {
        int eventCount = 0;
        try {
            while (!stopped && eventCount < count) {
                eventCount++;
                consumer.accept.push(Collections.singletonMap("message",
                        prefix + " " + StringUtils.center(eventCount + " of " + count, 20)));
            }
        } finally {
            stopped = true;
            done.countDown();
        }
    }

    @Override
    public void stop() {
        stopped = true;
    }

    @Override
    public void awaitStop() throws InterruptedException {
        done.await();
    }

    @Override
    public Collection<PluginConfigSpec<?>> configSchema() {
        return Arrays.asList(EVENT_COUNT_CONFIG, PREFIX_CONFIG);
    }

    @Override
    public String getId() {
        return this.id;
    }
}
		
  1. 设置标志以请求协作停止输入
  2. 阻塞直到输入停止

让我们逐步检查该类的每个部分。

@LogstashPlugin(name="java_input_example")
public class JavaInputExample implements Input {
		

关于类声明的注意事项

  • 所有 Java 插件都必须使用 @LogstashPlugin 注解进行标注。此外

    • 必须提供注解的 name 属性,并定义插件名称,该名称将在 Logstash 管道定义中使用。例如,此输入将在 Logstash 管道定义的输入部分中被引用为 input { java_input_example => { .... } }
    • name 属性的值必须与类名匹配(不区分大小写和下划线)。
  • 该类必须实现 co.elastic.logstash.api.Input 接口。

  • Java 插件不得在 org.logstashco.elastic.logstash 包中创建,以防止与 Logstash 本身中的类发生潜在冲突。

下面的代码片段包含了设置定义和引用它的方法。

public static final PluginConfigSpec<Long> EVENT_COUNT_CONFIG =
        PluginConfigSpec.numSetting("count", 3);

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

@Override
public Collection<PluginConfigSpec<?>> configSchema() {
    return Arrays.asList(EVENT_COUNT_CONFIG, PREFIX_CONFIG);
}
		

PluginConfigSpec 类允许开发人员指定插件支持的设置,包括设置名称、数据类型、弃用状态、必需状态和默认值。在此示例中,count 设置定义将生成的事件数量,prefix 设置定义要包含在事件字段中的可选前缀。这两个设置都不是必需的,如果未显式设置,则设置默认分别为 3message

configSchema 方法必须返回插件支持的所有设置列表。在 Java 插件项目的未来阶段,Logstash 执行引擎将验证是否存在所有必需的设置,并确保不存在不受支持的设置。

private String id;
private long count;
private String prefix;

public JavaInputExample(String id, Configuration config, Context context) {
    this.id = id;
    count = config.get(EVENT_COUNT_CONFIG);
    prefix = config.get(PREFIX_CONFIG);
}
		

所有 Java 输入插件都必须有一个包含 String id 以及 ConfigurationContext 参数的构造函数。这是将在运行时实例化它们时使用的构造函数。所有插件设置的检索和验证都应在此构造函数中进行。在此示例中,两个插件设置的值被检索并存储在局部变量中,以便稍后在 start 方法中使用。

任何额外的初始化也可以在构造函数中进行。如果在输入插件的配置或初始化过程中遇到任何不可恢复的错误,则应抛出描述性异常。该异常将被记录,并阻止 Logstash 启动。

@Override
public void start(Consumer<Map<String, Object>> consumer) {
    int eventCount = 0;
    try {
        while (!stopped && eventCount < count) {
            eventCount++;
            consumer.accept.push(Collections.singletonMap("message",
                    prefix + " " + StringUtils.center(eventCount + " of " + count, 20)));
        }
    } finally {
        stopped = true;
        done.countDown();
    }
}
		

start 方法开始输入中的事件生成循环。输入是灵活的,可以通过多种不同的机制生成事件,包括

  • 拉取机制,例如定期查询外部数据库
  • 推送机制,例如从客户端发送到本地网络端口的事件
  • 定时计算,例如心跳
  • 任何其他产生有用事件流的机制。事件流可以是有限的或无限的。如果输入产生无限的事件流,则此方法应循环执行,直到通过 stop 方法发出停止请求。如果输入产生有限的事件流,则此方法应在生成流中的最后一个事件或收到停止请求时(以先发生者为准)终止。

事件应构建为 Map<String, Object> 的实例,并通过 Consumer<Map<String, Object>>.accept() 方法推送到事件管道中。为了减少分配和 GC 压力,输入可以重复使用同一个 map 实例,只需在调用 Consumer<Map<String, Object>>.accept() 之间修改其字段即可,因为事件管道会根据 map 数据的一个副本创建事件。

private final CountDownLatch done = new CountDownLatch(1);
private volatile boolean stopped;

@Override
public void stop() {
    stopped = true;
}

@Override
public void awaitStop() throws InterruptedException {
    done.await();
}
		
  1. 设置标志以请求协作停止输入
  2. 阻塞直到输入停止

stop 方法通知输入停止生成事件。停止机制可以用任何符合 API 契约的方式实现,不过 volatile boolean 标志对许多用例来说都很有效。

输入以异步和协作方式停止。使用 awaitStop 方法进行阻塞,直到输入完成停止过程。注意,此方法**不应**像 stop 方法那样向输入发出停止信号。awaitStop 机制可以用任何符合 API 契约的方式实现,不过 CountDownLatch 对许多用例来说都很有效。

@Override
public String getId() {
    return id;
}
		

对于输入插件,getId 方法应始终返回在实例化时通过构造函数提供给插件的 id。

最后,但同样重要的是,强烈建议进行单元测试。示例输入插件包含一个 示例单元测试,你可以将其用作自己的模板。

Java 插件被打包为 Ruby gem,用于依赖管理和与 Ruby 插件的互操作性。一旦它们被打包为 gem,就可以像 Ruby 插件一样使用 logstash-plugin 实用程序进行安装。由于 Java 插件开发不需要具备 Ruby 或其工具链的知识,因此将 Java 插件打包为 Ruby gem 的过程已通过示例 Java 插件中提供的 Gradle 构建文件中的自定义任务实现了自动化。以下各节介绍如何配置和执行该打包任务,以及如何将打包好的 Java 插件安装到 Logstash 中。

以下部分出现在示例 Java 插件随附的 build.gradle 文件顶部附近

// ===========================================================================
// plugin info
// ===========================================================================
group                      'org.logstashplugins'
version                    "${file("VERSION").text.trim()}"
description                = "Example Java filter implementation"
pluginInfo.licenses        = ['Apache-2.0']
pluginInfo.longDescription = "This gem is a Logstash plugin required to be installed on top of the Logstash core pipeline using \$LS_HOME/bin/logstash-plugin install gemname. This gem is not a stand-alone program"
pluginInfo.authors         = ['Elasticsearch']
pluginInfo.email           = ['info@elastic.co']
pluginInfo.homepage        = "https://esdocs.cn/guide/en/logstash/current/index.html"
pluginInfo.pluginType      = "filter"
pluginInfo.pluginClass     = "JavaFilterExample"
pluginInfo.pluginName      = "java_filter_example"
// ===========================================================================
		
  1. 必须与主插件类的包匹配
  2. 从必需的 VERSION 文件中读取
  3. SPDX 许可证 ID 列表

你应该为你的插件配置上述值。

  • version 值将自动从插件代码库根目录中的 VERSION 文件中读取。
  • pluginInfo.pluginType 应设置为 inputfiltercodecoutput 之一。
  • pluginInfo.pluginName 必须与主插件类上 @LogstashPlugin 注解中指定的名称匹配。Gradle 打包任务将验证这一点,如果它们不匹配,将返回错误。

将插件打包为 Ruby gem 需要几个 Ruby 源文件以及一个 gemspec 文件和一个 Gemfile。这些 Ruby 文件仅用于定义 Ruby gem 结构或在 Logstash 启动时注册 Java 插件。它们在运行时事件处理期间不会被使用。Gradle 打包任务会根据上述部分中配置的值自动生成所有这些文件。

你可以使用以下命令运行 Gradle 打包任务

./gradlew gem
		

对于 Windows 平台:根据需要在命令中将 ./gradlew 替换为 gradlew.bat

该任务将在你的插件代码库的根目录中生成一个 gem 文件,文件名为 logstash-{{plugintype}}-<pluginName>-<version>.gem

将 Java 插件打包为 Ruby gem 后,你可以使用以下命令将其安装在 Logstash 中

bin/logstash-plugin install --no-verify --local /path/to/javaPlugin.gem
		

对于 Windows 平台:根据需要在命令中用反斜杠替换正斜杠。

以下是一个最小的 Logstash 配置,可用于测试 Java 输入插件是否已正确安装并正常运行。

input {
  java_input_example {}
}
output {
  stdout { codec => rubydebug }
}
		

将上述 Logstash 配置复制到文件(例如 java_input.conf)。启动 Logstash 并使用

bin/logstash -f /path/to/java_input.conf
		

使用上述配置时,预期的 Logstash 输出(不包括初始化信息)为

{
      "@version" => "1",
       "message" => "message        1 of 3       ",
    "@timestamp" => yyyy-MM-ddThh:mm:ss.SSSZ
}
{
      "@version" => "1",
       "message" => "message        2 of 3       ",
    "@timestamp" => yyyy-MM-ddThh:mm:ss.SSSZ
}
{
      "@version" => "1",
       "message" => "message        3 of 3       ",
    "@timestamp" => yyyy-MM-ddThh:mm:ss.SSSZ
}
		

如果你对 Logstash 中的 Java 插件支持有任何反馈,请在我们的 GitHub 主问题页面上发表评论或在 Logstash 论坛上发帖。

© . This website operates independently and is not affiliated with or endorsed by Elasticsearch B.V. All brand names, logos, and trademarks are the property of their respective owners.