19、多线程与并发:线程安全队列,无锁数据结构,协程(asyncio/gevent),GIL问题与解决方案
做衍生品做市,说白了就是跟时间赛跑。行情数据一来,你得在微秒级别内做出反应。我刚开始做这套系统时,天真地以为多开几个线程就能搞定一切。结果呢?数据错乱、死锁、性能不升反降……踩过的坑能写一本血泪史。
今天咱们就聊聊并发编程里的几个核心问题。嗯,都是我在实战中反复摔打过的经验。
19.1 线程安全队列:别让数据打架
做市系统里,行情线程、策略线程、风控线程、下单线程……它们之间要传递数据。最简单的办法就是共享一个列表,然后加锁。但锁用不好,性能就崩了。
我个人习惯用 queue.Queue。它是 Python 标准库自带的线程安全队列,内部已经帮你处理好了锁的细节。
import queue
import threading
# 创建一个最大容量为1000的队列
order_queue = queue.Queue(maxsize=1000)
def producer():
for i in range(10):
order = {"id": i, "price": 100.5 + i}
order_queue.put(order)
print(f"生产订单: {order}")
def consumer():
while True:
order = order_queue.get()
print(f"消费订单: {order}")
order_queue.task_done()
# 启动生产者和消费者线程
t1 = threading.Thread(target=producer)
t2 = threading.Thread(target=consumer)
t1.start()
t2.start()
t1.join()
order_queue.join() # 等待所有任务完成
你看,代码很简洁。但要注意一点:maxsize 一定要设置。我曾经在生产环境里忘了设这个参数,结果行情爆发时,队列被撑爆了,内存直接打满。嗯,那是个惨痛的教训。
19.2 无锁数据结构:性能的终极追求
锁虽然好用,但竞争激烈时,线程会被阻塞,上下文切换的成本很高。在微秒级别的交易系统里,这可能是致命的。
无锁数据结构,说白了就是利用 CPU 的原子操作(比如 CAS——Compare And Swap)来保证线程安全,不加锁。
Python 里有个 atomicwrites 库,但更常见的是用 multiprocessing 的 Value 和 Array,它们底层就是无锁的。
from multiprocessing import Value, Process
# 创建一个共享的整数,初始值为0
counter = Value('i', 0)
def increment():
for _ in range(100000):
with counter.get_lock(): # 这里其实还是用了锁
counter.value += 1
# 但真正的无锁实现需要自己用 ctypes 和原子操作
# 这里只是示意,实际项目中我会用 C 扩展
说实话,Python 里做真正的无锁数据结构挺麻烦的。因为 GIL 的存在,很多原子操作被限制了。我一般会在 C++ 层面实现无锁队列,然后通过 Python 的 C 扩展调用。你想想看,行情数据从网卡到策略引擎,中间每少一次锁,就能快几微秒。
collections.deque 配合 threading.Event 做简单的生产者-消费者模型。虽然不完全是「无锁」,但性能比 Queue 好不少。
19.3 协程:asyncio 与 gevent
线程切换是操作系统干的,开销大。协程切换是程序自己控制的,开销极小。做市系统里,大量时间花在等待网络 I/O 上——比如等待交易所的行情推送、等待订单回报。
这时候用协程,效率极高。
19.3.1 asyncio:官方推荐
Python 3.4 之后,asyncio 成了标准库。我个人习惯用它来写网络 I/O 密集型的模块,比如行情接收器。
import asyncio
async def fetch_market_data(symbol):
# 模拟网络请求
await asyncio.sleep(0.1)
return {"symbol": symbol, "price": 100.5}
async def main():
symbols = ["BTC", "ETH", "SOL"]
tasks = [fetch_market_data(s) for s in symbols]
results = await asyncio.gather(*tasks)
print(results)
asyncio.run(main())
你看,asyncio.gather 可以并发执行多个协程。但要注意,它只适合 I/O 密集型任务。如果你在协程里做 CPU 计算,它会阻塞整个事件循环。嗯,这里要特别小心。
19.3.2 gevent:隐式协程
gevent 跟 asyncio 不同,它通过 monkey patch 把标准库里的阻塞调用变成非阻塞的。说白了,你写同步代码,它自动帮你转成异步。
from gevent import monkey
monkey.patch_all() # 打补丁,替换标准库
import gevent
import time
def task(name):
print(f"开始任务: {name}")
time.sleep(1) # 这里本来是阻塞的,但被 gevent 替换成了非阻塞
print(f"完成任务: {name}")
# 并发执行三个任务
gevent.joinall([
gevent.spawn(task, "A"),
gevent.spawn(task, "B"),
gevent.spawn(task, "C"),
])
gevent 的好处是代码改动小。但坏处也很明显——monkey patch 可能引发奇怪的 bug。我曾经在生产环境里遇到过一次,因为 patch 了 socket 模块,导致某个第三方库崩溃。排查了整整两天……从那以后,我尽量用 asyncio。
19.4 GIL 问题与解决方案
GIL,全局解释器锁。这是 Python 的「原罪」。它保证同一时刻只有一个线程在执行 Python 字节码。说白了,多线程在 CPU 密集型任务上,不仅不快,反而更慢。
为什么会这样?因为 GIL 的存在,让 Python 的多线程变成了「伪并发」。
19.4.1 解决方案一:多进程
每个进程有自己的 GIL,所以可以真正并行。用 multiprocessing 模块,把 CPU 密集型任务分到多个进程里。
from multiprocessing import Pool
def calculate_volatility(prices):
# 假设这是一个 CPU 密集型的计算
return sum(prices) / len(prices)
if __name__ == "__main__":
with Pool(4) as p: # 启动4个进程
results = p.map(calculate_volatility, [price_list1, price_list2])
但多进程也有代价:进程间通信(IPC)比线程间通信慢得多。我一般只在策略回测或风险计算时用多进程,实时交易还是用线程+协程的组合。
19.4.2 解决方案二:C 扩展
把 CPU 密集型的代码用 C/C++ 写,然后通过 Python 调用。C 扩展在执行时可以释放 GIL。
// 伪代码示意
static PyObject* fast_calculation(PyObject* self, PyObject* args) {
Py_BEGIN_ALLOW_THREADS // 释放 GIL
// 在这里做密集计算
Py_END_ALLOW_THREADS // 重新获取 GIL
Py_RETURN_NONE;
}
我团队里有个同事,把做市策略里的定价模型用 C++ 重写了一遍,性能提升了 10 倍。嗯,这就是 C 扩展的魅力。
19.4.3 解决方案三:协程 + 异步 I/O
对于 I/O 密集型任务,协程完全绕过了 GIL 的问题。因为协程在等待 I/O 时,根本不需要 CPU。说白了,GIL 只在 CPU 计算时才有影响。
19.5 知识体系总览
下面这张图,是我自己总结的并发编程知识体系。你看一眼,心里就有数了。
19.6 实战中的选择
说了这么多,到底怎么选?我个人的经验是这样的:
| 场景 | 推荐方案 | 原因 |
|---|---|---|
| 行情数据接收(I/O密集) | asyncio 协程 | 轻量、高效、标准库支持 |
| 策略计算(CPU密集) | 多进程 或 C扩展 | 绕过GIL,真正并行 |
| 订单管理(混合型) | 线程 + 协程组合 | 线程处理状态同步,协程处理I/O |
| 风控检查(低延迟) | 无锁数据结构 + C扩展 | 微秒级响应,不能有锁竞争 |
你想想看,做市系统里,每一微秒都可能是利润。选对并发模型,比选对策略参数还重要。嗯,今天就聊到这儿,下次咱们聊聊内存管理和对象池——那也是性能优化的重头戏。