异步锁等待期如何实时监测DataFrame最后一行值
解决方案:多线程+协程场景下的实时条件等待
核心问题分析
asyncio.Lock仅适用于协程间同步,无法跨线程生效,线程间共享资源需用threading.Lock或threading.Condition。- Thread2中每次创建的
df是_download的快照,无法实时获取Thread1的更新,导致等待条件时看不到最新状态。 - 用
time.sleep轮询会阻塞线程/协程,无法及时响应条件变化,还会浪费资源。
解决方法
- 用
threading.Condition实现线程安全的条件等待:结合锁与等待通知机制,条件满足时唤醒任务,避免无效轮询。 - 共享资源线程安全访问:全局
_download的读写加锁保护,确保多线程数据一致性。 - 协程实时获取最新数据:不传递快照,需要时加锁读取最新共享资源。
- 避免阻塞协程:用
asyncio.to_thread将阻塞同步操作放到线程池,不阻塞事件循环。
修改后的完整代码
import pandas as pd import numpy as np import threading import asyncio import time # 全局共享资源和同步对象 _download = {} _shared_lock = threading.Lock() _condition = threading.Condition(_shared_lock) # 基于共享锁的条件变量 stop_flag = False # 终止标志 class Thread1(threading.Thread): def __init__(self): super().__init__() def run(self): # 初始化测试用DataFrame ind_df = pd.DataFrame({ 'value': np.arange(100), 'flag': np.random.rand(100) }) data = self.get_data(ind_df) count = 0 global _download, stop_flag while not stop_flag: with _shared_lock: # 更新共享资源时加锁 _download[count] = next(data).values _condition.notify_all() # 更新后通知等待的线程 time.sleep(0.5) count += 1 def get_data(self, df): for idx in range(df.shape[0]): yield df.iloc[idx] # 循环生成数据(避免迭代完停止) while True: yield df.iloc[np.random.randint(0, df.shape[0])] class Thread2(threading.Thread): def __init__(self): super().__init__() def run(self): global stop_flag while not stop_flag: asyncio.run(self.apply_coroutines()) time.sleep(1) async def apply_coroutines(self): await asyncio.gather( self.coroutine1(), self.coroutine2(), ) async def coroutine1(self): # 把阻塞的同步操作放到线程池,避免阻塞协程事件循环 await asyncio.to_thread(self._wait_for_flag_condition) def _wait_for_flag_condition(self): """线程安全的条件等待逻辑""" global _download with _condition: while True: # 加锁读取最新数据 latest_entry = list(_download.values())[-1] if _download else None # 匹配flag≈0.6(容错范围) if latest_entry is not None and abs(latest_entry[1] - 0.6) < 0.01: print(f"条件满足!最新数据: {latest_entry}") break # 等待通知,超时1秒自动唤醒检查 _condition.wait(timeout=1) async def coroutine2(self): print('coroutine 2') await asyncio.sleep(0.5) # 启动线程 if __name__ == "__main__": t = Thread1() t.start() s = Thread2() s.start() try: while True: time.sleep(1) except KeyboardInterrupt: # 终止线程 stop_flag = True t.join() s.join() print("线程已终止")
关键修改说明
threading.Condition通知机制:Thread1更新共享资源后调用notify_all(),能立即唤醒等待的线程,无需等待超时。asyncio.to_thread的作用:将阻塞的条件等待逻辑移到线程执行,保证协程事件循环不被阻塞,coroutine2能正常运行。- 共享资源锁保护:所有对
_download的读写操作都在锁范围内,避免多线程数据竞争。 - 实时数据读取:等待条件时每次读取最新的
_download,确保判断基于实时状态。
内容的提问来源于stack exchange,提问作者bugrahaskan
相关产品推荐
相关产品推荐

