跳转到主要内容

创建扩展

Druid 使用模块化系统,允许在运行时添加扩展。

编写自己的扩展

Druid 的扩展利用 Guice 在运行时添加功能。简单来说,Guice 是一个依赖注入(Dependency Injection)框架,我们使用它来维护 Druid 进程的对象图。扩展可以通过添加 Guice 绑定,对对象图进行任何它们想要/需要的更改。虽然扩展确实赋予了你几乎随心所欲更改任何内容的权限,但通常我们希望用户扩展下列列出的内容。这意味着我们承诺对本页面提到的接口变更遵循我们的版本控制策略,而其他接口被视为“内部”接口,即使在补丁版本之间也可能以不兼容的方式进行更改。

  1. 通过扩展 org.apache.druid.segment.loading.DataSegment*org.apache.druid.tasklogs.TaskLog* 类来添加新的深度存储实现。
  2. 通过扩展 org.apache.druid.data.input.InputSource 来添加新的输入源。
  3. 通过扩展 org.apache.druid.data.input.InputEntity 来添加新的输入实体。
  4. 如有必要,通过扩展 org.apache.druid.data.input.InputSourceReader 来添加新的输入源读取器。在大多数情况下,你可以使用 org.apache.druid.data.input.impl.InputEntityIteratingReader
  5. 通过扩展 org.apache.druid.data.input.InputFormat 来添加新的输入格式。
  6. 通过扩展用于文本格式的 org.apache.druid.data.input.TextReader 或用于二进制格式的 org.apache.druid.data.input.IntermediateRowParsingReader 来添加新的输入实体读取器。
  7. 通过扩展 org.apache.druid.query.aggregation.AggregatorFactoryorg.apache.druid.query.aggregation.Aggregatororg.apache.druid.query.aggregation.BufferAggregator 来添加聚合器。
  8. 通过扩展 org.apache.druid.query.aggregation.PostAggregator 来添加后聚合器。
  9. 通过扩展 org.apache.druid.query.extraction.ExtractionFn 来添加提取函数。
  10. 通过扩展 org.apache.druid.segment.serde.ComplexMetricSerde 来添加复杂指标。
  11. 通过扩展 org.apache.druid.query.QueryRunnerFactoryorg.apache.druid.query.QueryToolChestorg.apache.druid.query.Query 来添加新的查询类型。
  12. 通过调用 Jerseys.addResource(binder, clazz) 来添加新的 Jersey 资源。
  13. 通过扩展 org.apache.druid.server.initialization.jetty.ServletFilterHolder 来添加新的 Jetty 过滤器。
  14. 通过扩展 org.apache.druid.metadata.PasswordProvider 来添加新的密码提供程序。
  15. 通过扩展 org.apache.druid.metadata.DynamicConfigProvider 来添加新的动态配置提供程序。
  16. 通过实现 druid-processing 包中的 org.apache.druid.segment.transform.Transform 接口来添加新的摄取转换。
  17. 将你的扩展与其他 Druid 扩展捆绑在一起

扩展通过实现 org.apache.druid.initialization.DruidModule 来添加到系统中。

创建 Druid 模块

DruidModule 类有两个方法

  1. 一个 configure(Binder) 方法
  2. 一个 getJacksonModules() 方法

configure(Binder) 方法与普通 Guice 模块拥有的方法相同。

getJacksonModules() 方法提供了一个 Jackson 模块列表,用于帮助初始化 Druid 所使用的 Jackson ObjectMapper 实例。这就是你如何将通过 Jackson 实例化的扩展(如 AggregatorFactory 和 InputSource 对象)添加到 Druid 的方式。

注册你的 Druid 模块

一旦创建了 DruidModule,你需要在 jar 包的 META-INF/services 目录中打包一个额外文件。对于 Maven 项目,最简单的方法是在 src/main/resources 目录中创建文件。在 Druid 代码库中的 cassandra-storagehdfs-storages3-extensions 模块下有相关示例。

你的 jar 包中应该存在的文件是

META-INF/services/org.apache.druid.initialization.DruidModule

它应该是一个文本文件,包含换行符分隔的实现了 DruidModule 的包限定类名列表,例如

org.apache.druid.storage.cassandra.CassandraDruidModule

如果你的 jar 包包含此文件,那么当它被添加到类路径或作为扩展加载时,Druid 会识别该文件并实例化该模块。你的模块应该有一个默认的构造函数;但如果你需要访问运行时配置属性,可以添加一个带有 @Inject 注解的方法,从 Guice 中注入一个 Properties 对象。

添加新的深度存储实现

