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

解决APScheduler多线程中SQLite事务嵌套报错问题

问题:APScheduler定时脚本中的SQLite事务嵌套错误修复

问题描述

编写了一个基于APScheduler的定时脚本,通过SQLite存储数据,对比变更后用Telegram API发送通知。脚本每小时执行任务,但出现cannot start a transaction within a transaction错误,此前为解决跨线程游标问题将游标创建放入compareChanges()函数,但引发事务冲突,且任务可能同时执行(不会更新同一行)。

原代码:

import sqlite3
conn = sqlite3.connect('db.db', check_same_thread=False)

from apscheduler.schedulers.background import BackgroundScheduler
sched = BackgroundScheduler()

def compareChanges(site, new_value, url, subname = ""):
    cur = conn.cursor()

    # Get old value (if exist)
    old_value = cur.execute("SELECT Data FROM Data WHERE Name = ? AND Subname = ? LIMIT 1", (site, subname)).fetchone()
    
    # Insert/initialize if not exist (data will be updated later in this function)
    if not old_value:
        old_value = "NOT SET"
        cur.execute("INSERT INTO `Data` (Name, Subname, Data, Timestamp) VALUES (?, ?, ?, ?)", (site, subname, old_value, datetime.now().isoformat()))
        #conn.commit()

    new_value = str(new_value) # Convert to string because sometimes the data passed is a dict
    if old_value != new_value:
        # Log and notify of change
        logger.info(f"[{site}] has new value of {new_value}")
        pushMsg(f"{site} => {new_value}", url) # Send Telegram message of update
        # Set change and save file
        cur.execute("UPDATE `Data` SET `Data` = ?, `Timestamp` = ? WHERE Name = ?", (new_value, datetime.now().isoformat(), site))
        conn.commit()
        cur.close()
        return True # Return True if changed
    else:
        #logger.info(f"{site} - No Change: {old_value}")
        cur.close()
        return False # Return false on no change
    
def Function1():
    import requests
    response = requests.get("https://example.com/api/endpoint").json()
    compareChanges("Website Status Check", response['status'], "https://example.com")

def Bestbuy_Latest_Price():
    import requests
    response = requests.get("https://bestbuy.com/product/page")
    # This is pseudocode but it loads products on page and compares price for each product find
    for item in response:
        compareChanges("Bestbuy_Latest_Price", {product: price}, "https://bestbuy.com/product", subname = sku)


sched.add_job(Function1, 'cron', hour='*')
sched.add_job(Bestbuy_Latest_Price, 'cron', hour='*')
sched.start()

错误原因

  1. 全局连接的线程不安全问题:虽然设置了check_same_thread=False,但SQLite连接本身并非线程安全,多个任务线程同时操作同一个连接时,会导致事务状态混乱,触发嵌套事务错误。
  2. 未正确管理事务:原代码中INSERT操作后未提交,后续UPDATE又执行commit,导致事务状态残留;同时UPDATE语句的WHERE条件遗漏了Subname,会错误更新同Name下所有Subname的记录。

修复方案

关键修改点

  • 每个任务/数据库操作请求使用独立的数据库连接,避免多线程共享连接
  • 确保每个数据库操作的事务完整,使用with语句自动管理连接和游标,避免资源泄漏
  • 修复UPDATE语句的WHERE条件,包含Subname保证更新准确
  • 统一事务提交逻辑,避免未提交的操作残留

修正后的代码

import sqlite3
from apscheduler.schedulers.background import BackgroundScheduler
import requests
from datetime import datetime
import logging

# 初始化日志
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)

sched = BackgroundScheduler()

def get_db_connection():
    # 每次请求创建新连接,避免多线程冲突
    conn = sqlite3.connect('db.db')
    conn.row_factory = sqlite3.Row  # 方便获取字段值
    return conn

def compareChanges(site, new_value, url, subname=""):
    new_value = str(new_value)
    changed = False

    # 使用with语句自动管理连接和游标,自动关闭资源
    with get_db_connection() as conn:
        cur = conn.cursor()
        # 获取旧值
        cur.execute("SELECT Data FROM Data WHERE Name = ? AND Subname = ? LIMIT 1", (site, subname))
        old_value_row = cur.fetchone()
        
        old_value = old_value_row['Data'] if old_value_row else "NOT SET"

        if not old_value_row:
            # 初始化数据
            cur.execute("INSERT INTO Data (Name, Subname, Data, Timestamp) VALUES (?, ?, ?, ?)", 
                       (site, subname, old_value, datetime.now().isoformat()))
            conn.commit()

        if old_value != new_value:
            # 记录变更并发送通知
            logger.info(f"[{site}] has new value of {new_value}")
            pushMsg(f"{site} => {new_value}", url)
            # 更新数据,修正WHERE条件包含Subname
            cur.execute("UPDATE Data SET Data = ?, Timestamp = ? WHERE Name = ? AND Subname = ?", 
                       (new_value, datetime.now().isoformat(), site, subname))
            conn.commit()
            changed = True

    return changed

def Function1():
    try:
        response = requests.get("https://example.com/api/endpoint").json()
        compareChanges("Website Status Check", response['status'], "https://example.com")
    except Exception as e:
        logger.error(f"Function1 failed: {str(e)}")

def Bestbuy_Latest_Price():
    try:
        response = requests.get("https://bestbuy.com/product/page")
        # 伪代码:遍历产品
        for item in response:
            product = item.get('product')
            price = item.get('price')
            sku = item.get('sku')
            compareChanges("Bestbuy_Latest_Price", {product: price}, "https://bestbuy.com/product", subname=sku)
    except Exception as e:
        logger.error(f"Bestbuy_Latest_Price failed: {str(e)}")

# 添加定时任务
sched.add_job(Function1, 'cron', hour='*')
sched.add_job(Bestbuy_Latest_Price, 'cron', hour='*')

if __name__ == "__main__":
    sched.start()
    try:
        # 保持主进程运行
        while True:
            pass
    except KeyboardInterrupt:
        sched.shutdown()

额外说明

  1. 独立连接:get_db_connection()函数每次调用创建新连接,彻底避免多线程共享同一连接导致的事务冲突。
  2. with语句管理资源:使用with get_db_connection()自动处理连接的关闭和事务提交,无需手动close或commit(with块结束时自动提交,异常时自动回滚)。
  3. 错误处理:为每个任务添加try-except块,避免单个任务失败导致整个脚本崩溃。
  4. 修正WHERE条件:UPDATE语句加入Subname,确保只更新对应子名称的记录,避免误操作。

内容的提问来源于stack exchange,提问作者Bijan

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 00:55:12