Apache Druid 数据摄入排障实战指南:从事件丢失到 Segment 交接的完整排查手册
2026/9/23 2:57:21 网站建设 项目流程

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 问题都围绕它们展开:

  1. 实时摄入(Realtime Ingestion):事件以流(如 Kafka、HTTP push 或 pull)方式进入 Realtime 节点或索引服务任务,先经过内存缓冲(in-memory buffer)与中间持久化(intermediate persist),最终在 Segment 交接(handoff)时交给 Historical 节点加载并提供查询。
  2. 批量摄入(Batch Ingestion):通过 Index Task(本地)或 Hadoop 任务(index_hadoop)从静态文件、已有 Segment(IngestSegmentFirehosedataSourceinputSpec)中读取数据,直接产出 Segment 写入深度存储(Deep Storage)。

两条路径的公共失败模式都可以从三处获取证据:摄入进程日志摄入指标(ingest 系列 metrics)Coordinator 控制台。下表是本 FAQ 涉及问题与排查入口的速查:

问题现象首选排查入口对应章节
实时摄入无事件 / 事件被丢弃日志中ingest/events/*指标第二节
批量摄入零事件检查 ingestion spec 的intervals是否覆盖数据时间范围第二节
部分事件未摄入ingest/events/thrownAwayingest/events/unparseable指标第三节
Segment 上传到哪druid.storage.type配置第四节
流式摄入不交接元数据存储、Historical 容量、深度存储配置第五节
Historical 上无 SegmentCoordinator 控制台 + 容量配置第八节
查询返回空结果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 节点加载。如果交接失败,先确认两件事:

  1. 摄入进程日志中无异常
  2. 运行分布式集群时druid.storage.type不能是local——本地存储只适合单机/测试环境,分布式集群下 Historical 无法跨节点访问本地磁盘。

其余常见交接失败原因,官方 FAQ 明确列出以下四类:

  1. Druid 无法写入元数据存储(metadata storage):检查 MySQL / PostgreSQL 等元数据存储的配置是否正确。元数据存储承载 Segment 的版本、加载状态等信息,是交接成功的前提。
  2. Historical 节点容量不足,无法下载更多 Segment:此时 Coordinator 日志会出现异常,Coordinator 控制台也会显示 Historical 节点接近容量上限。需要调大 Historical 的容量配置(见第八节)。
  3. Segment 损坏无法下载:Historical 节点日志中会出现异常。
  4. 深度存储配置不当:确认 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 给出三条硬性要求:

  1. druid-hdfs-storage扩展包含进 classpath(按 including-extensions 的说明加载扩展);
  2. 把全部 Hadoop 配置和依赖放进 classpath——在一台已配置 Hadoop 的机器上执行以下命令即可获得依赖清单:
    hadoop classpath
  3. 按深度存储文档提供必要的 HDFS 设置(见 docs/content/dependencies/deep-storage.md 与 HDFS 扩展文档):
属性可能值说明
druid.storage.typehdfs必须设置
druid.storage.storageDirectory存放 Segment 的目录必须设置
druid.hadoop.security.kerberos.principaldruid@EXAMPLE.COMKerberos 主体(可选)
druid.hadoop.security.kerberos.keytab/etc/security/keytabs/druid.headlessUser.keytabkeytab 路径(可选)

如果使用 Hadoop indexer,把输出目录设置为 Hadoop 上的路径即可直接工作。若集群开启了 Kerberos 安全认证,可通过设置druid.hadoop.security.kerberos.principaldruid.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)

查询为空时,按以下三步排查:

  1. 用 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"],还可指定sizetimestampSpecqueryGranularityaggregatorsrollup等分析类型(完整说明)。特别地,aggregators分析会返回“可用于查询指标列”的聚合器列表——用它核对聚合器名最直接。

  2. 确认查询中使用的聚合器名称与上述指标名完全一致。名称不匹配是“有数据但查不到”的高频原因。
  3. 确认查询区间(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要读取行的数据源(类似关系型数据库的表)
intervalISO-8601 区间,定义要读取的数据时间范围
dimensions要选取的维度列表;为空数组则返回空维度;为 null 或不定义则返回全部维度
metrics要选取的指标列表;为空数组则返回空指标;为 null 或不定义则选取全部指标
filter维度过滤器(DimFilter),可在回灌时过滤掉想删除的行

这些参数与 IngestSegmentFirehoseFactory.java 的@JsonProperty定义一一对应(dataSourceintervalfilterdimensionsmetrics)。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 本质上是同一件事,操作方式完全一致:

  1. 使用IngestSegmentFirehose运行一个 Indexer Task。该 Firehose 会把已有 Segment 读出来、按新粒度聚合(aggregate),再写回 Druid;
  2. 回灌过程中可以用filter过滤掉想删除的行(例如清理有问题的数据);
  3. 通常以批量任务方式运行,例如每天喂入一块数据并聚合。

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):

参数默认值说明
maxRowsInMemory75000内存中聚合的最大行数(rollup 后行数),用于控制 JVM 堆占用
intermediatePersistPeriodPT10M中间持久化周期
windowPeriodPT10M实时摄入接受事件的时间窗口
maxPendingPersists0最多可排队的未完成持久化数

十三、结语:把 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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询