You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

SQLite3多进程并发访问时读取失败问题求助

SQLite3并行读写偶发读取跳过问题解决方案

问题背景

两个并行Python脚本操作同一SQLite数据库:

  • Script1:每30秒从MQTT接收消息并写入database1.db的table1
  • Script2:每120秒读取table1,偶发出现读取完全跳过的情况,单独运行Script2无此问题。已尝试设置连接timeout=20和pragma busy_timeout=10000,问题依旧。

现有代码的核心问题

  1. SQL注入风险+参数拼接错误:读取函数中用%s拼接SQL语句,不仅存在注入风险,还可能因参数格式问题导致查询无结果,被误认为是锁的问题
  2. 连接/游标管理不规范:未确保连接在异常情况下也能关闭,可能残留锁;写入函数重复创建游标,属于冗余操作
  3. 隔离级别过严:SQLite默认隔离级别是SERIALIZABLE,会加剧读写冲突
  4. 错误处理缺失:读取时若因锁导致查询返回空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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.03 11:01:38