Apache Druid 数据摄入排障实战指南:从事件丢失到 Segment 交接的完整排查手册
【免费下载链接】druidApache Druid: a high performance real-time analytics database.项目地址: https://gitcode.com/gh_mirrors/druid7/druid
本指南基于 Apache Druid 仓库的官方摄入 FAQ(docs/content/ingestion/faq.md)整理而成,围绕数据摄入(ingestion)全链路中最常见的高频问题展开:实时摄入事件被拒、批量摄入零事件、事件丢失、Segment 落盘位置、流式摄入不交接、HDFS 深度存储配置、Historical 节点无 Segment、查询返回空结果,以及如何通过重新索引(reindex)完成 Schema 变更与粒度调整。读完本文,你将掌握一套可复现的排查命令、关键配置参数和底层原理(含源码级证据),能够独立定位并解决 Druid 摄入链路上的绝大多数故障。
一、Druid 摄入的两条主路径与排障总览
在进入具体问题之前,先明确 Druid 数据摄入的两条主路径,因为几乎所有 FAQ 问题都围绕它们展开:
- 实时摄入(Realtime Ingestion):事件以流(如 Kafka、HTTP push 或 pull)方式进入 Realtime 节点或索引服务任务,先经过内存缓冲(in-memory buffer)与中间持久化(intermediate persist),最终在 Segment 交接(handoff)时交给 Historical 节点加载并提供查询。
- 批量摄入(Batch Ingestion):通过 Index Task(本地)或 Hadoop 任务(
index_hadoop)从静态文件、已有 Segment(IngestSegmentFirehose或dataSourceinputSpec)中读取数据,直接产出 Segment 写入深度存储(Deep Storage)。
两条路径的公共失败模式都可以从三处获取证据:摄入进程日志、摄入指标(ingest 系列 metrics)与Coordinator 控制台。下表是本 FAQ 涉及问题与排查入口的速查:
| 问题现象 | 首选排查入口 | 对应章节 |
|---|---|---|
| 实时摄入无事件 / 事件被丢弃 | 日志中ingest/events/*指标 | 第二节 |
| 批量摄入零事件 | 检查 ingestion spec 的intervals是否覆盖数据时间范围 | 第二节 |
| 部分事件未摄入 | ingest/events/thrownAway、ingest/events/unparseable指标 | 第三节 |
| Segment 上传到哪 | druid.storage.type配置 | 第四节 |
| 流式摄入不交接 | 元数据存储、Historical 容量、深度存储配置 | 第五节 |
| Historical 上无 Segment | Coordinator 控制台 + 容量配置 | 第八节 |
| 查询返回空结果 | Segment Metadata Query + 聚合器名 + 查询区间 | 第九节 |
| 需要改 Schema / 粒度 | IngestSegmentFirehose 或 HadoopdataSourceinputSpec | 第十、十一节 |
| 实时摄入卡住 | 检查 persist / handoff 是否超时、列构建耗时 | 第十二节 |
二、数据没有被加载(My Data isn't being loaded)
2.1 实时摄入:windowPeriod之外的乱序事件被拒绝
实时摄入最常见的失败原因是事件时间落在 Druid 配置的windowPeriod之外。Druid 实时摄入只接受“距离当前时间在可配置窗口内”的事件——这一机制用于防止乱序(out-of-order)或过期事件无限期占用内存。
从源码可以确认该行为的具体实现。RealtimeTuningConfig.java 定义了默认配置:
private static final int defaultMaxRowsInMemory = 75000; private static final Period defaultIntermediatePersistPeriod = new Period("PT10M"); private static final Period defaultWindowPeriod = new Period("PT10M");即默认windowPeriod为 10 分钟,在 tuningConfig 中通过"windowPeriod": "PT10M"覆盖。而实际拒绝判定逻辑位于 MessageTimeRejectionPolicyFactory.java:
@Override public boolean accept(long timestamp) { long maxTimestamp = this.maxTimestamp; if (timestamp > maxTimestamp) { maxTimestamp = tryUpdateMaxTimestamp(timestamp); } return timestamp >= (maxTimestamp - windowMillis); }从代码可以推断:该策略以已见到的最大事件时间为基准,只接受落在maxTimestamp - windowMillis之后的事件,更早的事件一律拒绝。这正是“迟到事件被丢弃”的根源。如何确认:查看实时进程日志中带有ingest/events/*的行,这些指标会告诉你事件的 ingested(已摄入)、rejected(被拒)等情况。如果被拒事件大量存在,说明你的数据带时间戳与当前系统时间偏离过大,或需要调大windowPeriod。
生产建议(官方明确推荐):历史数据请使用批量摄入方式(batch ingestion),不要走实时摄入——实时摄入的窗口机制天然不适合回填历史。
2.2 批量摄入:intervals未覆盖数据时间范围
如果批量加载历史数据时没有任何事件被加载,首先确认 ingestion spec 中granularitySpec.intervals是否真正包含数据的时间范围。Druid 会直接丢弃区间之外的事件。一个典型正确示例(见 batch-ingestion.md):
"granularitySpec" : { "type" : "uniform", "segmentGranularity" : "DAY", "queryGranularity" : "NONE", "intervals" : [ "2013-08-31/2013-09-01" ] }intervals是 ISO-8601 区间(start/end),必须覆盖你数据中事件时间戳的实际范围,否则事件在摄入阶段就被静默丢弃。
三、并非所有事件都被摄入(Not all of my events were ingested)
3.1 用摄入指标定位被拒事件
Druid 会拒绝windowPeriod之外的事件,最可靠的判断方式是查看Druid ingest 指标(完整指标表见 docs/content/operations/metrics.md)。下表为与“事件被拒”直接相关的指标:
| 指标 | 含义 | 正常值 |
|---|---|---|
ingest/events/thrownAway | 因超出windowPeriod被拒绝的事件数 | 0 |
ingest/events/unparseable | 因无法解析被拒绝的事件数 | 0 |
ingest/events/processed | 每个上报周期内成功处理的事件数 | 等于周期内事件数 |
ingest/rows/output | 持久化的 Druid 行数(rollup 后) | 事件数(含 rollup) |
ingest/events/messageGap | 事件数据时间与当前系统时间的差距 | 取决于事件携带时间 |
这些指标的实现在 RealtimeMetricsMonitor.java 中,且仅当 Realtime 节点的 monitors 列表包含RealtimeMetricsMonitor时才会输出,配置时需注意。此外,ingest/persists/backPressure(创建 persist 任务并阻塞等待的毫秒数)正常应为 0 或极低,若持续偏高说明持久化链路存在瓶颈。
3.2 摄入数正确但查询结果不对:聚合器陷阱
如果摄入的事件数看起来正确,请确认你的查询是否构造正确。如果你在 ingestion spec 中定义了count聚合器,查询时必须用longSum聚合器去聚合这个字段。若在查询中直接使用count聚合器,统计的是 Druid 行的数量,而不是原始事件数——因为 Druid 摄入时会进行roll-up(按维度聚合压缩行数)。关于 rollup 对行数的影响,可参考 docs/content/design/segments.md 中关于 Segment 结构与 rollup 的说明。
四、摄入完成后 Segment 去了哪里(Where do my segments end up)
Segment 的去向完全取决于druid.storage.type的取值。摄入完成后,Druid 会把 Segment 上传到深度存储(Deep Storage)。默认深度存储是本地磁盘(local mount),配置项见 docs/content/dependencies/deep-storage.md:
| 属性 | 说明 | 默认 |
|---|---|---|
druid.storage.type | 深度存储类型,必须设置 | local |
druid.storage.storageDirectory | 存放 Segment 的目录 | 必须设置 |
深度存储的持久性决定了数据的安全性:只要 Druid 节点能访问到该存储层中的 Segment,无论丢失多少个 Druid 节点数据都不会丢;一旦 Segment 从该存储层消失,其代表的数据即永久丢失。生产环境建议使用 S3(druid-s3-extensions)、HDFS(druid-hdfs-storage)等分布式存储,详见 扩展列表。
五、流式摄入不交接 Segment(My stream ingest is not handing off segments)
实时摄入的最后一环是Segment 交接(handoff):Realtime 节点把 Segment 推送到深度存储,并通过元数据存储与 Zookeeper 通知 Historical 节点加载。如果交接失败,先确认两件事:
- 摄入进程日志中无异常;
- 运行分布式集群时
druid.storage.type不能是local——本地存储只适合单机/测试环境,分布式集群下 Historical 无法跨节点访问本地磁盘。
其余常见交接失败原因,官方 FAQ 明确列出以下四类:
- Druid 无法写入元数据存储(metadata storage):检查 MySQL / PostgreSQL 等元数据存储的配置是否正确。元数据存储承载 Segment 的版本、加载状态等信息,是交接成功的前提。
- Historical 节点容量不足,无法下载更多 Segment:此时 Coordinator 日志会出现异常,Coordinator 控制台也会显示 Historical 节点接近容量上限。需要调大 Historical 的容量配置(见第八节)。
- Segment 损坏无法下载:Historical 节点日志中会出现异常。
- 深度存储配置不当:确认 Segment 确实存在于深度存储中,且 Coordinator 日志无报错。
交接相关指标(同样在 docs/content/operations/metrics.md)包括:
| 指标 | 含义 | 正常值 |
|---|---|---|
ingest/handoff/failed | 交接失败次数 | 0 |
ingest/handoff/count | 已发生的交接次数 | 每个 Segment 粒度周期至少 1 次 |
ingest/sink/count | 未交接的 sink 数 | 1~3 |
ingest/persists/failed | 持久化失败次数 | 0 |
从源码看,交接由 RealtimePlumber.java 中的SegmentHandoffNotifier驱动,它负责等待 Historical 节点确认加载完成后才推进数据生命周期。
六、如何启用 HDFS 深度存储(How do I get HDFS to work)
要让 HDFS 作为深度存储工作,官方 FAQ 给出三条硬性要求:
- 把
druid-hdfs-storage扩展包含进 classpath(按 including-extensions 的说明加载扩展); - 把全部 Hadoop 配置和依赖放进 classpath——在一台已配置 Hadoop 的机器上执行以下命令即可获得依赖清单:
hadoop classpath - 按深度存储文档提供必要的 HDFS 设置(见 docs/content/dependencies/deep-storage.md 与 HDFS 扩展文档):
| 属性 | 可能值 | 说明 |
|---|---|---|
druid.storage.type | hdfs | 必须设置 |
druid.storage.storageDirectory | 存放 Segment 的目录 | 必须设置 |
druid.hadoop.security.kerberos.principal | druid@EXAMPLE.COM | Kerberos 主体(可选) |
druid.hadoop.security.kerberos.keytab | /etc/security/keytabs/druid.headlessUser.keytab | keytab 路径(可选) |
如果使用 Hadoop indexer,把输出目录设置为 Hadoop 上的路径即可直接工作。若集群开启了 Kerberos 安全认证,可通过设置druid.hadoop.security.kerberos.principal与druid.hadoop.security.kerberos.keytab主动认证(替代周期性执行kinit的 cron 方案)。该扩展还支持将 Google Cloud Storage(gs://bucket/...)作为深度存储使用。
七、Coordinator 控制台:检查 Segment 分配的第一现场
Coordinator 控制台位于http://<COORDINATOR_IP>:<PORT>(默认端口 8081),它是检查 Segment 是否被分配到 Historical 节点的第一入口。Coordinator 的核心职责是维护全局拓扑:它周期性(druid.coordinator.period,默认PT60S)对比“可用 Segment 集合”与“正在服务的 Segment 集合”,并决定 Segment 的加载/卸载/复制/均衡(相关参数见 docs/content/configuration/coordinator.md)。
Historical 节点的加载流程(详见 docs/content/design/historical.md)为:Coordinator 在 Zookeeper 中为 Historical 创建 load queue 条目 → Historical 检查本地磁盘缓存(segment cache)→ 若无缓存则从深度存储下载 Segment 元数据与数据 → 加载完成后在 Zookeeper 的 served segments 路径上宣布 → 该 Segment 立即可查询。Historical 还提供两个有用的 HTTP 端点:
GET /druid/historical/v1/loadstatus:返回本地缓存中的所有 Segment 是否都已加载,可用于判断节点重启后是否已就绪;GET /druid/historical/v1/readiness:判断节点是否可查询。
八、Historical 节点上看不到 Segment(I don't see my Druid segments on my historical nodes)
在 Coordinator 控制台确认 Segment 是否真的已加载到 Historical 节点。若 Segment 未出现,检查 Coordinator 日志中关于容量(capacity)或复制(replication)错误的消息。一个常见原因是:Historical 节点的maxSizes太小,导致无法下载更多数据。官方 FAQ 给出调整示例:
-Ddruid.segmentCache.locations=[{"path":"/tmp/druid/storageLocation","maxSize":"500000000000"}] -Ddruid.server.maxSize=500000000000对应配置的完整定义(见 docs/content/configuration/historical.md):
| 属性 | 说明 | 默认 |
|---|---|---|
druid.server.maxSize | 该节点希望被分配的 Segment 总字节数上限。注意这不是 Historical 实际强制执行的硬限制,而是发布给 Coordinator 供其规划分配的值 | 0 |
druid.segmentCache.locations | 分配给 Historical 的 Segment 先落到本地文件系统(磁盘缓存),这些位置定义本地缓存落在哪里 | 无(默认不缓存) |
同时,Historical 节点还有一组健康指标可监控(见 docs/content/operations/metrics.md):segment/max(Segment 最大字节限制)、segment/used(已用字节)、segment/usedPercent(已用百分比,应 < 100%)、segment/count(已服务 Segment 数)、segment/pendingDelete(待清理的磁盘字节数)。segment/usedPercent接近 100% 就意味着容量吃紧。
另外,官方推荐 Segment 文件大小控制在300MB–700MB之间(见 docs/content/design/segments.md),过大时应调整segmentGranularity或在partitioningSpec中调小targetPartitionSize(建议从 500 万行起步),这也与容量规划直接相关。
九、查询返回空结果(My queries are returning empty results)
查询为空时,按以下三步排查:
- 用 Segment Metadata Query 检查 datasource 实际建了哪些维度与指标。示例查询(见 docs/content/querying/segmentmetadataquery.md):
{ "queryType":"segmentMetadata", "dataSource":"sample_datasource", "intervals":["2013-01-01/2014-01-01"] }analysisTypes默认为["cardinality", "interval", "minmax"],还可指定size、timestampSpec、queryGranularity、aggregators、rollup等分析类型(完整说明)。特别地,aggregators分析会返回“可用于查询指标列”的聚合器列表——用它核对聚合器名最直接。 - 确认查询中使用的聚合器名称与上述指标名完全一致。名称不匹配是“有数据但查不到”的高频原因。
- 确认查询区间(intervals)落在存在数据的有效时间范围内。Segment 按时间分区,区间对不上自然查不到。
十、如何用新 Schema 重新索引已有数据(Reindexing with schema changes)
需要修改 Segment 的名称、维度、指标、rollup 等属性时,使用IngestSegmentFirehose + Index Task把已有 Druid Segment 按新 Schema 重新摄入。它允许你从 Druid 已有 Segment 读取数据、按新 Schema 聚合后写回 Druid。
IngestSegmentFirehose 的 spec 格式与参数(见 docs/content/ingestion/firehose.md):
{ "type" : "ingestSegment", "dataSource" : "wikipedia", "interval" : "2013-01-01/2013-01-02" }| 属性 | 说明 | 必填 |
|---|---|---|
| type | 固定为"ingestSegment" | 是 |
| dataSource | 要读取行的数据源(类似关系型数据库的表) | 是 |
| interval | ISO-8601 区间,定义要读取的数据时间范围 | 是 |
| dimensions | 要选取的维度列表;为空数组则返回空维度;为 null 或不定义则返回全部维度 | 否 |
| metrics | 要选取的指标列表;为空数组则返回空指标;为 null 或不定义则选取全部指标 | 否 |
| filter | 维度过滤器(DimFilter),可在回灌时过滤掉想删除的行 | 是 |
这些参数与 IngestSegmentFirehoseFactory.java 的@JsonProperty定义一一对应(dataSource、interval、filter、dimensions、metrics)。Index Task 通过firehoseFactory.connect(parser)建立读取管道,逐行读取已有 Segment 的数据并重新聚合(IndexTask.java)。
如果使用 Hadoop 批量摄入,则改用dataSourceinputSpec 做 reindexing(详见 docs/content/ingestion/batch-ingestion.md 与 update-existing-data.md)。reindex 是数据治理的兜底手段,官方同时提醒:建议始终保留一份原始数据副本,以备未来再次 reindex。
十一、如何改变已有数据的粒度(Changing granularity of existing data)
常见场景:希望降低老数据的粒度——例如超过 1 个月的数据只保留小时级粒度,而新数据保持分钟级粒度。这与第十节的 reindexing 本质上是同一件事,操作方式完全一致:
- 使用IngestSegmentFirehose运行一个 Indexer Task。该 Firehose 会把已有 Segment 读出来、按新粒度聚合(aggregate),再写回 Druid;
- 回灌过程中可以用
filter过滤掉想删除的行(例如清理有问题的数据); - 通常以批量任务方式运行,例如每天喂入一块数据并聚合。
Hadoop 批量摄入路径则使用dataSourceinputSpec 完成同样的重索引(详见 docs/content/ingestion/batch-ingestion.md)。
需要注意粒度调整的成本:降低粒度意味着跨 Segment 合并与重算聚合,耗时与数据量成正比,建议在低峰窗口执行,并提前评估目标 Segment 规模(300MB–700MB 区间)。
十二、实时摄入看起来卡住了(Real-time ingestion seems to be stuck)
实时摄入“卡住”在多数情况下是有意的背压(backpressure)机制在起作用。Druid 会在以下两种情况下主动限流(throttle)以防止 OOM:
- 中间持久化(intermediate persist)耗时过长;
- 交接(handoff)耗时过长。
源码证据:ingest/persists/backPressure指标衡量“创建 persist 任务并阻塞等待其完成”的毫秒数(RealtimeMetricsMonitor.java),正常情况下应为 0 或极低,偏高即代表背压正在生效。
排查建议:如果节点日志显示某些列构建耗时极长(例如 Segment 粒度是小时级,但构建某一列就花了 30 分钟),则应重新评估配置或扩容实时摄入节点。列构建慢通常意味着高基数(high cardinality)维度、过大的内存行缓冲或过小的 persist 周期,可从以下 tuningConfig 参数入手(默认值见 RealtimeTuningConfig.java):
| 参数 | 默认值 | 说明 |
|---|---|---|
maxRowsInMemory | 75000 | 内存中聚合的最大行数(rollup 后行数),用于控制 JVM 堆占用 |
intermediatePersistPeriod | PT10M | 中间持久化周期 |
windowPeriod | PT10M | 实时摄入接受事件的时间窗口 |
maxPendingPersists | 0 | 最多可排队的未完成持久化数 |
十三、结语:把 FAQ 变成你的排障 SOP
本 FAQ 覆盖的十二个问题,本质上是同一套排障循环在不同环节的投影:先看日志与指标确认“数据到底有没有进来”(ingest/events/*、ingest/persists/*、ingest/handoff/*),再看存储与容量确认“数据落到了哪”(深度存储、元数据存储、Historical 磁盘缓存),最后回到查询层核对 Schema 匹配(Segment Metadata Query、聚合器名、查询区间)。官方 FAQ 的末尾也强调,数据摄入对初次使用者确有门槛,遇到问题可以到社区交流渠道(IRC、Druid 用户 Google Group)求助;同时,掌握本文所列的日志、指标与配置检查点,绝大多数摄入问题都能在 10 分钟内定位到根因。
延伸阅读(仓库内路径):
- 摄入指标全表:docs/content/operations/metrics.md
- 批量摄入与 inputSpec:docs/content/ingestion/batch-ingestion.md
- Firehose 与 IngestSegmentFirehose:docs/content/ingestion/firehose.md
- 更新已有数据(reindex / delta):docs/content/ingestion/update-existing-data.md
- 深度存储配置:docs/content/dependencies/deep-storage.md
- Segment 结构与 rollup:docs/content/design/segments.md
- Historical 加载流程:docs/content/design/historical.md
- Segment Metadata Query:docs/content/querying/segmentmetadataquery.md
【免费下载链接】druidApache Druid: a high performance real-time analytics database.项目地址: https://gitcode.com/gh_mirrors/druid7/druid
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考