我拿sqlite3写了一个线程安全的持久化任务队列,还挺好用

2026-08-16 0 795

前段时间我写了个小工具,需要处理大批量的文件转码任务。起初直接丢进线程池,但程序一重启,没跑完的任务全丢了,我当时就想:要是有个简单的持久化队列就好了。上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 IMMEDIATEbusy_timeout,完全能支撑一个小型任务的调度。如果你是Python开发者,手边又正好有任务队列的需求,不妨试试这个土办法。

我拿sqlite3写了一个线程安全的持久化任务队列,还挺好用
收藏 (0) 打赏

感谢您的支持,我会继续努力的!

打开微信/支付宝扫一扫,即可进行扫码打赏哦,分享从这里开始,精彩与您同在
点赞 (0)

版权声明:
本站资源有的来自互联网收集整理,本站纯免费分享提供学习使用,如果侵犯了您的合法权益,请联系本站我们会及时删除。
本站资源仅供研究、学习交流之用,免费开源项目不代表完全可商用,若商业用途请先咨询开发企业能否商用,否则产生的一切后果将由下载用户自行承担。
原创板块未经允许不得转载,否则将追究法律责任。

淘吗网 python 我拿sqlite3写了一个线程安全的持久化任务队列,还挺好用 https://www.taomawang.com/server/python/2550.html

常见问题

相关文章

猜你喜欢
发表评论
暂无评论
官方客服团队

为您解决烦忧 - 24小时在线 专业服务