news 2026/9/23 10:27:07

Apache Druid InfluxDB Emitter 扩展实战:将 Druid 服务指标实时写入 InfluxDB

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Druid InfluxDB Emitter 扩展实战:将 Druid 服务指标实时写入 InfluxDB
  • 数据库
  • OLAP
  • 大数据
  • 后端

【免费下载链接】druid

Apache Druid: a high performance real-time analytics database.

项目地址:https://gitcode.com/gh_mirrors/druid6/druid
点击查看免费下载

本篇技术指南以 Apache Druid 仓库中的社区扩展druid-influxdb-emitter为核心(官方文档),系统讲解如何将该扩展加载进 Druid、配置全部参数、理解其基于队列的异步发送模型,并深入剖析 Druid 指标事件到 InfluxDB Line Protocol 的转换规则。读完本文,你将能够在自己的 Druid 集群中启用该扩展,把服务指标(Service Metric Events)稳定地投递到 InfluxDB 时序数据库,供 Grafana 等可视化平台进一步分析与告警。

扩展简介与工作流程

druid-influxdb-emitter是一个社区贡献的 Druid Emitter 扩展,作用是把 Druid 产生的监控指标通过 HTTP 协议发送到 InfluxDB 时序数据库。与 Druid 内置的 logging、http 等 emitter 相比,它的目标消费者是 InfluxDB,因此数据以 InfluxDB 原生的Line Protocol(行协议)文本格式投递。

从源码实现看(InfluxdbEmitter.java),该扩展的工作流程分为四个阶段:

  1. 接收事件入队:Druid 的监控框架每产生一个事件,emit(Event event)方法被调用。该方法只处理ServiceMetricEvent类型的事件(其余事件类型被忽略),并将其放入一个有界阻塞队列LinkedBlockingQueue(容量由maxQueueSize决定),见 InfluxdbEmitter.java。
  2. 定时批量转换start()启动一个单线程的ScheduledExecutorService,按照flushDelay(首次延迟)和flushPeriod(固定周期)调度transformAndSendToInfluxdb(),见 InfluxdbEmitter.java。
  3. 全量冲刷发送transformAndSendToInfluxdb()一次性取出当前队列中的全部事件(eventsQueue.size()决定取出数量,逐个poll()),逐条转换为 Line Protocol 后拼接,再 POST 到 InfluxDB 的 HTTP 写接口,见 InfluxdbEmitter.java。整个队列在一次发送中全部清空,而不是逐条发送。
  4. 关闭时兜底冲刷close()被调用(进程关闭/优雅退出)时,会先调用flush()将队列中剩余事件全部发送,随后关闭调度线程,见 InfluxdbEmitter.java。

注意:该扩展目前只发送 Service Metric Events(服务指标事件),即 Druid metrics 文档中列举的那一类指标;其他 feed(如 alerts)不会进入该队列。

加载扩展

druid-influxdb-emitter属于社区扩展(Community Extension),不随默认 Druid 发行包(tarball)打包。要使用它,需要先把它下载安装到 Druid 的extensions目录,再在common.runtime.propertiesdruid.extensions.loadList中加入扩展名。

根据 扩展加载文档,社区扩展通常可以通过pull-deps工具按 Maven 坐标下载。对于该扩展,其 Maven groupId 为org.apache.druid.extensions.contrib,artifactId 为druid-influxdb-emitter,版本为当前 Druid 稳定版本:

java \ -cp "lib/*" \ -Ddruid.extensions.directory="extensions" \ -Ddruid.extensions.hadoopDependenciesDir="hadoop-dependencies" \ org.apache.druid.cli.Main tools pull-deps \ --no-default-hadoop \ -c "org.apache.druid.extensions.contrib:druid-influxdb-emitter:<DRUID_VERSION>"

安装完成后,在common.runtime.properties中加载扩展并启用该 emitter:

