如何在Django的wallet_verify视图中实现多进程并行处理
问题背景
能不能在Django请求处理流程里用多进程?
当向http://127.0.0.1:8000/wallet_verify发请求时,会执行下面的wallet_verify视图函数:
def wallet_verify(request): walelts = botactive.objects.all()
第一步筛选botactive表中active字段为True的用户:
for active in walelts: check_active = active.active if check_active == True: user_is_active = active.user
接着获取活跃用户在Bybitapidatas表中的API密钥信息:
database = Bybitapidatas.objects.filter(user=user_is_active) for apikey in database: apikey = apikey.apikey for apisecret in database: apisecret = apisecret.apisecret
因为调用Bybit交易所接口每次只能用一组API密钥,所以想并行处理每个用户的操作:
for a, b in zip(list(Bybitapidatas.objects.filter(user=user_is_active).values("apikey")), list(Bybitapidatas.objects.filter(user=user_is_active).values("apisecret"))): session =spot.HTTP(endpoint='https://api-testnet.bybit.com/', api_key=a['apikey'], api_secret=b['apisecret'])
之后检查用户USDT余额,余额不足11就跳过,否则执行BTCUSDT市价买入:
GET_USDT_BALANCE = session.get_wallet_balance()['result']['balances'] for i in GET_USDT_BALANCE: if 'USDT' in i.values(): GET_USDT_BALANCE = session.get_wallet_balance()['result']['balances'] idx_USDT = GET_USDT_BALANCE.index(i) GET_USDTBALANCE = session.get_wallet_balance()['result']['balances'][idx_USDT]['free'] print(round(float(GET_USDTBALANCE),2)) if round(float(GET_USDTBALANCE),2) < 11 : pass else: session.place_active_order( symbol="BTCUSDT", side="Buy", type="MARKET", qty=10, timeInForce="GTC" )
遇到的问题
我想遍历数据库时并行处理每个用户的上述操作,但试了用multiprocessing和Pool,提示应用未启动,没法在wallet_verify函数里执行。作为编程新手,想问下怎么在这个处理Post请求的视图函数里实现多进程并行?
解决方案
在Django视图中直接使用multiprocessing会触发应用未初始化的问题,因为子进程不会自动加载Django配置。以下是可行的实现方案:
1. 提取独立处理函数
将单个用户的操作逻辑抽离为独立函数,并且在函数内部初始化Django环境(确保子进程能正常使用Django相关功能):
import os import django from bybit_api import spot # 替换为你实际的Bybit库导入路径 def process_user_operation(apikey, apisecret): # 子进程必须初始化Django环境 os.environ.setdefault("DJANGO_SETTINGS_MODULE", "你的项目名.settings") django.setup() # 执行Bybit操作逻辑 session = spot.HTTP(endpoint='https://api-testnet.bybit.com/', api_key=apikey, api_secret=apisecret) try: # 优化原逻辑:避免重复调用接口,直接遍历余额数据 balance_data = session.get_wallet_balance()['result']['balances'] usdt_balance = 0.0 for item in balance_data: if item['coin'] == 'USDT': usdt_balance = float(item['free']) break print(round(usdt_balance, 2)) if usdt_balance >= 11: session.place_active_order( symbol="BTCUSDT", side="Buy", type="MARKET", qty=10, timeInForce="GTC" ) except Exception as e: # 捕获异常,避免单个进程崩溃影响其他任务 print(f"处理API时出错: {str(e)}")
2. 修改视图函数,使用多进程池
在视图中提前完成所有ORM查询,将结果传递给子进程,避免子进程直接操作数据库:
from django.http import HttpResponse from multiprocessing import Pool from .models import botactive, Bybitapidatas def wallet_verify(request): if request.method != 'POST': return HttpResponse("仅支持POST请求", status=400) # 1. 主进程中查询所有需要处理的API密钥对 active_user_ids = botactive.objects.filter(active=True).values_list('user', flat=True) api_pairs = [] for user_id in active_user_ids: api_datas = Bybitapidatas.objects.filter(user_id=user_id).values('apikey', 'apisecret') api_pairs.extend( (data['apikey'], data['apisecret']) for data in api_datas ) # 2. 启动多进程池并行处理 with Pool(processes=4) as pool: # processes设为CPU核心数或合适值 pool.starmap(process_user_operation, api_pairs) return HttpResponse("批量操作已启动", status=200)
3. 关键注意事项
- Django环境初始化:子进程必须调用
django.setup(),否则无法使用Django的ORM或其他功能。 - 避免子进程操作ORM:尽量在主进程完成所有数据库查询,将结果传递给子进程,减少数据库连接冲突。
- 异常捕获:子进程逻辑必须加异常处理,防止单个任务失败导致整个进程池崩溃。
- 进程数控制:
Pool(processes)不要设置过大,建议等于CPU核心数,避免占用过多系统资源。 - 请求超时问题:如果并行任务耗时较长,Django请求可能超时。若出现这种情况,建议改用Celery等异步任务队列,但多进程方案对新手更友好直接。
4. 替代方案:使用concurrent.futures(更简洁)
如果觉得multiprocessing繁琐,可以用ProcessPoolExecutor,语法更直观:
from concurrent.futures import ProcessPoolExecutor def wallet_verify(request): if request.method != 'POST': return HttpResponse("仅支持POST请求", status=400) # 收集API密钥对 active_user_ids = botactive.objects.filter(active=True).values_list('user', flat=True) api_pairs = [] for user_id in active_user_ids: api_datas = Bybitapidatas.objects.filter(user_id=user_id).values('apikey', 'apisecret') api_pairs.extend( (data['apikey'], data['apisecret']) for data in api_datas ) # 并行执行 with ProcessPoolExecutor(max_workers=4) as executor: executor.map(process_user_operation, [p[0] for p in api_pairs], [p[1] for p in api_pairs]) return HttpResponse("批量操作已启动", status=200)
内容的提问来源于stack exchange,提问作者user18377368
相关产品推荐
相关产品推荐

