解决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()
错误原因
- 全局连接的线程不安全问题:虽然设置了
check_same_thread=False,但SQLite连接本身并非线程安全,多个任务线程同时操作同一个连接时,会导致事务状态混乱,触发嵌套事务错误。 - 未正确管理事务:原代码中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()
额外说明
- 独立连接:
get_db_connection()函数每次调用创建新连接,彻底避免多线程共享同一连接导致的事务冲突。 - with语句管理资源:使用
with get_db_connection()自动处理连接的关闭和事务提交,无需手动close或commit(with块结束时自动提交,异常时自动回滚)。 - 错误处理:为每个任务添加try-except块,避免单个任务失败导致整个脚本崩溃。
- 修正WHERE条件:UPDATE语句加入
Subname,确保只更新对应子名称的记录,避免误操作。
内容的提问来源于stack exchange,提问作者Bijan
相关产品推荐
相关产品推荐

