如何解决ThreadPoolExecutor线程因网络请求无响应导致的阻塞问题
我运行着7×24小时的脚本,使用concurrent.futures.ThreadPoolExecutor同时向Broker发送请求,每个线程执行查询账户余额或发送订单这类简单任务。但当线程因请求未收到响应时,有时会陷入阻塞。后续操作依赖这些线程获取的数据,因此必须等线程执行完成或失败才能继续推进。
以下是订单函数及concurrent.futures的示例代码:
def sell_orders_0(): order = client.new_order(symbol = ticker, side = 'SELL', type = 'LIMIT_MAKER', quantity = xquantity, price = xprice) def buy_orders_0(): order = client.new_order(symbol = ticker, side = 'SELL', type = 'LIMIT_MAKER', quantity = xquantity, price = xprice) with concurrent.futures.ThreadPoolExecutor() as executor: sells_0 = executor.submit(sell_orders_0) buys_0 = executor.submit(buy_orders_0)
通过hanging Threads模块捕获到的线程挂起调用栈如下:
Thread 139646566659840 "ThreadPoolExecutor-666849_1" hangs - File "/usr/lib/python3.9/threading.py", line 912, in _bootstrap self._bootstrap_inner() File "/usr/lib/python3.9/threading.py", line 954, in _bootstrap_inner self.run() File "/usr/lib/python3.9/threading.py", line 892, in run self._target(*self._args, **self._kwargs) File "/usr/lib/python3.9/concurrent/futures/thread.py", line 77, in _worker work_item.run() File "/usr/lib/python3.9/concurrent/futures/thread.py", line 52, in run result = self.fn(*self.args, **self.kwargs) File "/home/user/binance_bot.py", line 1346, in sell_orders_3 order = client.new_order(symbol = ticker, side = 'SELL', type = 'LIMIT_MAKER', quantity = xquantity, price = xprice) File "/home/user/.local/lib/python3.9/site-packages/binance/spot/trade.py", line 68, in new_order return self.sign_request("POST", url_path, params) File "/home/user/.local/lib/python3.9/site-packages/binance/api.py", line 83, in sign_request return self.send_request(http_method, url_path, payload) File "/home/user/.local/lib/python3.9/site-packages/binance/api.py", line 115, in send_request response = self._dispatch_request(http_method)(**params) File "/home/user/.local/lib/python3.9/site-packages/requests/sessions.py", line 635, in post return self.request("POST", url, data=data, json=json, **kwargs) File "/home/user/.local/lib/python3.9/site-packages/requests/sessions.py", line 587, in request resp = self.send(prep, **send_kwargs) File "/home/user/.local/lib/python3.9/site-packages/requests/sessions.py", line 701, in send r = adapter.send(request, **kwargs) File "/home/user/.local/lib/python3.9/site-packages/requests/adapters.py", line 489, in send resp = conn.urlopen( File "/home/user/.local/lib/python3.9/site-packages/urllib3/connectionpool.py", line 703, in urlopen httplib_response = self._make_request( File "/home/user/.local/lib/python3.9/site-packages/urllib3/connectionpool.py", line 444, in _make_request httplib_response = conn.getresponse() File "/usr/lib/python3.9/http/client.py", line 1347, in getresponse response.begin() File "/usr/lib/python3.9/http/client.py", line 307, in begin version, status, reason = self._read_status() File "/usr/lib/python3.9/http/client.py", line 268, in _read_status line = str(self.fp.readline(_MAXLINE + 1), "iso-8859-1") File "/usr/lib/python3.9/socket.py", line 704, in readinto [0/1893] return self._sock.recv_into(b) File "/usr/lib/python3.9/ssl.py", line 1241, in recv_into return self.read(nbytes, buffer) File "/usr/lib/python3.9/ssl.py", line 1099, in read return self._sslobj.read(len, buffer) Thread 139646533089024 "ThreadPoolExecutor-666849_0" hangs - File "/usr/lib/python3.9/threading.py", line 912, in _bootstrap self._bootstrap_inner() File "/usr/lib/python3.9/threading.py", line 954, in _bootstrap_inner self.run() File "/usr/lib/python3.9/threading.py", line 892, in run self._target(*self._args, **self._kwargs) File "/usr/lib/python3.9/concurrent/futures/thread.py", line 77, in _worker work_item.run() File "/usr/lib/python3.9/concurrent/futures/thread.py", line 52, in run result = self.fn(*self.args, **self.kwargs) File "/home/user/binance_bot.py", line 1298, in sell_orders_0 order = client.new_order(symbol = ticker, side = 'SELL', type = 'LIMIT_MAKER', quantity = xquantity, price = xprice) File "/home/user/.local/lib/python3.9/site-packages/binance/spot/trade.py", line 68, in new_order return self.sign_request("POST", url_path, params) File "/home/user/.local/lib/python3.9/site-packages/binance/api.py", line 83, in sign_request return self.send_request(http_method, url_path, payload) File "/home/user/.local/lib/python3.9/site-packages/binance/api.py", line 115, in send_request response = self._dispatch_request(http_method)(**params) File "/home/user/.local/lib/python3.9/site-packages/requests/sessions.py", line 635, in post return self.request("POST", url, data=data, json=json, **kwargs) File "/home/user/.local/lib/python3.9/site-packages/requests/sessions.py", line 587, in request resp = self.send(prep, **send_kwargs) File "/home/user/.local/lib/python3.9/site-packages/requests/sessions.py", line 701, in send r = adapter.send(request, **kwargs) File "/home/user/.local/lib/python3.9/site-packages/requests/adapters.py", line 489, in send resp = conn.urlopen( File "/home/user/.local/lib/python3.9/site-packages/urllib3/connectionpool.py", line 703, in urlopen httplib_response = self._make_request( File "/home/user/.local/lib/python3.9/site-packages/urllib3/connectionpool.py", line 444, in _make_request httplib_response = conn.getresponse() File "/usr/lib/python3.9/http/client.py", line 1347, in getresponse response.begin() File "/usr/lib/python3.9/http/client.py", line 307, in begin version, status, reason = self._read_status() File "/usr/lib/python3.9/http/client.py", line 268, in _read_status line = str(self.fp.readline(_MAXLINE + 1), "iso-8859-1") File "/usr/lib/python3.9/socket.py", line 704, in readinto return self._sock.recv_into(b) File "/usr/lib/python3.9/ssl.py", line 1241, in recv_into return self.read(nbytes, buffer) File "/usr/lib/python3.9/ssl.py", line 1099, in read return self._sslobj.read(len, buffer) ---------Thread 139646838966080 "MainThread" hangs --------- File "/home/user/binance_bot.py", line 1596, in <module> buys_0 = executor.submit(buy_orders_0) File "/usr/lib/python3.9/concurrent/futures/_base.py", line 628, in __exit__ self.shutdown(wait=True) File "/usr/lib/python3.9/concurrent/futures/thread.py", line 229, in shutdown t.join() File "/usr/lib/python3.9/threading.py", line 1033, in join self._wait_for_tstate_lock() File "/usr/lib/python3.9/threading.py", line 1049, in _wait_for_tstate_lock elif lock.acquire(block, timeout):
我曾尝试让代码挂起时抛出异常,但由于无法从外部取消运行中的线程,该方案不可行。现在想知道如何让代码实现非阻塞,避免线程挂起导致整个脚本卡住。
1. 为HTTP请求设置超时时间
从调用栈可以看到,阻塞发生在底层Socket的read操作,根源是client.new_order发起的HTTP请求没有设置超时。直接在请求层面添加超时,是最有效的解决方式:
如果使用Binance官方SDK,可在初始化客户端时设置全局超时,或在调用new_order时单独传入超时参数:
# 初始化客户端时设置全局超时(单位:秒) client = Client(api_key, api_secret, timeout=10) # 或者单独为new_order请求设置超时 def sell_orders_0(): order = client.new_order( symbol=ticker, side='SELL', type='LIMIT_MAKER', quantity=xquantity, price=xprice, timeout=10 # 单独设置超时 )
如果SDK不支持直接传超时,可通过修改requests会话配置添加全局超时:
import requests from binance import Client session = requests.Session() session.timeout = 10 # 设置全局超时 client = Client(api_key, api_secret, session=session)
2. 为Future添加超时获取逻辑
即使请求本身未设置超时,也可在获取Future结果时添加超时,避免主线程无限等待:
with concurrent.futures.ThreadPoolExecutor() as executor: sells_0 = executor.submit(sell_orders_0) buys_0 = executor.submit(buy_orders_0) try: # 等待结果,最多等待15秒 sell_result = sells_0.result(timeout=15) except concurrent.futures.TimeoutError: print("卖单请求超时,标记为失败") # 添加失败处理逻辑,比如重试或记录日志 try: buy_result = buys_0.result(timeout=15) except concurrent.futures.TimeoutError: print("买单请求超时,标记为失败")
注意:这种方式仅阻止主线程等待,挂起的线程仍会占用资源,长期运行可能耗尽线程池。因此必须配合请求层面的超时一起使用。
3. 改用异步HTTP库
如果脚本架构允许,可使用aiohttp这类异步HTTP库替代同步的requests,结合asyncio实现非阻塞并发请求,从根源避免线程阻塞:
import asyncio import aiohttp async def sell_orders_0(): async with aiohttp.ClientSession(timeout=aiohttp.ClientTimeout(total=10)) as session: # 构造Binance API请求(需自行处理签名,或使用异步版Binance SDK) url = "https://api.binance.com/api/v3/order" params = { "symbol": ticker, "side": "SELL", "type": "LIMIT_MAKER", "quantity": xquantity, "price": xprice, # 添加签名参数... } async with session.post(url, params=params) as resp: return await resp.json() async def buy_orders_0(): # 类似卖单逻辑 pass async def main(): sell_task = asyncio.create_task(sell_orders_0()) buy_task = asyncio.create_task(buy_orders_0()) await asyncio.gather(sell_task, buy_task, return_exceptions=True) asyncio.run(main())
异步模式下所有IO操作都是非阻塞的,不会出现线程挂起,更适合7×24运行的脚本。
4. 自定义线程超时回收机制
如果无法修改请求超时,可使用守护线程+超时重启的方式回收挂起线程:
import threading import time def run_with_timeout(func, timeout=10): result = None exception = None def wrapper(): nonlocal result, exception try: result = func() except Exception as e: exception = e thread = threading.Thread(target=wrapper, daemon=True) thread.start() thread.join(timeout) if thread.is_alive(): raise TimeoutError(f"函数{func.__name__}执行超时") if exception: raise exception return result # 使用方式 try: sell_result = run_with_timeout(sell_orders_0, timeout=10) except TimeoutError: print("卖单超时")
这种方式通过守护线程运行任务,超时后主线程不再等待,但如果主线程长期运行,挂起的守护线程仍会占用资源,仅作为临时方案,优先推荐前两种方法。
内容的提问来源于stack exchange,提问作者PyNoob

