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

异步锁等待期如何实时监测DataFrame最后一行值

解决方案:多线程+协程场景下的实时条件等待

核心问题分析

  • asyncio.Lock仅适用于协程间同步,无法跨线程生效,线程间共享资源需用threading.Lock或threading.Condition。
  • Thread2中每次创建的df是_download的快照,无法实时获取Thread1的更新,导致等待条件时看不到最新状态。
  • 用time.sleep轮询会阻塞线程/协程,无法及时响应条件变化,还会浪费资源。

解决方法

  1. 用threading.Condition实现线程安全的条件等待:结合锁与等待通知机制,条件满足时唤醒任务,避免无效轮询。
  2. 共享资源线程安全访问:全局_download的读写加锁保护,确保多线程数据一致性。
  3. 协程实时获取最新数据:不传递快照,需要时加锁读取最新共享资源。
  4. 避免阻塞协程:用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 14:06:56