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
相关产品推荐
相关产品推荐

