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

Threading/Queue应用中consumer函数finally块无法执行的解决问询

问题分析与解决方案

嘿,这个问题我之前也碰到过!核心原因是你把worker线程设置成了守护线程(Daemon Thread)——当主线程执行完urls.join()后会直接退出,守护线程会被强制终止,根本没机会走到finally块里执行数据库关闭、旧数据清理这些收尾逻辑。

下面给你分步骤解决:

1. 取消守护线程设置,等待所有worker线程自然结束

守护线程的特性是“主线程死我就死”,我们需要把worker改成非守护线程,让它们自己跑完所有逻辑(包括finally块),再让主线程等待它们全部完成。修改主线程代码:

# 主线程修改
workers = []
num_threads = 5  # 替换成你实际的线程数
for i in range(num_threads):
    worker = Thread(target=import_mo, args=(urls,))
    # 删掉 worker.setDaemon(True) 这一行,使用默认的非守护线程
    worker.start()
    workers.append(worker)

create_url()
urls.join()  # 等待队列中所有任务被标记为完成

# 等待所有worker线程执行完毕(包括finally块的清理逻辑)
for worker in workers:
    worker.join()

2. 给worker线程添加退出信号,避免无限阻塞

原来的while True循环里,urls.get()会一直阻塞等待新任务,即使队列已经空了。我们需要给每个worker线程发送一个“终止符”,让它们知道可以退出循环了:

第一步:修改create_url(),添加终止信号

在所有URL放入队列后,给每个worker线程放一个None作为终止标记:

def create_url():
    try:
        mariadb_connection = mariadb.connect(host='host', database='db', user='user', password='pw')
        cursor = mariadb_connection.cursor()
        cursor.execute('SELECT type_id from tbl_items')
        item_list = cursor.fetchall()
        print("Create URL - Record retrieved successfully")
        for row in item_list:
            url = 'https://someinternet.com/type_id=' + str(row[0])
            urls.put(url)
        
        # 给每个worker线程添加终止信号
        for _ in range(num_threads):
            urls.put(None)
        
        return urls
    except mariadb.Error as error:
        mariadb_connection.rollback()
        print("Failed retrieving itemtypes from tbl_items table {}".format(error))
    finally:
        if mariadb_connection.is_connected():
            cursor.close()
            mariadb_connection.close()

第二步:修改import_mo()函数,处理终止信号

在循环里检查拿到的URL是否是终止符,是的话就退出循环:

def import_mo(urls):
    list_mo_esi = []
    try:
        mariadb_connection = mariadb.connect(host='host', database='db', user='user', password='pw')
        cursor = mariadb_connection.cursor()
        while True:
            url = urls.get()
            # 收到终止信号,退出循环
            if url is None:
                urls.task_done()
                break
            
            # 原有的请求逻辑
            s = requests.Session()
            retries = Retry(total=5, backoff_factor=1, status_forcelist=[502, 503, 504])
            s.mount('https://', HTTPAdapter(max_retries=retries))
            jsonraw = s.get(url)
            jsondata = ujson.loads(jsonraw.text)
            
            # 原有的数据库新增/更新逻辑
            for row in jsondata:
                cursor.execute('SELECT order_id from tbl_mo WHERE order_id = %s', (row['order_id'], ))
                exists_mo = cursor.fetchall()
                list_mo_esi.append(row['order_id'])
                if len(exists_mo) != 0:
                    print("updating order#", row['order_id'])
                    cursor.execute('UPDATE tbl_mo SET volume = %s, price = %s WHERE order_id = %s', (row['volume_remain'], row['price'], row['order_id'], ))
                    mariadb_connection.commit()
                else:
                    cursor.execute('INSERT INTO tbl_mo (type_id, order_id, ordertype,volume, price) VALUES (%s,%s,%s,%s,%s)', (row['type_id'], row['order_id'], row['is_buy_order'], row['volume_remain'], row['price'], ))
                    mariadb_connection.commit()
            
            urls.task_done()
    except mariadb.Error as error:
        mariadb_connection.rollback()
        print("Failed retrieving itemtypes from tbl_items table {}".format(error))
    finally:
        # 修正变量顺序:先查询数据库里的订单列表,再计算要清理的内容
        cursor.execute('SELECT order_id FROM tbl_mo')
        # fetchall返回的是元组列表,转成普通列表
        list_mo_sql = [row[0] for row in cursor.fetchall()]
        list_mo_purge = list(set(list_mo_sql)-set(list_mo_esi))
        
        # 执行实际的清理逻辑
        if list_mo_purge:
            placeholders = ', '.join(['%s']*len(list_mo_purge))
            cursor.execute(f'DELETE FROM tbl_mo WHERE order_id IN ({placeholders})', list_mo_purge)
            mariadb_connection.commit()
            print(f"清理了{len(list_mo_purge)}条过期订单")
        
        print(f"本次同步了{len(list_mo_esi)}条订单")
        print(f"数据库现有{len(list_mo_sql)}条订单")
        
        if mariadb_connection.is_connected():
            cursor.close()
            mariadb_connection.close()

3. 修复finally块里的变量顺序问题

你原来的finally块里,list_mo_purge = list(set(list_mo_sql)-set(list_mo_esi))在list_mo_sql定义之前,这会抛出NameError,上面的代码已经把顺序调整过来了,并且把fetchall()的结果转成了普通列表(因为fetchall()返回的是(order_id,)这样的元组)。

这样修改后,worker线程会在处理完所有URL并收到终止信号后,正常退出循环,走到finally块执行清理逻辑,主线程会等待所有worker线程完成后才退出,你的数据库关闭、旧数据清理代码就能正常运行了。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.12 05:19:04