SQLite3多进程并发访问时读取失败问题求助
SQLite3并行读写偶发读取跳过问题解决方案
问题背景
两个并行Python脚本操作同一SQLite数据库:
- Script1:每30秒从MQTT接收消息并写入
database1.db的table1 - Script2:每120秒读取
table1,偶发出现读取完全跳过的情况,单独运行Script2无此问题。已尝试设置连接timeout=20和pragma busy_timeout=10000,问题依旧。
现有代码的核心问题
- SQL注入风险+参数拼接错误:读取函数中用
%s拼接SQL语句,不仅存在注入风险,还可能因参数格式问题导致查询无结果,被误认为是锁的问题 - 连接/游标管理不规范:未确保连接在异常情况下也能关闭,可能残留锁;写入函数重复创建游标,属于冗余操作
- 隔离级别过严:SQLite默认隔离级别是
SERIALIZABLE,会加剧读写冲突 - 错误处理缺失:读取时若因锁导致查询返回空DataFrame,不会触发异常,无法区分是无数据还是锁问题
具体解决方案
- 使用上下文管理器自动处理连接的打开和关闭,避免资源泄漏和锁残留
- 用参数化查询替代字符串拼接,彻底避免SQL注入并保证参数格式正确
- 调整隔离级别为
READ COMMITTED,允许读取已提交的数据,降低读写冲突概率 - 确保
busy_timeout在连接创建后立即生效 - 针对锁繁忙的异常增加手动重试逻辑,增强容错性
- 优化写入事务,尽量缩短锁持有时间(如批量写入)
修改后的代码示例
读取函数(Script2)
import sqlite3 import pandas as pd import os from sqlite3 import Error def db_read_function(param1, param2, param3): temp_df = pd.DataFrame() # 初始化空DataFrame而非字符串'null' db_path = os.path.join(os.getcwd(), 'database1.db') retry_count = 3 # 设置重试次数 for _ in range(retry_count): try: with sqlite3.connect(db_path, timeout=20) as con: con.execute('PRAGMA busy_timeout = 10000') con.isolation_level = 'READ COMMITTED' # 参数化查询,避免SQL注入和格式问题 query = ''' SELECT * FROM table1 WHERE device_id = ? AND payload_timestamp_utc = (SELECT MAX(payload_timestamp_utc) FROM table1 WHERE device_id = ?) AND start_time_utc < ? AND end_time_utc > ? ORDER BY start_time_utc ASC ''' temp_df = pd.read_sql(query, con, params=(param1, param1, param3, param2)) break # 成功读取则跳出重试循环 except Error as e: if 'database is locked' in str(e): print(f"读取遇到锁,重试中... 错误信息: {e}") continue else: print(f"读取错误: {e}") break return temp_df
写入函数(Script1)
import sqlite3 import os from sqlite3 import Error def db_insert_function(row): last_row_id = None db_path = os.path.join(os.getcwd(), 'database1.db') try: with sqlite3.connect(db_path, timeout=20) as con: con.execute('PRAGMA busy_timeout = 10000') con.isolation_level = 'READ COMMITTED' sql = ''' INSERT INTO table1(site_name, payload_timestamp_utc, device_id, start_time_utc, end_time_utc, value) VALUES(?, ?, ?, ?, ?, ?) ''' cursor = con.cursor() cursor.execute(sql, row) con.commit() last_row_id = cursor.lastrowid except Error as e: print(f"写入错误: {e}") return last_row_id
表结构优化建议
将时间字段改为DATETIME类型或存储整数时间戳,方便排序和查询:
c.execute(''' CREATE TABLE IF NOT EXISTS table1 (id INTEGER NOT NULL PRIMARY KEY AUTOINCREMENT UNIQUE, site_name TEXT, payload_timestamp_utc DATETIME, device_id TEXT, start_time_utc DATETIME, end_time_utc DATETIME, value TEXT) ''')
额外优化点
如果Script1的写入频率较高,可考虑批量收集MQTT消息后一次性写入,减少事务次数和锁持有时间;若仍存在锁问题,可启用SQLite的WAL(Write-Ahead Logging)模式,大幅提升并发读写性能:
在连接后执行con.execute('PRAGMA journal_mode=WAL'),WAL模式允许多个读操作和一个写操作同时进行
内容的提问来源于stack exchange,提问作者Flamby
相关产品推荐
相关产品推荐

