SQLite 和 Apache Druid 经常被放在一起讨论,不是因为它们能互相替代,而是因为它们在时序数据链路中恰好形成互补:SQLite 解决数据可靠落盘和本地缓存,Druid 解决海量数据的预聚合与多维分析。以一个设备温度监控项目为例,边缘设备每分钟产生数百条读数,如果直接通过 HTTP 推送 Druid,小批次写入会迅速拖垮实时索引;如果只存 SQLite,随着设备数量增加,跨天范围查询和聚合会变得非常缓慢。因此更务实的架构是本地 SQLite 缓存,周期性批量导入 Druid,再用 Druid 提供分析接口。

一、明确分工:SQLite负责可靠写入,Druid负责分析查询
在这个实战项目中,SQLite 扮演的第一个角色是边缘端或采集端的本地持久化层。设备数据往往先经过网关或小型主机,这些环境可能网络抖动、带宽受限,甚至短暂离线。SQLite 的嵌入式特性让它不需要独立服务进程,只要一个文件就能完成事务写入。开启 WAL 模式后,读操作不会阻塞写操作,写入性能也能满足每分钟数千条以内的场景。对于设备温度、湿度这类数据,建议使用 ts 字段保存 ISO 格式时间,device_id 作为维度字段,temperature 和 humidity 作为指标字段。下表结构可以直接落地使用:
CREATE TABLE IF NOT EXISTS readings (
id INTEGER PRIMARY KEY AUTOINCREMENT,
device_id TEXT NOT NULL,
ts TEXT NOT NULL,
temperature REAL,
humidity REAL,
shipped INTEGER NOT NULL DEFAULT 0,
created_at TEXT NOT NULL DEFAULT (datetime('now'))
);
CREATE INDEX idx_readings_device_ts ON readings(device_id, ts);
CREATE INDEX idx_readings_shipped ON readings(shipped, id);
PRAGMA journal_mode=WAL;
PRAGMA synchronous=NORMAL;
其中 shipped 字段是同步状态标记,0 表示尚未推送到 Druid,1 表示已经推送成功。这个设计让桥接程序可以随时中断并恢复,不会因为进程重启而丢失数据。另一个容易忽略的问题是时间格式,Druid 摄入时建议直接使用 ISO 8601 字符串,避免在边缘端做复杂的时间转换。SQLite 自身没有原生的时间类型,但文本存储 ISO 时间既方便排序,也方便 Druid 解析。
Druid 则完全不适合作为单行更新的目标存储。它的数据按 segment 组织,每个 segment 通常按小时或天划分,写入后通过索引文件提供快速聚合。频繁的小批次写入会产生大量小 segment,导致查询性能下降和集群存储膨胀。因此,不要让设备直接连接 Druid 写入原始数据,而是先经过 SQLite 合并、去重,再批量提交。这样做还有一个额外好处:当 Druid 集群升级或短暂不可用时,本地 SQLite 仍然能继续接收数据,故障恢复后再补偿同步。
二、Druid摄入规范与预聚合设计
Druid 的数据摄入需要定义数据源结构、时间列、维度列和指标列。在温度监控项目中,device_id 是典型的维度,temperature 和 humidity 是指标。Druid 支持在摄入阶段进行 rollup,也就是按照时间粒度和维度组合预先聚合,从而大幅减少存储量。例如按分钟粒度聚合后,原本一分钟内同一设备的 60 条原始读数可以合并成一条记录,同时保留计数、求和、最大值等聚合结果。
下面是一份完整的 Druid 摄入规范,它从本地目录读取 JSON 文件,将数据写入名为 sensor_metrics 的数据源。数据源按小时切分 segment,查询粒度为一分钟,并对 temperature 和 humidity 做求和与最大值聚合:
{
"type": "index_parallel",
"spec": {
"dataSchema": {
"dataSource": "sensor_metrics",
"timestampSpec": {
"column": "ts",
"format": "iso"
},
"dimensionsSpec": {
"dimensions": ["device_id", "factory", "line"]
},
"metricsSpec": [
{"type": "count", "name": "count"},
{"type": "doubleSum", "fieldName": "temperature", "name": "temperature_sum"},
{"type": "doubleMax", "fieldName": "temperature", "name": "temperature_max"},
{"type": "doubleSum", "fieldName": "humidity", "name": "humidity_sum"}
],
"granularitySpec": {
"segmentGranularity": "hour",
"queryGranularity": "minute",
"rollup": true
}
},
"ioConfig": {
"type": "index_parallel",
"inputSource": {
"type": "local",
"baseDir": "/data/ingest",
"filter": "*.json"
},
"inputFormat": {
"type": "json"
}
},
"tuningConfig": {
"type": "index_parallel",
"maxRowsPerSegment": 5000000
}
}
}
这份规范里,rollup 设置为 true 表示启用预聚合,queryGranularity 控制查询时的最小时间粒度。需要注意的是,一旦启用 rollup,原始明细数据会丢失,只能查询到聚合结果。如果业务需要保留原始读数,可以把 rollup 改成 false,并让 queryGranularity 设置为 none,但存储成本会明显增加。对于温度监控这种场景,分钟级聚合通常已经足够,保留计数和最大值可以还原大部分异常判断需求。
摄入完成后,Druid 的查询 SQL 可以直接做多维过滤和时间范围聚合。例如查询某个工厂过去 24 小时各设备的平均温度,Druid 会在秒级返回结果。相比之下,同样的查询在 SQLite 上需要扫描大量行并实时计算平均值,数据量超过几百万行后延迟会显著上升。这也解释了为什么分析层交给 Druid 而不是让应用直接查询 SQLite。
三、桥接程序:批量同步、断点续传与去重
桥接程序的核心任务是把 SQLite 中 shipped=0 的记录批量读出来,转换格式后推送到 Druid 或 Kafka,再把这些记录标记为已同步。使用批量游标读取可以控制内存占用,每次读取 2000 条,推送成功后再更新状态。这样即使程序中途退出,下一轮启动时仍未标记的记录会被重新读取,实现断点续传。
下面是一个 Python 版本的桥接程序骨架,它模拟了从 SQLite 读取未同步数据并调用 Druid 索引接口的过程。实际生产环境中,建议把 Druid HTTP 接口替换为 Kafka 生产者,由 Kafka 缓冲后再由 Druid 实时或批量消费,这样能进一步降低背压风险:
import sqlite3
import time
import requests
DB_PATH = "/var/lib/sensor/sensor.db"
DRUID_INDEX_URL = "http://druid-router:8888/druid/indexer/v1/task"
def get_unshipped(cursor, limit=2000):
cursor.execute(
"SELECT id, device_id, ts, temperature, humidity FROM readings WHERE shipped=0 ORDER BY id LIMIT ?",
(limit,)
)
return cursor.fetchall()
def mark_shipped(conn, ids):
if not ids:
return
placeholders = ",".join("?" * len(ids))
conn.execute(
"UPDATE readings SET shipped=1 WHERE id IN (%s)" % placeholders,
ids
)
conn.commit()
def main():
conn = sqlite3.connect(DB_PATH)
cursor = conn.cursor()
while True:
rows = get_unshipped(cursor)
if not rows:
time.sleep(10)
continue
payload = []
for row in rows:
payload.append({
"ts": row[2],
"device_id": row[1],
"temperature": row[3],
"humidity": row[4]
})
# 这里应调用 Druid 索引接口或发送到 Kafka,示例省略真实请求
# requests.post(DRUID_INDEX_URL, json=payload)
mark_shipped(conn, [row[0] for row in rows])
time.sleep(1)
if __name__ == "__main__":
main()
这个程序在 mark_shipped 之前可以加入校验逻辑,例如通过 Druid 查询接口确认新 segment 已经可查询,或者将数据写入 Kafka 后等待成功确认。如果使用 Kafka,去重机制可以结合 Kafka 的幂等生产者或消费者端去重表。还有一种常见做法是在 Druid 摄入时指定固定的 task id,并对同一批数据使用相同的 segment 标识,避免重复摄入。
需要注意的是,SQLite 的 UPDATE readings SET shipped=1 操作如果单次更新几十万条会锁库较久。分批更新并按 LIMIT 读取可以避免长事务。对于同步延迟要求不高的场景,每 10 秒同步一批、每批 1000 到 5000 条是比较稳妥的参数。若边缘设备数量特别多,还可以使用多个 SQLite 文件分片,或者将历史已同步数据定期归档到另一张表,主表只保留最近 7 天未同步和待补偿记录。
四、性能对比与运维注意事项
从写入延迟来看,SQLite 单条插入在普通磁盘上可以做到毫秒级,而 Druid 的小批次实时写入往往需要几秒到几十秒才能被查询到。从聚合查询来看,Druid 对千万级甚至亿级行的范围聚合可以在几百毫秒内完成,SQLite 在百万级数据时已经需要数秒。这种差异来自两者完全不同的存储结构:SQLite 基于 B-tree 行存储,适合事务和点查;Druid 使用列式存储和预聚合 segment,专门优化扫描与聚合。
下面是一个简化对比表,帮助选型时快速判断:
| 维度 | SQLite | Druid |
|---|---|---|
| 写入方式 | 单行事务,毫秒级 | 批量摄入,适合大数据量 |
| 查询模式 | 点查、小范围扫描 | 大范围扫描、多维聚合 |
| 存储成本 | 低,单文件易复制 | 高,需要集群与深存储 |
| 运维复杂度 | 极低 | 较高,需维护多个服务节点 |
| 适合数据规模 | GB 以内 | TB 及以上 |
运维上,SQLite 需要定期清理已同步的历史数据,否则文件会持续增长。可以结合 shipped=1 和 created_at 条件将超过 7 天的数据删除或归档,并执行 VACUUM 回收空间。Druid 则需要关注 segment 数量和大小,必要时开启定时 compaction 来自动合并小 segment。数据摄入失败时,及时查看 Overlord 和 MiddleManager 日志,通常可以从 task 报错中定位是时间格式问题还是维度字段缺失问题。
最后需要强调,SQLite 与 Druid 的组合并不是要替代 Kafka、Flink 等流处理组件,而是给出一个轻量化的落地路径。当设备规模从几百增长到几万时,可以在 SQLite 与 Druid 之间加入 Kafka 作为缓冲层,桥接程序只负责把本地数据推入 Kafka,Druid 通过 Kafka 索引服务持续消费。这样既保留了边缘端的可靠落盘,又让分析集群能够平滑扩展,适合大多数中小型物联网项目的演进节奏。
SQLiteDruid时序数据库时序数据实战修改时间:2026-08-25 21:17:59