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

如何测试Pandas DataFrame的竞态条件?锁机制线程安全验证

问题描述

我想用schedule库每隔若干秒运行一些修改全局DataFrame的函数。我知道Pandas不是线程安全的,所以在每个函数调用里加了锁来规避风险。下面的极简示例代码运行符合预期,但我不确定怎么检查代码是否会出现竞态条件。能不能指导如何正确测试?或者仅因为用了with lock,就可以直接认定代码是线程安全的?

示例代码:

import schedule
import time, datetime
import pandas as pd
from threading import Lock

data = [(3,5,7), (2,4,6),(5,8,9)]
df = pd.DataFrame(data, columns = ['A','B','C'])
lock = Lock()

def job(lock):
    global df
    with lock:
        df = pd.concat([df, df.iloc[0:2]])
        print('J1', datetime.datetime.utcnow(), len(df), df.A.sum())
    
def job2(lock):
    global df
    with lock:
        df = pd.concat([df, df.iloc[0:2]])
        print('J2', datetime.datetime.utcnow(), len(df), df.A.sum())

schedule.every(0.75).seconds.do(job, lock=lock)
schedule.every(0.25).seconds.do(job2, lock=lock)

while True:
    schedule.run_pending()
    time.sleep(0.1)

解答

1. 仅用with lock是否能保证线程安全?

你的代码里,所有修改全局df的操作都被包裹在同一个Lock的with块中,这确实能保证这些操作互斥执行——同一时间只有一个线程能进入锁保护的代码段,不会出现多个线程同时修改df的情况。

不过要注意两个关键细节,确保真正的线程安全:

  • 必须保证所有访问或修改df的代码路径都使用同一个锁。如果后续新增了其他操作df的函数,也必须用这个锁包裹逻辑,否则仍会出现竞态条件。
  • 代码中df = pd.concat(...)是重新赋值全局变量,这个完整流程(读取旧df→执行concat→赋值新df)都在锁内完成,不会出现“一个线程读取旧值还没完成赋值,另一个线程又读取旧值”的冲突场景。

就当前代码来看,只要锁的使用逻辑保持一致(所有共享资源操作都在锁内),可以认定是线程安全的。

2. 如何测试竞态条件?

如果要验证代码是否真的不存在竞态条件,可以从以下几个方向测试:

(1)放大并发负载

把任务执行间隔调至更短,或增加更多并发任务,让线程冲突的概率变大。比如将job2的间隔改成0.01秒,再新增job3、job4等同类任务,运行一段时间后检查结果:

  • 核对每次操作后df的长度、A列总和是否符合逻辑(比如每次concat两行,长度应+2,总和应加上前两行A的和)。
  • 若出现长度或总和不符合预期的情况,说明存在竞态条件。

(2)添加断言验证一致性

在锁内操作完成后,添加断言来验证数据的正确性。示例如下:

def job(lock):
    global df
    with lock:
        prev_len = len(df)
        prev_sum = df.A.sum()
        # 记录要添加的行的A列和
        added_rows_sum = df.iloc[0:2].A.sum()
        
        df = pd.concat([df, df.iloc[0:2]])
        
        # 验证长度是否符合预期
        assert len(df) == prev_len + 2, f"长度异常:预期{prev_len+2},实际{len(df)}"
        # 验证总和是否符合预期
        assert df.A.sum() == prev_sum + added_rows_sum, f"总和异常:预期{prev_sum+added_rows_sum},实际{df.A.sum()}"
        print('J1', datetime.datetime.utcnow(), len(df), df.A.sum())

给job2也加上类似断言,运行代码后若断言未触发,说明无竞态条件;若断言报错,则存在线程安全问题。

(3)模拟极端并发场景

用threading.Thread手动创建多个线程同时执行修改操作,直接测试并发冲突:

import threading

def run_job_repeatedly(lock, repeat_times):
    global df
    for _ in range(repeat_times):
        with lock:
            df = pd.concat([df, df.iloc[0:2]])

# 创建10个线程,每个线程执行100次修改
threads = []
for _ in range(10):
    t = threading.Thread(target=run_job_repeatedly, args=(lock, 100))
    threads.append(t)
    t.start()

# 等待所有线程执行完毕
for t in threads:
    t.join()

# 验证最终长度:初始3行,总操作次数10*100=1000次,每次加2行,最终应为3+1000*2=2003
assert len(df) == 2003, f"最终长度异常:预期2003,实际{len(df)}"

如果断言通过,说明锁的保护有效;若不通过,则存在竞态条件。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 05:03:23