查看 druid-azure-extensionsdruid-google-extensionsdruid-cassandra-storagedruid-hdfs-storagedruid-s3-extensions 模块,了解如何执行此操作的示例。

扩展的基本思想是,你需要为你的 DataSegmentPusherURIDataPuller 对象添加绑定。添加它们的方式如下(摘自 HdfsStorageDruidModule)

Binders.dataSegmentPullerBinder(binder)
.addBinding("hdfs")
.to(HdfsDataSegmentPuller.class).in(LazySingleton.class);

Binders.dataSegmentPusherBinder(binder)
.addBinding("hdfs")
.to(HdfsDataSegmentPusher.class).in(LazySingleton.class);

Binders.dataSegment*Binder() 是 druid-core jar 提供的一个调用,它设置了一个 Guice multibind "MapBinder"。如果这听起来很困惑,不必担心;只需将其视为一个神奇的咒语。

对于 Puller 绑定器,addBinding("hdfs") 为类型为 "hdfs" 的 loadSpec 对象创建了一个新的处理程序。对于 Pusher 绑定器,它为 druid.storage.type 参数创建了一个你可以指定的新类型值。

to(...).in(...); 是标准的 Guice 用法。

除了 DataSegmentPusherURIDataPuller,你还可以绑定

  • DataSegmentKiller:移除段,用作 Kill 任务的一部分以删除未使用的段,即对被更新版本取代或从集群中丢弃的段执行垃圾回收。
  • DataSegmentMover:允许将段从一个位置迁移到另一个位置,目前这仅用作 MoveTask 的一部分,将未使用的段移动到不同的 S3 存储桶或前缀,通常是为了降低未使用数据的存储成本(例如移至 Glacier 或更便宜的存储)。
  • DataSegmentArchiver:只是 Mover 的包装器,但预配置了目标存储桶/路径,因此不需要在 ArchiveTask 运行时进行指定。

验证你的深度存储实现

警告! 这不是一个正式的程序,而是一组验证你的新深度存储实现是否能够推送、拉取和删除段的提示。

建议使用批量摄取任务来验证你的实现。段会在约 1 分钟后自动汇总到 Historical 节点。通过这种方式,你可以验证段的推送(在实时处理进程中)和拉取(在 Historical 进程中)。

DataSegmentPusher

无论你的数据存储(云存储服务、分布式文件系统等)在哪里,在摄取任务结束后,你应该能够看到一个新文件:index.zip(对于 HDFS 数据存储为 partitionNum_index.zip)。

URIDataPuller

在摄取任务结束约 1 分钟后,你应该能够看到你的 Historical 进程尝试加载新段。

DataSegmentKiller

测试段删除的最简单方法是将段标记为未使用,然后在 Web 控制台中启动删除任务(killing task)。

要将段标记为未使用,你需要连接到元数据存储,并将段表行上的 used 列更新为 false

要启动段删除任务,你需要访问 Web 控制台,然后为相应的数据源选择 issue kill task

删除任务结束后,index.zip(对于 HDFS 数据存储为 partitionNum_index.zip)文件应该会从数据存储中删除。

添加对新输入源的支持

添加对新输入源的支持需要实现三个接口:InputSourceInputEntityInputSourceReaderInputSource 用于定义输入数据存储在哪里。InputEntity 用于定义如何在原生并行索引中并行读取数据。InputSourceReader 定义了如何读取你的新输入源,在大多数情况下,你可以直接使用提供的 InputEntityIteratingReader

druid-s3-extensions 模块中的 S3InputSourceS3Entity 有一个相关示例。

添加 InputSource 几乎完全通过 Jackson 模块而不是 Guice 完成。请特别注意实现

@Override
public List<? extends Module> getJacksonModules()
{
return ImmutableList.of(
new SimpleModule().registerSubtypes(new NamedType(S3InputSource.class, "s3"))
);
}

这是将 InputSource 注册到 Jackson 的多态序列化/反序列化层。具体来说,有了这个,意味着如果你在 IO 配置中指定 "inputSource": { "type": "s3", ... },系统就会为你的 InputSource 实现加载此 InputSource。

请注意,在 Druid 内部,我们已经使 Jackson 反序列化对象的 @JacksonInject 注解实际上使用底层的 Guice 注入器来解析要注入的对象。因此,如果你的 InputSource 需要访问某个对象,你可以在 setter 上添加 @JacksonInject 注解,它将在实例化时被设置。

添加对新数据格式的支持

添加对新数据格式的支持需要实现两个接口:InputFormatInputEntityReaderInputFormat 用于定义数据格式。InputEntityReader 用于定义如何解析数据并将其转换为 Druid InputRow

druid-orc-extensions 模块中的 OrcInputFormatOrcReader 有一个示例。

