Cryptofeed多进程模式下无法向Queue推送数据问题排查
问题描述
使用cryptofeed通过多进程获取代币行情数据,再用worker进程推送至InfluxDB时遇到两个问题:
- feed进程正常运行,但worker执行
queue.get(block=False)时始终提示队列为空,抛出Empty异常; - 修改代码后出现
FileNotFoundError,报错显示SemLock重建失败。
原代码:
queue = multiprocessing.Queue() async def ticker_function(data, time): # Send data to queue queue.put(data, block=False) def new_feed(coin): f = FeedHandler() f.add_feed(cryptofeed.exchanges.Binance( channels=[TICKER], symbols=[f"{coin}-USDT"], callbacks={TICKER: ticker_function} )) f.run() def worker(queue): print("Starting worker") while True: if not queue.empty(): data = queue.get(block=False) # THIS IS WHERE IT SAYS QUEUE EMPTY json_payload = { "measurement": data.symbol, "time": int(time.time()*1000000000), "fields": { 'bid': float(data.bid), 'ask': float(data.ask), } } # Push to influxdb write_api.write(bucket='test_bucket', org='mydb', record=json_payload) else: time.sleep(0.01) if __name__ == '__main__': data_process = multiprocessing.Process(target=worker, args=(queue,)) data_process.start() feed_process = multiprocessing.Process(target=new_feed, args=('BTC',)) feed_process.start() # To stop the main process from ending, adding a join (won't ever trigger) feed_process.join()
首次报错:
Exception has occurred: Empty exception: no description File "/Users/user/sandbox/live-pnl/main.py", line 90, in worker data = queue.get(block=False) File "<string>", line 1, in <module>
修改后报错:
Traceback (most recent call last): File "<string>", line 1, in <module> File "/opt/homebrew/Cellar/python@3.10/3.10.6_2/Frameworks/Python.framework/Versions/3.10/lib/python3.10/multiprocessing/spawn.py", line 116, in spawn_main exitcode = _main(fd, parent_sentinel) File "/opt/homebrew/Cellar/python@3.10/3.10.6_2/Frameworks/Python.framework/Versions/3.10/lib/python3.10/multiprocessing/spawn.py", line 126, in _main self = reduction.pickle.load(from_parent) File "/opt/homebrew/Cellar/python@3.10/3.10.6_2/Frameworks/Python.framework/Versions/3.10/lib/python3.10/multiprocessing/synchronize.py", line 110, in __setstate__ self._semlock = _multiprocessing.SemLock._rebuild(*state) FileNotFoundError: [Errno 2] No such file or directory
问题原因
- 队列实例不共享:macOS默认使用
spawn模式启动子进程,全局定义的multiprocessing.Queue()会在子进程重新导入模块时被重新创建,导致feed进程和worker进程操作的是完全独立的两个队列,所以worker永远取不到数据。 - SemLock报错:本质是跨进程同步对象的序列化/重建问题。如果修改时错误处理了队列的传递逻辑,spawn模式下无法正确重建同步锁相关资源,就会触发该错误。
修复方案及代码
核心改动是确保队列实例在主进程创建并传递给子进程,同时优化队列操作和InfluxDB客户端的初始化逻辑:
import multiprocessing import time from cryptofeed import FeedHandler from cryptofeed.exchanges import Binance from cryptofeed.defines import TICKER from influxdb_client import InfluxDBClient, WriteOptions async def ticker_function(data, timestamp, queue): queue.put(data) def new_feed(coin, queue): # 用偏函数绑定队列到回调,避免闭包作用域问题 from functools import partial bound_callback = partial(ticker_function, queue=queue) f = FeedHandler() f.add_feed(Binance( channels=[TICKER], symbols=[f"{coin}-USDT"], callbacks={TICKER: bound_callback} )) f.run() def worker(queue): print("Starting worker") # 在worker内部独立初始化InfluxDB客户端,避免跨进程共享 client = InfluxDBClient(url="你的InfluxDB地址", token="你的访问token", org="mydb") write_api = client.write_api(write_options=WriteOptions(batch_size=10)) while True: # 阻塞等待队列数据,避免竞态条件和空轮询 data = queue.get() json_payload = { "measurement": data.symbol, "time": int(time.time() * 1000000000), "fields": { 'bid': float(data.bid), 'ask': float(data.ask), } } try: write_api.write(bucket='test_bucket', org='mydb', record=json_payload) except Exception as e: print(f"写入InfluxDB失败: {e}") if __name__ == '__main__': # 队列必须在主进程创建,再传递给子进程 queue = multiprocessing.Queue() data_process = multiprocessing.Process(target=worker, args=(queue,)) data_process.start() feed_process = multiprocessing.Process(target=new_feed, args=('BTC', queue)) feed_process.start() feed_process.join() data_process.join()
关键改动说明:
- 将队列创建移至
if __name__ == '__main__':块内,作为参数传递给子进程,确保所有进程共享同一个队列实例。 - 用
functools.partial将队列绑定到回调函数,避免闭包作用域导致的队列访问错误。 - 移除
queue.empty()判断和block=False参数,改用阻塞式queue.get(),既避免竞态条件,又减少空轮询的资源消耗。 - 在worker进程内独立初始化InfluxDB客户端,避免跨进程共享客户端对象引发的异常。
内容的提问来源于stack exchange,提问作者fvim
相关产品推荐
相关产品推荐

