导读:本期聚焦于小何创作的《如何结合SQLite与Druid时序数据库构建一个可靠的物联网监控项目?》,敬请观看详情。传感器数据先落SQLite再批量送入Druid,是不少物联网项目控制写入成本的有效路径。SQLite在边缘端提供WAL持久化和断点续传能力,Druid在服务端承担海量时序数据的预聚合与多维分析。本文围绕设备温度监控这一实战项目,说明SQLite本地表结构、Druid摄入规范、Python桥接同步和去重机制。文章还会对比两种存储在小批量写入、范围扫描、聚合查询中的表现,并指出Druid不适合频繁更新、SQLite需要定期清理等关键边界。读完可以掌握从边缘缓存到分析集群的完整落地方法,避免数据积压、重复导入和查询延迟偏高的问题。

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

如何结合SQLite与Druid时序数据库构建一个可靠的物联网监控项目?

一、明确分工:SQLite负责可靠写入,Druid负责分析查询

在这个实战项目中,SQLite 扮演的第一个角色是边缘端或采集端的本地持久化层。设备数据往往先经过网关或小型主机,这些环境可能网络抖动、带宽受限,甚至短暂离线。SQLite 的嵌入式特性让它不需要独立服务进程,只要一个文件就能完成事务写入。开启 WAL 模式后,读操作不会阻塞写操作,写入性能也能满足每分钟数千条以内的场景。对于设备温度、湿度这类数据,建议使用 ts 字段保存 ISO 格式时间,device_id 作为维度字段,temperaturehumidity 作为指标字段。下表结构可以直接落地使用:

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 是典型的维度,temperaturehumidity 是指标。Druid 支持在摄入阶段进行 rollup,也就是按照时间粒度和维度组合预先聚合,从而大幅减少存储量。例如按分钟粒度聚合后,原本一分钟内同一设备的 60 条原始读数可以合并成一条记录,同时保留计数、求和、最大值等聚合结果。

下面是一份完整的 Druid 摄入规范,它从本地目录读取 JSON 文件,将数据写入名为 sensor_metrics 的数据源。数据源按小时切分 segment,查询粒度为一分钟,并对 temperaturehumidity 做求和与最大值聚合:

{
  "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,专门优化扫描与聚合。

下面是一个简化对比表,帮助选型时快速判断:

维度SQLiteDruid
写入方式单行事务,毫秒级批量摄入,适合大数据量
查询模式点查、小范围扫描大范围扫描、多维聚合
存储成本低,单文件易复制高,需要集群与深存储
运维复杂度极低较高,需维护多个服务节点
适合数据规模GB 以内TB 及以上

运维上,SQLite 需要定期清理已同步的历史数据,否则文件会持续增长。可以结合 shipped=1created_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

免责声明:​ 已尽一切努力确保本网站所含信息的准确性。网站内容多为原创整理与精心编撰,观点力求客观中立。本站旨在免费分享,内容仅供个人学习、研究或参考使用。若引用了第三方作品,版权归原作者所有。如内容涉及您的权益,请联系我们处理。
内容垂直聚焦
专注技术核心技术栏目,确保每篇文章深度聚焦于实用技能。从代码技巧到架构设计,为用户提供无干扰的纯技术知识沉淀,精准满足专业提升需求。
知识结构清晰
覆盖从开发到部署的全链路。AI、前端、编程、数据库、服务器、建站、系统层层递进,构建清晰学习路径,帮助用户系统化掌握开发与运维所需的核心技术。
深度技术解析
拒绝泛泛而谈,深入技术细节与实践难点。无论是数据库优化还是服务器配置,均结合真实场景与代码示例进行剖析,致力于提供可直接应用于工作的解决方案。
专业领域覆盖
精准对应开发生命周期。从前端界面到后端编程,从数据库操作到服务器运维,形成完整闭环,一站式满足全栈工程师和运维人员的技术需求。
即学即用高效
内容强调实操性,步骤清晰、代码完整。用户可根据教程直接复现和应用于自身项目,显著缩短从学习到实践的距离,快速解决开发中的具体问题。
持续更新保障
专注既定技术方向进行长期、稳定的内容输出。确保各栏目技术文章持续更新迭代,紧跟主流技术发展趋势,为用户提供经久不衰的学习价值。