Polars DataFrame 本身不是线程安全的,尤其在多线程并发修改(如 df[i, col] += 1)时存在数据竞争风险;虽因 CPython GIL 可能偶然“看似正确”,但行为未定义,必须通过显式同步(如 threading.Lock 或任务队列)保障一致性。
polars dataframe 本身
不是线程安全的
,尤其在多线程并发修改(如 df[i, col] += 1)时存在数据竞争风险;虽因 cpython gil 可能偶然“看似正确”,但行为未定义,必须通过显式同步(如 threading.lock 或任务队列)保障一致性。
Polars 是一个以高性能和并行计算见长的 DataFrame 库,其底层由 Rust 实现,并通过 Arrow 内存模型管理数据。然而,
线程安全性不能仅凭“运行结果看起来正常”来判断
。正如你在示例中观察到的 df[2, 'b'] += 1 在多线程下输出了预期数值(100002),这本质上是 GIL(全局解释器锁)的副作用——它强制 Python 字节码层面的 += 操作串行化,
掩盖了底层内存操作的竞争问题
。
⚠️ 关键事实:
GIL 仅保护 Python 解释器状态,不保护 Polars 的 Rust 数据结构
;
df[row, col] += value 等就地(in-place)操作会触发 Polars 内部的内存写入,而 Rust 层无 Python 级别的锁保护;
多线程并发执行此类操作,可能导致原子性丢失、越界写入、引用计数错误,甚至段错误(segfault)——尤其在启用 polars.polars_enable_string_cache() 或涉及 list/struct 等复杂类型时更易暴露。
✅ 安全实践推荐:
读多写少场景 → 使用 threading.Lock
对共享 DataFrame 的写操作加锁,简单直接:
高吞吐写入场景 → 使用生产者-消费者队列
将修改操作序列化为指令(如 (op, row, col, delta)),由单一线程消费并批量执行,既规避竞争,又保留 Polars 自身的向量化并行能力:
.collect() 是线程安全的
LazyFrame.collect() 仅读取数据并返回新 DataFrame,不修改原对象,且内部 Rust 执行引擎天然支持多线程调度(通过 Rayon)。多个线程可安全并发调用 .collect(),无需额外同步。
? 总结:
永远不要依赖“看起来正常”的多线程就地修改
;
Polars 的性能优势源于
单线程内向量化 + Rust 并行计算
,而非跨 Python 线程的并发写入;
正确做法是:
将并发逻辑收口到 Python 层同步,再交由 Polars 高效执行
——既安全,又不失性能。
import threading
import polars as pl
df = pl.DataFrame({"a": range(5), "b": range(5)})
lock = threading.Lock()
def safe_increment(col: str):
for _ in range(100_000):
with lock:
df[2, col] += 1
t1 = threading.Thread(target=safe_increment, args=("a",))
t2 = threading.Thread(target=safe_increment, args=("b",))
t1.start(); t2.start()
t1.join(); t2.join()
print(df) # ✅ 确定性结果from queue import Queue
import threading
q = Queue()
_STOP = object()
def worker():
while True:
item = q.get()
if item is _STOP:
break
df[item[0], item[1]] += item[2] # 原子性单次更新
# 启动工作线程
worker_t = threading.Thread(target=worker)
worker_t.start()
# 多线程提交任务(无锁)
def submit_updates(col: str):
for _ in range(100_000):
q.put((2, col, 1))
threading.Thread(target=submit_updates, args=("a",)).start()
threading.Thread(target=submit_updates, args=("b",)).start()
# 结束信号
q.put(_STOP)
worker_t.join()