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

Cryptofeed多进程模式下无法向Queue推送数据问题排查

问题描述

使用cryptofeed通过多进程获取代币行情数据,再用worker进程推送至InfluxDB时遇到两个问题:

  1. feed进程正常运行,但worker执行queue.get(block=False)时始终提示队列为空,抛出Empty异常;
  2. 修改代码后出现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()

关键改动说明:

  1. 将队列创建移至if __name__ == '__main__':块内,作为参数传递给子进程,确保所有进程共享同一个队列实例。
  2. 用functools.partial将队列绑定到回调函数,避免闭包作用域导致的队列访问错误。
  3. 移除queue.empty()判断和block=False参数,改用阻塞式queue.get(),既避免竞态条件,又减少空轮询的资源消耗。
  4. 在worker进程内独立初始化InfluxDB客户端,避免跨进程共享客户端对象引发的异常。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 16:15:26