Apache Druid 深入
概述
Apache Druid 是专为时序事件数据设计的实时分析数据库:以时间为中心分区、Segment 不可变存储、原生实时摄入。广泛用于监控、用户行为分析、实时大屏。本文讲透列式存储、时间分区、Segment、数据摄入与实时/批量架构。
一、Druid 定位
| 特性 | 说明 |
|---|---|
| 类型 | 时序 OLAP |
| 核心 | 实时摄入 + 时间分区 |
| 数据 | 事件流(带时间戳) |
| 场景 | 监控、行为分析、实时大屏 |
| 与普通 OLAP 差异 | 说明 |
|---|---|
| 时间维度 | 数据以时间组织 |
| 摄入 | 原生流式(推/拉) |
| 查询 | 时间范围 + 聚合 |
二、列式存储
2.1 存储格式
数据按列存储(列式):
只读所需列
高压缩| 列类型 | 说明 |
|---|---|
| 维度列 | 字典编码 |
| 指标列 | 数值压缩 |
| 时间列 | 分区依据 |
2.2 优势
| 优势 | 说明 |
|---|---|
| 扫描快 | 少读列 |
| 压缩高 | 字典/压缩 |
| 聚合快 | 列计算 |
三、时间分区
3.1 分区方式
数据按时间范围分区:
Segment 按时间间隔切分
查询按时间裁剪| 粒度 | 说明 |
|---|---|
| 分区粒度 | 小时/天 |
| Segment | 时间范围内一个分区 |
| 裁剪 | 跳过无关时间 |
查询 WHERE 时间范围
→ 只扫相关 Segment3.2 段内组织
Segment 内:
按时间排序
列式 + 索引
不可变(追加新段)四、Segment
4.1 概念
Segment = 数据不可变存储单元:
一段时间范围的数据(按时间粒度切分)
自包含(元数据 + 数据)
分布在 Historical 节点| 特性 | 说明 |
|---|---|
| 不可变 | 写后只读 |
| 自包含 | 独立加载 |
| 分布 | 多副本 |
| 生命周期 | 从实时到历史 |
4.2 Segment 流转
实时段 → 定时落盘 → 历史段
实时节点持有 → 发布 → Historical 节点| 阶段 | 节点 |
|---|---|
| 摄入 | MiddleManager |
| 实时查询 | 实时节点 |
| 历史存储 | Historical |
| 合并 | Coordinator |
五、数据摄入
5.1 摄入方式
| 方式 | 说明 |
|---|---|
| 推送(Push) | 直接发给 MiddleManager |
| 拉取(Pull) | 从 Kafka 等拉取 |
| 批摄入 | 文件批量 |
实时摄入:
Kafka 主题 → 实时任务 → Segment
延迟秒级可见5.2 摄入配置(以 Kafka 为例)
json
{
"type": "kafka",
"dataSchema": {
"dataSource": "events",
"timestampColumn": "ts",
"dimensions": ["city", "event_type"],
"metrics": [{"name": "count", "type": "long"}]
},
"ioConfig": {
"topic": "events",
"consumerProperties": {"bootstrap.servers": "kafka:9092"}
},
"granularitySpec": {
"segmentGranularity": "HOUR",
"queryGranularity": "MINUTE"
}
}| 配置 | 说明 |
|---|---|
| timestampColumn | 时间列 |
| dimensions | 维度 |
| metrics | 指标 |
| segmentGranularity | 段粒度 |
| queryGranularity | 聚合粒度 |
5.3 摄入优化
| 优化 | 说明 |
|---|---|
| 分区数 | 与并行匹配 |
| 批量 | 批量提交 |
| 指标预聚合 | 减少数据量 |
| 采样 | 可选 |
六、实时/批量摄入架构
6.1 实时链路
Kafka → MiddleManager(实时任务)
→ 实时 Segment(查询可见)
→ 落盘发布 → Historical6.2 批量链路
HDFS/文件 → 批摄入任务
→ Segment → Historical| 对比 | 实时 | 批量 |
|---|---|---|
| 来源 | Kafka | 文件/湖 |
| 延迟 | 秒级 | 分钟+ |
| 用途 | 实时指标 | 历史加载 |
6.3 混合
实时覆盖近期 + 批量加载历史:
统一查询(时间范围合并)七、查询
7.1 查询类型
| 类型 | 说明 |
|---|---|
| 原生 JSON 查询 | Timeseries/TopN/GroupBy |
| SQL | 类 SQL(Calcite) |
| 过滤 | 时间 + 维度 |
SQL 示例:
SELECT city, COUNT(*) FROM events
WHERE ts >= '2026-08-04 10:00:00'
GROUP BY city;7.2 查询流程
Broker 节点:
1. 接收查询
2. 分发到各节点
3. 合并结果| 节点 | 职责 |
|---|---|
| Broker | 查询入口 |
| Historical | 历史段查询 |
| 实时节点 | 实时段查询 |
八、架构组件
| 组件 | 职责 |
|---|---|
| Coordinator | 段分配/管理 |
| Overlord | 任务调度 |
| MiddleManager | 摄入执行 |
| Historical | 历史数据 |
| Broker | 查询路由 |
| ZK/Metastore | 协调元数据 |
多节点扩展:
Historical 扩存储
MiddleManager 扩摄入
Broker 扩查询并发九、适用与限制
9.1 适用
| 场景 | 说明 |
|---|---|
| 监控指标 | 时间序列 |
| 用户行为 | 事件分析 |
| 实时大屏 | 秒级聚合 |
| 日志分析 | 时间查询 |
9.2 限制
| 限制 | 说明 |
|---|---|
| 复杂 Join 弱 | 需预宽表 |
| 更新难 | Segment 不可变 |
| 运维复杂 | 组件多 |
| 明细查询弱 | 聚合为主 |
常见问题速查
| 问题 | 原因与处理 |
|---|---|
| 查询慢 | 检查时间裁剪/粒度 |
| 摄入积压 | 扩 MiddleManager |
| 段太多 | 调整段粒度/合并 |
| 实时数据丢 | 检查 Kafka 位点 |
| Join 弱 | 预宽表摄入 |