druid.extensions.loadList=["druid-influxdb-emitter", ...其他扩展] druid.emitter=influxdb

其中druid.emitter=influxdb用于激活该 emitter 模块(对应源码中的 InfluxdbEmitterModule.java,其内部常量EMITTER_TYPE = "influxdb"),druid.emitter的取值决定使用哪个 emitter 模块,可参考 配置总览。

配置参数详解

该扩展的全部配置项都位于druid.emitter.influxdb前缀之下,写入common.runtime.properties。下表完整列出了所有参数(来自 官方文档):

属性说明是否必填默认值
druid.emitter.influxdb.hostnameInfluxDB 服务器的主机名
druid.emitter.influxdb.portInfluxDB 服务器端口8086
druid.emitter.influxdb.protocol发送指标使用的协议,httphttps二选一http
druid.emitter.influxdb.trustStorePathhttps 场景下使用的 trustStore 路径
druid.emitter.influxdb.trustStoreTypehttps 场景下使用的 trustStore 类型jks(实际取java.security.KeyStore.getDefaultType()
druid.emitter.influxdb.trustStorePasswordhttps 场景下使用的 trustStore 密码
druid.emitter.influxdb.databaseNameInfluxDB 中的数据库名
druid.emitter.influxdb.maxQueueSize存放待发送事件的队列容量Integer.MAX_VALUE(2^31-1)
druid.emitter.influxdb.flushPeriod每隔多少毫秒把事件队列解析为 Line Protocol 并 POST 到 InfluxDB60000
druid.emitter.influxdb.flushDelay定时任务首次执行前等待多少毫秒60000
druid.emitter.influxdb.influxdbUserName访问 InfluxDB 数据库的用户名
druid.emitter.influxdb.influxdbPassword该用户的密码
druid.emitter.influxdb.dimensionWhitelist允许作为 tag 写入 Line Protocol 的维度白名单["dataSource","type","numMetrics","numDimensions","threshold","dimension","taskType","taskStatus","tier"]

默认值与源码验证

上述默认值并非仅见于文档,在配置类 InfluxdbEmitterConfig.java 中有直接对应的常量定义:

private static final int DEFAULT_PORT = 8086; private static final int DEFAULT_QUEUE_SIZE = Integer.MAX_VALUE; private static final int DEFAULT_FLUSH_PERIOD = 60000; // milliseconds private static final List<String> DEFAULT_DIMENSION_WHITELIST = Arrays.asList( "dataSource", "type", "numMetrics", "numDimensions", "threshold", "dimension", "taskType", "taskStatus", "tier"); private static final String DEFAULT_PROTOCOL = "http";

几个值得注意的细节:

  • hostnamedatabaseNameinfluxdbUserNameinfluxdbPassword四个字段在构造函数中通过Preconditions.checkNotNull强制校验,缺失任何一个都会在启动阶段直接抛出异常(配置类构造失败),相关行为被 InfluxdbEmitterConfigTest.java 的testConfigWithNullHostnametestConfigWithNullInfluxdbUserNametestConfigWithNullInfluxdbPassword等用例覆盖。
  • flushDelay未配置时默认值同样取 60000(与flushPeriod相同),即首次冲刷延迟 1 分钟。
  • trustStoreType未配置时取KeyStore.getDefaultType(),在绝大多数 JDK 上即jks,与文档表格中标注的jks一致。
  • port未配置时回退为 8086(InfluxDB 默认 HTTP 端口),对应测试testConfigWithNullPort

一个完整的配置示例

# 加载扩展并启用 influxdb emitter druid.extensions.loadList=["druid-influxdb-emitter"] druid.emitter=influxdb # InfluxDB 连接与认证 druid.emitter.influxdb.hostname=influxdb.example.com druid.emitter.influxdb.port=8086 druid.emitter.influxdb.protocol=http druid.emitter.influxdb.databaseName=druid_metrics druid.emitter.influxdb.influxdbUserName=druid_writer druid.emitter.influxdb.influxdbPassword=your_password_here # 队列与冲刷策略(可选,以下为默认值) druid.emitter.influxdb.maxQueueSize=2147483647 druid.emitter.influxdb.flushPeriod=60000 druid.emitter.influxdb.flushDelay=60000 # 维度白名单(可选,以下为默认值) druid.emitter.influxdb.dimensionWhitelist=["dataSource","type","numMetrics","numDimensions","threshold","dimension","taskType","taskStatus","tier"]

参数调优建议(基于实现机制)

  • flushPeriod / maxQueueSize 的取舍:因为发送逻辑是"整队列一次性清空",maxQueueSize相当于高水位缓冲。若指标量大而flushPeriod过长,队列可能积压较多事件,单次 HTTP 请求体随之变大;反之flushPeriod过短会增大对 InfluxDB 的请求频率。生产环境建议结合 Druid 侧指标量观察 InfluxDB 的写入负载后调整。
  • HTTPS 场景protocol=https时,trustStorePathtrustStorePassword二者必须同时提供,否则构造客户端时会抛出IllegalStateException("Can't load TrustStore. Truststore path or password is not set."),这一点在 InfluxdbEmitter.java 的buildInfluxdbClient()中实现,并被InfluxdbEmitterTest的三个异常用例验证(缺路径、缺密码、路径无效均抛异常)。

InfluxDB 侧前置要求:启用认证与授权

使用该扩展前,必须在InfluxDB 服务端启用认证与授权(authentication and authorization)。因为 emitter 在构造写入请求时会把用户名密码直接拼进 HTTP 查询串:

POST /write?db=<databaseName>&u=<influxdbUserName>&p=<influxdbPassword>

对应源码见 InfluxdbEmitter.java 的postToInflux()。若 InfluxDB 未开启认证,该用户名/密码参数会被忽略,但按照官方文档说明,正确姿势仍是先在 InfluxDB 侧创建具备对应数据库写入权限的用户,再配置到 Druid 中。

Line Protocol 转换规则(核心机制)

InfluxDB 通过 HTTP 写入数据时使用Line Protocol文本格式。其语法为:

<measurement>[,<tag_key>=<tag_value>[,<tag_key>=<tag_value>]] <field_key>=<field_value>[,<field_key>=<field_value>] [<timestamp>]

其中<timestamp>自 epoch 起的纳秒数

转换规则

Druid 的 Service Metric Event 中,metric字段是一个以/分隔的多段字符串(例如query/cache/total/hits)。扩展按照以下规则把它映射为 Line Protocol(见 InfluxdbEmitter.java 的transformForInfluxSystems()):

  • measurement(测量名)druid_+ metric 的第一段。例如query/cache/total/hitsdruid_query
  • tags(标签)
    • service:事件中的 service 字段原样保留(如druid/historical);
    • metric:metric 的中间段(去掉首尾段),用_拼接,并冠以druid_前缀。例如query/cache/total/hitsdruid_cache_total;如果 metric 只有两段(如query/time),则没有中间段,也就不生成 metric tag
    • hostname:事件中的 host 字段去掉端口部分(按:切分取第一段),如historical001:8083historical001
    • 白名单维度:事件携带的用户维度(user dims)中,凡是命中dimensionWhitelist的维度,都会追加为 tag。
  • field(字段)druid_+ metric 的最后一段,值为事件中的 value。例如query/cache/total/hitsdruid_hits=34787256
  • timestamp(时间戳):事件创建时间(event.getCreatedTime())换算为纳秒:毫秒值 × 1000000。

官方文档示例

一个由 Druid logging emitter 记录的典型服务指标事件:

Event [{"feed":"metrics","timestamp":"2017-10-31T09:09:06.857Z","service":"druid/historical","host":"historical001:8083","version":"0.11.0-SNAPSHOT","metric":"query/cache/total/hits","value":34787256}]

按上述规则转换后,得到可直接 POST 到 InfluxDB 的字符串:

druid_query,service=druid/historical,hostname=historical001,metric=druid_cache_total druid_hits=34787256 1509440946857000000

逐段拆解:

组成部分内容来源
measurementdruid_querymetric 首段query
tagservice=druid/historical事件 service
taghostname=historical001事件 host 去掉端口
tagmetric=druid_cache_totalmetric 中间段cache/total_拼接并加druid_前缀
fielddruid_hits=34787256metric 末段hits+ 事件 value
timestamp1509440946857000000事件时间戳的纳秒表示

注意:每条转换后的记录末尾会追加一个换行符\n(源码中payload.append(StringUtils.format(" %d\n", ...))),多条事件拼接后形成多行 Line Protocol 主体,一次性 POST。

源码级验证:测试用例即规范

仓库中的单元测试 InfluxdbEmitterTest.java 用实际断言固化了上述规则,是最直观的"可执行文档":

  • testTransformForInfluxWithLongMetric:metric 为metric/te/st/value(4 段),期望输出druid_metric,service=druid/historical,metric=druid_te_st,hostname=localhost,dataSource=test_datasource druid_value=1234 1509357600000000000,同时验证了命中白名单的dataSource维度被追加为 tag,而未在白名单中的nonWhiteListedDim维度被丢弃。
  • testTransformForInfluxWithShortMetric:metric 为metric/time(2 段),期望输出druid_metric,service=druid/historical,hostname=localhost druid_time=1234 ...,验证两段 metric 不产生 metric tag的规则。
  • testMetricIsInDimensionWhitelist/testMetricIsInDefaultDimensionWhitelist:分别验证自定义白名单与默认白名单下维度 tag 的生成行为。

维度值清洗规则

Line Protocol 的 tag 值中不允许出现点号(.)和空白字符,因此 emitter 对写入的维度值做了统一清洗:sanitize()方法使用正则[\s]+|[.]+(连续空白或连续点号)将所有匹配字符替换为下划线_,见 InfluxdbEmitter.java 与 InfluxdbEmitter.java。

维度白名单(dimensionWhitelist)

Druid 的事件可能携带大量用户自定义维度(user dims),如果全部写入 Line Protocol,会导致 InfluxDB 的 tag 基数爆炸,显著影响写入性能与查询效率。因此该扩展引入了白名单机制:

  • 事件中携带的维度,只有当维度名在dimensionWhitelist中时,才会作为 tag 追加到 Line Protocol;
  • 白名单默认包含:dataSourcetypenumMetricsnumDimensionsthresholddimensiontaskTypetaskStatustier
  • 可通过配置druid.emitter.influxdb.dimensionWhitelist覆盖默认值(JSON 数组格式),未配置时使用默认集合(见 InfluxdbEmitterConfig.java)。

实际判断逻辑位于transformForInfluxSystems()

for (String dimName : dimNames) { if (this.dimensionWhiteList.contains(dimName)) { tag.append(StringUtils.format(",%1$s=%2$s", dimName, sanitize(String.valueOf(event.getUserDims().get(dimName))))); } }

即:遍历事件全部用户维度,仅对白名单命中的维度执行"追加 tag + 值清洗"两步操作。

HTTPS 支持与 TrustStore 配置

druid.emitter.influxdb.protocol=https时,emitter 会使用 Apache HttpClient 构造一个带自定义SSLContext的客户端:

  1. 校验trustStorePathtrustStorePassword均已配置,否则抛IllegalStateException
  2. trustStorePath读取 KeyStore(类型为trustStoreType,默认jks),用密码加载;
  3. 基于该 TrustStore 初始化TrustManagerFactory并构造 TLS 的SSLContext
  4. 使用NoopHostnameVerifier(跳过主机名校验)构建 HttpClient,见 InfluxdbEmitter.java。

HTTPS 场景的最小配置示例:

druid.emitter.influxdb.protocol=https druid.emitter.influxdb.trustStorePath=/path/to/truststore.jks druid.emitter.influxdb.trustStoreType=jks druid.emitter.influxdb.trustStorePassword=truststore_password

若使用自签证书或内部 CA,需要把对应的 CA 证书导入到上述 trustStore 中,Druid 侧才能与 InfluxDB 完成 TLS 握手。

常见问题排查要点

  • 启动即报错:若hostnamedatabaseNameinfluxdbUserNameinfluxdbPassword缺失,配置类构造失败,进程启动失败,日志会提示对应字段不可为 null。
  • HTTPS 配置缺失protocol=https但未同时提供 trustStore 路径与密码,抛出IllegalStateException,日志提示 "Can't load TrustStore. Truststore path or password is not set."。
  • 写入失败postToInflux()中 POST 请求异常(网络不通、认证失败、数据库不存在等)会被捕获并记录 info 日志("Failed to post events to InfluxDB."),不会导致 Druid 进程崩溃,但数据会丢失(队列已被清空)。因此务必提前确认 InfluxDB 侧数据库已创建、认证用户具备写入权限、网络可达。
  • 数据未出现:确认druid.emitter=influxdb已设置、扩展已加入druid.extensions.loadList,并注意flushDelay/flushPeriod默认均为 60000ms,首次冲刷发生在启动约 1 分钟后。

小结

druid-influxdb-emitter为 Druid 提供了一条通向 InfluxDB 的轻量指标通路:事件先入有界队列,再按固定周期整体转换为 InfluxDB Line Protocol 批量 POST,关闭时兜底冲刷。本文从扩展加载、参数配置、InfluxDB 前置要求、Line Protocol 转换规则到 HTTPS 支持,完整还原了该扩展的官方文档与源码实现。核心实现可继续阅读 InfluxdbEmitter.java、InfluxdbEmitterConfig.java 与 InfluxdbEmitterTest.java 三个文件,其中测试用例对转换规则的逐字符断言是最值得信赖的行为规范。

  • 数据库
  • OLAP
  • 大数据
  • 后端

【免费下载链接】druid

Apache Druid: a high performance real-time analytics database.

项目地址:https://gitcode.com/gh_mirrors/druid6/druid
点击查看免费下载

相关推荐

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/23 10:25:04

ToClaw 实战手册:11 个技巧让 AI Agent 配置 TaoToken 更顺手

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/23 10:24:50

sensor OV9728的参数

OV9728 是 OmniVision&#xff08;豪威科技&#xff09;推出的一款 720p 高清 CMOS 图像传感器。需要特别留意的是&#xff0c;该产品目前状态为“已停产”&#xff08;End-of-Life&#xff09;。&#x1f4f7; 核心参数光学格式&#xff1a;1/6.5 英寸像素尺寸&#xff1a;1.7…

作者头像 李华
网站建设 2026/9/23 10:23:29

电竞选手国际锦标赛实战经验与策略优化

1. 比赛背景与整体回顾2026年2月2日这场赛事对我来说意义非凡——这是我转型专业选手后参加的首次国际级锦标赛。作为一项综合了策略规划、实时应变与心理博弈的竞技项目&#xff0c;这场比赛云集了32个国家的128名顶尖选手&#xff0c;赛程持续14小时&#xff0c;包含5个阶段的…

作者头像 李华
网站建设 2026/9/23 10:15:15

SWAT模型运行报错排查指南:从TxtInOut到成功运行

1. 从TxtInOut文件夹说起&#xff1a;SWAT跑起来之前的那道坎很多人以为SWAT模型最难的环节在数据准备和参数率定&#xff0c;但真正让新手卡住动弹不得的&#xff0c;往往是点击"Run SWAT"之后那几秒——要么弹出一个看不懂的报错框&#xff0c;要么运行进度条走到一…

作者头像 李华