- 数据库
- OLAP
- 大数据
- 后端
【免费下载链接】druid
Apache Druid: a high performance real-time analytics database.
本篇技术指南以 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),该扩展的工作流程分为四个阶段:
- 接收事件入队:Druid 的监控框架每产生一个事件,
emit(Event event)方法被调用。该方法只处理ServiceMetricEvent类型的事件(其余事件类型被忽略),并将其放入一个有界阻塞队列LinkedBlockingQueue(容量由maxQueueSize决定),见 InfluxdbEmitter.java。 - 定时批量转换:
start()启动一个单线程的ScheduledExecutorService,按照flushDelay(首次延迟)和flushPeriod(固定周期)调度transformAndSendToInfluxdb(),见 InfluxdbEmitter.java。 - 全量冲刷发送:
transformAndSendToInfluxdb()一次性取出当前队列中的全部事件(eventsQueue.size()决定取出数量,逐个poll()),逐条转换为 Line Protocol 后拼接,再 POST 到 InfluxDB 的 HTTP 写接口,见 InfluxdbEmitter.java。整个队列在一次发送中全部清空,而不是逐条发送。 - 关闭时兜底冲刷:
close()被调用(进程关闭/优雅退出)时,会先调用flush()将队列中剩余事件全部发送,随后关闭调度线程,见 InfluxdbEmitter.java。
注意:该扩展目前只发送 Service Metric Events(服务指标事件),即 Druid metrics 文档中列举的那一类指标;其他 feed(如 alerts)不会进入该队列。
加载扩展
druid-influxdb-emitter属于社区扩展(Community Extension),不随默认 Druid 发行包(tarball)打包。要使用它,需要先把它下载安装到 Druid 的extensions目录,再在common.runtime.properties的druid.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.hostname | InfluxDB 服务器的主机名 | 是 | 无 |
druid.emitter.influxdb.port | InfluxDB 服务器端口 | 否 | 8086 |
druid.emitter.influxdb.protocol | 发送指标使用的协议,http或https二选一 | 否 | http |
druid.emitter.influxdb.trustStorePath | https 场景下使用的 trustStore 路径 | 否 | 无 |
druid.emitter.influxdb.trustStoreType | https 场景下使用的 trustStore 类型 | 否 | jks(实际取java.security.KeyStore.getDefaultType()) |
druid.emitter.influxdb.trustStorePassword | https 场景下使用的 trustStore 密码 | 否 | 无 |
druid.emitter.influxdb.databaseName | InfluxDB 中的数据库名 | 是 | 无 |
druid.emitter.influxdb.maxQueueSize | 存放待发送事件的队列容量 | 否 | Integer.MAX_VALUE(2^31-1) |
druid.emitter.influxdb.flushPeriod | 每隔多少毫秒把事件队列解析为 Line Protocol 并 POST 到 InfluxDB | 否 | 60000 |
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";几个值得注意的细节:
hostname、databaseName、influxdbUserName、influxdbPassword四个字段在构造函数中通过Preconditions.checkNotNull强制校验,缺失任何一个都会在启动阶段直接抛出异常(配置类构造失败),相关行为被 InfluxdbEmitterConfigTest.java 的testConfigWithNullHostname、testConfigWithNullInfluxdbUserName、testConfigWithNullInfluxdbPassword等用例覆盖。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时,trustStorePath与trustStorePassword二者必须同时提供,否则构造客户端时会抛出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/hits→druid_query。 - tags(标签):
service:事件中的 service 字段原样保留(如druid/historical);metric:metric 的中间段(去掉首尾段),用_拼接,并冠以druid_前缀。例如query/cache/total/hits→druid_cache_total;如果 metric 只有两段(如query/time),则没有中间段,也就不生成 metric tag;hostname:事件中的 host 字段去掉端口部分(按:切分取第一段),如historical001:8083→historical001;- 白名单维度:事件携带的用户维度(user dims)中,凡是命中
dimensionWhitelist的维度,都会追加为 tag。
- field(字段):
druid_+ metric 的最后一段,值为事件中的 value。例如query/cache/total/hits→druid_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逐段拆解:
| 组成部分 | 内容 | 来源 |
|---|---|---|
| measurement | druid_query | metric 首段query |
| tag | service=druid/historical | 事件 service |
| tag | hostname=historical001 | 事件 host 去掉端口 |
| tag | metric=druid_cache_total | metric 中间段cache/total以_拼接并加druid_前缀 |
| field | druid_hits=34787256 | metric 末段hits+ 事件 value |
| timestamp | 1509440946857000000 | 事件时间戳的纳秒表示 |
注意:每条转换后的记录末尾会追加一个换行符
\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; - 白名单默认包含:
dataSource、type、numMetrics、numDimensions、threshold、dimension、taskType、taskStatus、tier; - 可通过配置
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的客户端:
- 校验
trustStorePath与trustStorePassword均已配置,否则抛IllegalStateException; - 从
trustStorePath读取 KeyStore(类型为trustStoreType,默认jks),用密码加载; - 基于该 TrustStore 初始化
TrustManagerFactory并构造 TLS 的SSLContext; - 使用
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 握手。
常见问题排查要点
- 启动即报错:若
hostname、databaseName、influxdbUserName、influxdbPassword缺失,配置类构造失败,进程启动失败,日志会提示对应字段不可为 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.
相关推荐
Apache Druid StatsD Emitter 扩展实战指南:将 Druid 指标实时推送至 StatsD / Statsite
Apache Druid StatsD Emitter 扩展实战指南:将 Druid 指标实时推送至 StatsD / Statsite 本文以当前仓库中 st
数据库数据分析OLAP大数据实时分析数据仓库后端Apache Druid Ambari Metrics Emitter 扩展:将 Druid 指标接入 Ambari Metrics 监控体系
Apache Druid Ambari Metrics Emitter 扩展:将 Druid 指标接入 Ambari Metrics 监控体系 导读 本文围绕
数据库数据分析OLAP大数据实时分析数据仓库后端Apache Druid Graphite Emitter 指南:将 Druid 指标通过 Pickle 协议送入 Graphite Carbon
Apache Druid Graphite Emitter 指南:将 Druid 指标通过 Pickle 协议送入 Graphite Carbon 本文以 Ap
数据库数据分析OLAP大数据实时分析数据仓库后端
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考