多线程股票交易程序下单频率控制问题求助
解决方案:控制多线程下单速率以适配Broker限制
针对你的问题,核心是要同时控制并发下单线程数和每秒下单总量,以下提供两种可行的实现方案,解决你用Semaphore未成功的问题:
方案一:Semaphore+时间窗口控制
这种方式既限制同时下单的线程数,又确保每秒不超过3笔订单:
import threading import time # 获取broker对象 kite = get_kite() # 限制同时最多3个线程执行下单逻辑 order_semaphore = threading.Semaphore(3) # 保护订单时间记录的锁,避免多线程篡改数据 time_lock = threading.Lock() # 存储最近1秒内的订单时间戳 recent_order_times = [] def place_order(): with order_semaphore: # 控制每秒订单量不超过3笔 with time_lock: current_time = time.time() # 清理1秒前的旧订单记录 recent_order_times[:] = [t for t in recent_order_times if current_time - t < 1] # 如果已达每秒上限,等待到时间窗口刷新 while len(recent_order_times) >= 3: wait_duration = 1 - (current_time - recent_order_times[0]) time.sleep(wait_duration) current_time = time.time() recent_order_times[:] = [t for t in recent_order_times if current_time - t < 1] # 记录当前订单时间 recent_order_times.append(current_time) # 执行实际下单操作 try: # 替换为你的kite下单逻辑,例如: # kite.place_order(transaction_type=kite.TRANSACTION_TYPE_BUY, ...) print("订单提交成功") except Exception as e: print(f"订单提交失败: {str(e)}") def entry(): while True: # 替换为你的实际条件扫描逻辑 condition = True if condition: place_order() # 添加扫描间隔,避免无意义的CPU占用 time.sleep(0.1) # 启动10个扫描线程 for i in range(10): thread = threading.Thread(target=entry) thread.start()
关键说明
- 全局共享的
order_semaphore确保同时最多3个线程进入下单流程 recent_order_times结合锁机制,严格控制1秒内的订单总量不超过3笔- 之前你用Semaphore失败大概率是因为没有结合时间窗口控制,或Semaphore不是全局共享对象
方案二:队列+单线程速率控制
这种方式将所有下单请求放入队列,由单独线程按固定速率处理,逻辑更简单稳定:
import threading import time from queue import Queue kite = get_kite() # 下单请求队列 order_queue = Queue() def process_orders(): while True: # 每秒处理3笔订单,每笔间隔约0.33秒 for _ in range(3): if not order_queue.empty(): # 取出并执行下单逻辑 order_task = order_queue.get() try: order_task() print("订单提交成功") except Exception as e: print(f"订单提交失败: {str(e)}") order_queue.task_done() time.sleep(1/3) # 启动订单处理线程(后台守护线程) threading.Thread(target=process_orders, daemon=True).start() def place_order(): # 封装实际下单逻辑,放入队列 def actual_order(): # 替换为你的kite下单代码 # kite.place_order(...) pass order_queue.put(actual_order) def entry(): while True: # 替换为你的条件扫描逻辑 condition = True if condition: place_order() time.sleep(0.1) # 启动10个扫描线程 for i in range(10): thread = threading.Thread(target=entry) thread.start()
关键说明
- 所有扫描线程仅负责将下单请求放入队列,不直接操作Broker
- 单独的处理线程按固定速率消费队列,完全避免并发下单冲突
- 适合对延迟容忍度较低,且希望严格控制下单节奏的场景
注意事项
- 务必确保Semaphore、队列等控制对象是全局共享的,不要在每个线程内创建新实例
- 实际下单逻辑中要增加异常捕获,处理Broker返回的错误(如订单被拒、网络异常等)
- 调整
entry函数中的time.sleep(0.1)参数,匹配你的条件扫描频率,避免过度消耗CPU
内容的提问来源于stack exchange,提问作者Aditya Wagh
相关产品推荐
相关产品推荐

