做小项目时SQLite几乎是默认选择:零部署、单文件、性能够用。但一旦老板要看实时大屏、要看每小时的漏斗转化、要按用户维度切片分析,问题就来了——SQLite是行式存储,聚合查询扫全表,数据量到千万级之后一条GROUP BY就能卡住整个进程。Apache Pinot则是专门为OLAP场景设计的列式存储引擎,毫秒级响应海量聚合查询。这篇文章就带你把这两者串起来:SQLite负责业务写入,Pinot负责实时分析,两者通过一条同步链路衔接,构成一套低成本的准实时分析架构。

为什么SQLite和Pinot是天然互补的组合
先明确两者的定位差异。SQLite是嵌入式关系型数据库,数据按行存储在单个文件里,写路径极轻,读路径适合点查和小范围查询,但对大范围聚合非常不友好。比如一张五千万行的订单表,执行SELECT city, SUM(amount) FROM orders GROUP BY city,SQLite需要完整扫描数据文件,几十秒都不奇怪。
Pinot正好相反。它是LinkedIn开源的分布式OLAP引擎,数据按列存储并且预先建立了各类索引(倒排索引、星型树索引等),同样的聚合查询在Pinot上通常几百毫秒内返回。Pinot的写入支持流式(对接Kafka)和批式(对接Hadoop、Spark或本地文件)两种模式,这为我们从SQLite同步数据提供了灵活的入口。
这套组合的价值在于成本。如果直接把业务迁到MySQL加ClickHouse,开发和运维成本都会翻倍;而SQLite加上单节点Pinot(Pinot也支持Quickstart模式单进程跑起来),一套VPS就能承载中小型项目的实时分析需求,等规模上来了再平滑扩容Pinot集群,业务侧的SQLite几乎不用动。
实战第一步:设计SQLite表与Pinot表结构
假设我们要分析的是订单数据。SQLite侧的表结构如下:
-- 业务库,负责写入
CREATE TABLE orders (
id INTEGER PRIMARY KEY AUTOINCREMENT,
user_id INTEGER NOT NULL,
city TEXT NOT NULL,
amount REAL NOT NULL,
status INTEGER NOT NULL,
created_at TEXT NOT NULL DEFAULT (datetime('now', 'localtime'))
);
CREATE INDEX idx_orders_created ON orders(created_at);Pinot侧要建一张对应的offline或realtime表。注意Pinot的schema是强类型的,并且建议把时间列显式声明出来,这样才能利用Pinot的时间分区裁剪能力:
{
"schemaName": "orders",
"dimensionFieldSpecs": [
{"name": "user_id", "dataType": "LONG"},
{"name": "city", "dataType": "STRING"},
{"name": "status", "dataType": "INT"}
],
"metricFieldSpecs": [
{"name": "amount", "dataType": "DOUBLE"}
],
"dateTimeFieldSpecs": [
{"name": "created_at", "dataType": "TIMESTAMP",
"format": "1:MILLISECONDS:EPOCH", "granularity": "1:MILLISECONDS"}
]
}一个容易被忽略的细节:SQLite的datetime('now')返回的是字符串,同步脚本里必须把它转成毫秒时间戳再写入Pinot,否则查询时时间过滤会全部失效。我在第一次搭的时候就在这里踩过坑,大屏上的时间范围筛选始终查不到数据,排查了半天才发现是时间格式问题。
实战第二步:搭建同步链路
同步方案有两条路:一是监听SQLite的变更再推到Kafka,走Pinot的realtime表;二是定时批量导出增量数据,走Pinot的batch摄入。对于中小项目,批量方案更简单可靠,我们用一个Python脚本演示核心逻辑:
import sqlite3
import time
import csv
DB_PATH = "app.db"
CHECKPOINT_FILE = "sync.checkpoint"
def load_checkpoint():
try:
with open(CHECKPOINT_FILE) as f:
return int(f.read().strip())
except FileNotFoundError:
return 0
def save_checkpoint(last_id):
with open(CHECKPOINT_FILE, "w") as f:
f.write(str(last_id))
def export_increment():
conn = sqlite3.connect(DB_PATH)
last_id = load_checkpoint()
rows = conn.execute(
"SELECT id, user_id, city, amount, status, "
"CAST(strftime('%s', created_at) AS INTEGER) * 1000 "
"FROM orders WHERE id > ? ORDER BY id LIMIT 50000",
(last_id,)
).fetchall()
if not rows:
return
# 写成CSV供Pinot batch摄入,或直接调Pinot的segment上传接口
with open("increment.csv", "w", newline="") as f:
writer = csv.writer(f)
writer.writerow(["user_id", "city", "amount", "status", "created_at"])
for r in rows:
writer.writerow(r[1:])
save_checkpoint(rows[-1][0])
conn.close()
while True:
export_increment()
time.sleep(30)脚本的关键点是检查点机制:用自增ID作为水位线,每次只导出新增数据,重启后能从断点继续,不会重复也不会遗漏。导出的CSV可以直接通过Pinot Controller的上传接口生成segment,也可以先落地成Avro文件走更正式的batch流水线。定时周期设为30秒,对绝大多数运营大屏来说,这个延迟完全够用。
如果对实时性要求更高,比如要秒级延迟,可以把time.sleep(30)换成对SQLite WAL日志的监听,或者干脆在业务写入代码里加一个Kafka生产者,双写SQLite和Kafka,Pinot的realtime表直接消费Kafka。这条路架构上更漂亮,但复杂度也上去了,建议先跑通批量链路,有需要再演进。
验证查询:同样的SQL在两边差多少
链路通了之后,用一组对比查询验证效果。SQLite侧:
-- SQLite:约42秒(1000万行测试数据) SELECT city, COUNT(*) AS cnt, SUM(amount) AS total FROM orders WHERE created_at >= '2024-01-01' GROUP BY city ORDER BY total DESC LIMIT 20;
Pinot侧的等价查询(通过PQL或标准SQL,Pinot新版本已支持SQL查询接口):
-- Pinot:约180毫秒(同样1000万行)
SELECT city, COUNT(*) AS cnt, SUM(amount) AS total
FROM orders
WHERE created_at >= FROMDATETIME('2024-01-01', 'yyyy-MM-dd')
GROUP BY city
ORDER BY total DESC
LIMIT 20;两百多倍的差距就是列式存储加索引的价值。更实际的做法是把Pinot对接Superset或Grafana做可视化,图表直连Pinot数据源,业务库完全不被分析查询打扰,线上接口的P99延迟也不会再因为运营临时跑报表而抖动。
常见坑点与适用边界
这套方案有几个坑值得提前知道。第一,Pinot不支持UPDATE和DELETE,如果业务数据有状态流转(比如订单从待支付变成已退款),要么同步时用最新状态覆盖(采用upsert语义,Pinot的realtime表支持upsert配置),要么在查询层用视图过滤。第二,SQLite在高并发写入时记得开启WAL模式,否则批量导出和业务写入会互相锁表。第三,Pinot的内存占用不低,单节点建议至少4GB内存,别和业务进程挤在一台小机器上。
适用边界也要说清楚:数据量在千万级以内、团队规模小、不需要复杂事务的分析场景,这套组合性价比极高;如果数据到了几十亿行,或者需要多表JOIN分析,就应该认真考虑完整的Pinot集群加Kafka流式链路,甚至引入数据湖方案。技术选型从来不是越重越好,SQLite加Pinot证明了轻量工具组合在一起,同样能解决看起来很重的实时分析问题。
SQLiteApache Pinot实时分析修改时间:2026-09-08 13:31:01