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

多线程股票交易程序下单频率控制问题求助

解决方案:控制多线程下单速率以适配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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 20:52:00