添加 InputFormat 与添加 InputSource 非常相似。它们完全通过 Jackson 运行,因此只需添加到 DruidModule 返回的 Jackson 模块中即可。

添加聚合器

添加 AggregatorFactory 对象与 InputSource 对象非常相似。它们完全通过 Jackson 运行,因此只需添加到 DruidModule 返回的 Jackson 模块中即可。

添加复杂指标

在当前版本中,添加 ComplexMetrics 有点繁琐。获取复杂指标的方法是通过 ComplexMetrics.registerSerde() 方法注册。要使其工作,没有特殊的 Guice 操作,只需在 configure(Binder) 方法中注册序列化/反序列化即可。

添加新的查询类型

添加新的查询类型需要实现三个接口。

  1. org.apache.druid.query.Query
  2. org.apache.druid.query.QueryToolChest
  3. org.apache.druid.query.QueryRunnerFactory

注册这些接口所使用的通用策略与深度存储机制相同。你需要执行类似下面的操作

DruidBinders.queryToolChestBinder(binder)
.addBinding(SegmentMetadataQuery.class)
.to(SegmentMetadataQueryQueryToolChest.class);

DruidBinders.queryRunnerFactoryBinder(binder)
.addBinding(SegmentMetadataQuery.class)
.to(SegmentMetadataQueryRunnerFactory.class);

第一个绑定 SegmentMetadataQueryQueryToolChest 以便在 SegmentMetadataQuery 被使用时调用。第二个对 QueryRunnerFactory 执行同样的操作。

添加新的 Jersey 资源

向模块添加新的 Jersey 资源需要调用以下代码来绑定模块中的资源

Jerseys.addResource(binder, NewResource.class);

添加新的密码提供程序实现

你需要实现 org.apache.druid.metadata.PasswordProvider 接口。Druid 使用 PasswordProvider 的每个地方都会创建一个实现的新实例,因此请确保在对象实例化期间提供了获取每个密码所需的所有必要信息。在你的 org.apache.druid.initialization.DruidModule 实现中,getJacksonModules 应该如下所示 -

    return ImmutableList.of(
new SimpleModule("SomePasswordProviderModule")
.registerSubtypes(
new NamedType(SomePasswordProvider.class, "some")
)
);

其中 SomePasswordProviderPasswordProvider 接口的实现,你可以参考 org.apache.druid.metadata.EnvironmentVariablePasswordProvider

添加新的 DynamicConfigProvider 实现

你需要实现 org.apache.druid.metadata.DynamicConfigProvider 接口。Druid 使用 DynamicConfigProvider 的每个地方都会创建一个实现的新实例,因此请确保在对象实例化期间提供了获取所有信息所需的所有必要信息。在你的 org.apache.druid.initialization.DruidModule 实现中,getJacksonModules 应该如下所示 -

    return ImmutableList.of(
new SimpleModule("SomeDynamicConfigProviderModule")
.registerSubtypes(
new NamedType(SomeDynamicConfigProvider.class, "some")
)
);

其中 SomeDynamicConfigProviderDynamicConfigProvider 接口的实现,你可以参考 org.apache.druid.metadata.MapStringDynamicConfigProvider

添加转换扩展

要创建转换扩展,请实现 org.apache.druid.segment.transform.Transform 接口。你需要安装 druid-processing 包来导入 org.apache.druid.segment.transform

import com.fasterxml.jackson.annotation.JsonCreator;
import com.fasterxml.jackson.annotation.JsonProperty;
import org.apache.druid.segment.transform.RowFunction;
import org.apache.druid.segment.transform.Transform;

public class MyTransform implements Transform {
private final String name;

@JsonCreator
public MyTransform(
@JsonProperty("name") final String name
) {
this.name = name;
}

@JsonProperty
@Override
public String getName() {
return name;
}

@Override
public RowFunction getRowFunction() {
return new MyRowFunction();
}

static class MyRowFunction implements RowFunction {
@Override
public Object eval(Row row) {
return "transformed-value";
}
}
}

然后将你的转换注册为 Jackson 模块。

import com.fasterxml.jackson.databind.Module;
import com.fasterxml.jackson.databind.jsontype.NamedModule;
import com.fasterxml.jackson.databind.module.SimpleModule;
import com.google.inject.Binder;
import com.google.common.collect.ImmutableList;
import org.apache.druid.initialization.DruidModule;

public class MyTransformModule implements DruidModule {
@Override
public List<? extends Module> getJacksonModules() {
return return ImmutableList.of(
new SimpleModule("MyTransformModule").registerSubtypes(
new NamedType(MyTransform.class, "my-transform")
)
):
}

@Override
public void configure(Binder binder) {
}
}

