前段时间我写了个小工具,需要处理大批量的文件转码任务。起初直接丢进线程池,但程序一重启,没跑完的任务全丢了,我当时就想:要是有个简单的持久化队列就好了。上Redis吧,感觉杀鸡用牛刀;用RabbitMQ吧,还得装Erlang。后来一拍脑门,SQLite不是现成的吗?轻量、单文件、还支持事务,于是我用Python标准库sqlite3折腾了一个任务队列,效果出乎意料地好。
这里不是说要替代真正的消息队列,但在某些不追求超高吞吐、又需要落盘的场景下,这个方案能省下不少事。下面直接看代码。
表结构设计
队列的基本需求就俩:往里放任务,从里面取任务。我建了一张很简单的表:
`tasks` (
`id` INTEGER PRIMARY KEY AUTOINCREMENT,
`payload` TEXT NOT NULL,
`status` TEXT NOT NULL DEFAULT 'pending',
`created_at` DATETIME DEFAULT CURRENT_TIMESTAMP
);
payload 用来存放任务内容,我一般放JSON字符串,具体是啥由调用方自己解释。status 有三个值:pending 表示等待执行,processing 表示被某个消费者领取了,done 表示已完成。严格来说,done 可以直接删掉,但留着方便排查。
一个朴素的队列类
import sqlite3
import time
import json
import threading
class Queue:
def __init__(self, db_path, max_connections=5):
# WAL模式允许读写并发,读和写不互相阻塞
self.db_path = db_path
self._init_db()
# 用一个连接池,让每个线程有自己的连接
self._local = threading.local()
self._max_connections = max_connections
self._lock = threading.Lock()
self._connections = []
def _init_db(self):
conn = sqlite3.connect(self.db_path, timeout=10)
conn.execute("PRAGMA journal_mode=WAL")
conn.execute("""
CREATE TABLE IF NOT EXISTS tasks (
id INTEGER PRIMARY KEY AUTOINCREMENT,
payload TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'pending',
created_at DATETIME DEFAULT CURRENT_TIMESTAMP
)
""")
conn.commit()
conn.close()
def _get_conn(self):
# 每个线程一个连接,避免跨线程使用
if hasattr(self._local, 'conn'):
return self._local.conn
with self._lock:
conn = sqlite3.connect(self.db_path, timeout=10)
conn.execute("PRAGMA busy_timeout=5000")
# 这里不能用row_factory?无所谓
self._local.conn = conn
self._connections.append(conn)
return conn
def push(self, payload):
conn = self._get_conn()
conn.execute(
"INSERT INTO tasks (payload, status) VALUES (?, 'pending')",
(json.dumps(payload),)
)
conn.commit()
def pull(self):
"""原子地取走一个pending任务,并标记为processing"""
conn = self._get_conn()
# 使用 IMMEDIATE 事务,防止两个线程同时取到同一行
try:
conn.execute("BEGIN IMMEDIATE")
cur = conn.execute(
"SELECT id, payload FROM tasks "
"WHERE status='pending' ORDER BY id LIMIT 1"
)
row = cur.fetchone()
if row is None:
conn.execute("ROLLBACK")
return None
task_id, payload = row
conn.execute(
"UPDATE tasks SET status='processing' WHERE id=?",
(task_id,)
)
conn.execute("COMMIT")
return task_id, json.loads(payload)
except sqlite3.OperationalError as e:
# 数据库被锁,先等一会再重试
if conn.in_transaction:
conn.execute("ROLLBACK")
time.sleep(0.02)
return self.pull()
def finish(self, task_id):
"""标记任务完成,直接删除这条记录"""
conn = self._get_conn()
conn.execute("DELETE FROM tasks WHERE id=?", (task_id,))
conn.commit()
def size(self):
conn = self._get_conn()
cur = conn.execute("SELECT COUNT(*) FROM tasks WHERE status='pending'")
return cur.fetchone()[0]
核心逻辑在 pull() 方法。用 BEGIN IMMEDIATE 拿到写锁,然后查询并更新状态。这个操作在SQLite里是串行的,所以不会出现两个消费者拿到同一个任务的情况。如果数据库正忙,就等一下然后递归重试,反正有超时兜底。
写个消费者跑起来
测试一下:我放100个任务,再用4个线程去消费,每个任务假装忙0.01秒。
def worker(q, worker_id):
while True:
got = q.pull()
if got is None:
break
task_id, payload = got
print(f"[{worker_id}] 处理任务 {task_id}: {payload}")
time.sleep(0.01)
q.finish(task_id)
# 处理完,继续取下一个
q = Queue("tasks.db")
# 先塞100个任务
for i in range(100):
q.push({"index": i, "sleep": 0.01})
# 开4个线程
threads = []
for wid in range(4):
t = threading.Thread(target=worker, args=(q, wid))
t.start()
threads.append(t)
for t in threads:
t.join()
print("剩余待处理任务:", q.size())
运行完你会发现,100个任务被4个线程基本平分了,而且没有重复消费。如果把程序中断,再重新启动,那些状态还是pending的任务会继续被处理,因为已经处理完的都被删掉了,没处理完的还在表里。
并发安全问题在哪?
最主要的坑是:Python的sqlite3连接默认不能跨线程使用,所以我用了threading.local来让每个线程持有独立的连接。另外,BEGIN IMMEDIATE会等待写锁,如果多个线程同时跑,可能会阻塞一小会儿,但我设置了busy_timeout为5秒,足够应付这种小场景。
还有一个细节:每次pull()都可能触发递归重试。在并发高的时候,偶尔会遇到database is locked,但SQLite官方建议的重试策略是随机等待并重试。我这里等20毫秒,效果还行。如果你的任务量更大,可以把等待时间调小,或者用指数退避。
但要注意,这个队列的吞吐量肯定比不了Redis的RPOPLPUSH。我做了一个快速测试:SQLite队列每秒大约能拉取和提交1000次任务(取决于硬盘)。如果你要一秒几万次,那还是老实上Redis吧。
进阶用法:延迟任务和定期清理
如果你想让任务延迟到未来某个时刻再执行,可以加一个run_at字段,pull()的查询条件加上run_at <= now。另外,表里面任务做完就删,所以表不会无限膨胀。但如果你的业务需要保留历史记录,那就不要删除,改成更新状态为done,再建个索引。
实际上我后来就是这么干的:为了能追溯执行情况,我把finish改成了UPDATE status='done' WHERE id=? AND status='processing',然后定时把三天前的done记录删掉。功能和灵活性都够。
意外收获:崩溃恢复
因为任务在执行前就被标记为processing,如果进程在执行到一半时崩溃,这个任务就永远停在processing了。解决思路也简单:启动消费者之前,先将所有processing重置为pending。因为只有当前程序在跑,一旦启动,说明之前残留的处理中任务都要重新来一遍。这个操作只需一行SQL:
UPDATE tasks SET status='pending' WHERE status='processing'
我第一次跑这个队列的时候,半夜断电了,第二天起来发现任务一个都没丢,全被重置了,那一瞬间真的感到SQLite带来的安全感。
什么时候不适合用这个方案
如果任务量特别大,或者需要多个进程同时写一个队列,SQLite也会成为瓶颈。因为数据库文件只有一个,写锁全局串行,多进程场景下性能会急剧下降。这种情况下,还是去装个Redis吧。
但在单机、多线程、中小任务量的场景下,用它代替“内存队列”可以换来持久化;代替“写文件+自己加锁”可以换来SQL查询和可靠性。我后来甚至把这个队列丢进了一个Flask服务里,接口push任务,后台线程消费,一切都很和谐。
最后的碎碎念
有人可能会觉得SQLite写队列有点“野路子”。实际上很多开源项目就用SQLite当队列,比如SQL底层的任务中间件也有这么干的。它好在简单,没有网络开销,一个文件就能带走。
这次实验也让我对SQLite的事务模型有了更深的认识。它的锁机制虽然简陋,但通过合理的BEGIN IMMEDIATE和busy_timeout,完全能支撑一个小型任务的调度。如果你是Python开发者,手边又正好有任务队列的需求,不妨试试这个土办法。