添加你自己的自定义可插拔 Coordinator Duty

协调器定期运行作业,即所谓的 CoordinatorDuty,其中包括加载新段、段平衡等。Druid 用户可以在不修改任何 Core Druid 类的情况下添加自定义的可插拔协调器职责(不属于 Core Druid 的部分)。用户可以通过编写实现 CoordinatorCustomDuty 接口并设置 JsonTypeName 的自定义协调器职责来做到这一点。接下来,用户需要在其模块的 DruidModule#getJacksonModules() 中将自定义协调器注册为子类型。完成这些步骤后,用户将能够使用以下属性加载其自定义协调器职责

druid.coordinator.dutyGroups=[<GROUP_NAME_1>, <GROUP_NAME_2>, ...]
druid.coordinator.<GROUP_NAME_1>.duties=[<DUTY_NAME_MATCHING_JSON_TYPE_NAME_1>, <DUTY_NAME_MATCHING_JSON_TYPE_NAME_2>, ...]
druid.coordinator.<GROUP_NAME_1>.period=<GROUP_NAME_1_RUN_PERIOD>

druid.coordinator.<GROUP_NAME_1>.duty.<DUTY_NAME_MATCHING_JSON_TYPE_NAME_1>.<SOME_CONFIG_1_KEY>=<SOME_CONFIG_1_VALUE>
druid.coordinator.<GROUP_NAME_1>.duty.<DUTY_NAME_MATCHING_JSON_TYPE_NAME_1>.<SOME_CONFIG_2_KEY>=<SOME_CONFIG_2_VALUE>

在新的可插拔 Coordinator 职责系统中,职责可以像今天协调器已经做的那样进行分组。职责将根据 druid.coordinator.dutyGroups 列表中的元素分为多个组。同一组中的所有职责将具有相同的运行周期,由 druid.coordinator.<GROUP_NAME>.period 配置。目前,有一个单线程为每个组顺序运行职责。

例如,请参阅 KillSupervisorsCustomDuty 获取自定义协调器职责实现,并查看 KillSupervisorsCustomDutyTest 了解可用于配置 KillSupervisorsCustomDuty 的示例属性。

通过 HTTP 代理为你的扩展路由数据

你可以为扩展的 HttpClient 添加通过 HTTP 代理连接的能力。

要支持扩展 HTTP 客户端的代理连接

  1. HttpClientProxyConfig 作为 @JsonProperty 添加到扩展的 HTTP 配置类中。
  2. 在扩展的模块类中,将 HttpProxyConfig 配置添加到 HttpClientConfig。例如,其中 config 变量是步骤 1 中的扩展 HTTP 配置
final HttpClientConfig.Builder builder = HttpClientConfig
.builder()
.withNumConnections(1)
.withReadTimeout(config.getReadTimeout().toStandardDuration())
.withHttpProxyConfig(config.getProxyConfig());

将你的扩展与其他 Druid 扩展捆绑在一起

当你执行 mvn install 时,Druid 扩展将被打包在 Druid tarball 和 extensions 目录中,它们都位于 distribution/target/ 下。

如果你想包含你的扩展,可以将你的扩展 Maven 坐标添加为 distribution/pom.xml 中的参数

mvn install 期间,Maven 会将你的扩展安装到本地 Maven 仓库,然后调用 pull-deps 从那里拉取你的扩展。最终,你应该能在 distribution/target/extensions 下和 Druid tarball 中看到你的扩展。

管理依赖

对于引入常用库的扩展来说,管理库冲突可能令人望而生畏。以下是建议指定为 provided 范围的库组 ID 列表,以防止与 Druid 中使用的版本发生冲突

"org.apache.druid",
"com.metamx.druid",
"asm",
"org.ow2.asm",
"org.jboss.netty",
"com.google.guava",
"com.google.code.findbugs",
"com.google.protobuf",
"com.esotericsoftware.minlog",
"log4j",
"org.slf4j",
"commons-logging",
"org.eclipse.jetty",
"org.mortbay.jetty",
"com.sun.jersey",
"com.sun.jersey.contribs",
"common-beanutils",
"commons-codec",
"commons-lang",
"commons-cli",
"commons-io",
"javax.activation",
"org.apache.httpcomponents",
"org.apache.zookeeper",
"org.codehaus.jackson",
"com.fasterxml.jackson",
"com.fasterxml.jackson.core",
"com.fasterxml.jackson.dataformat",
"com.fasterxml.jackson.datatype",
"org.roaringbitmap",
"net.java.dev.jets3t"

有关详细信息,请参阅 org.apache.druid.cli.PullDependencies 中的文档。