🚀 AI 一键生成 ptrade 策略代码
立即体验

ptrade 平台如何设置多个 run_interval 并行线程并实现安全的数据同步?

ptrade | 发布时间: 2026/8/10
以下内容由 EasyQuant 生成。

问题描述

ptrade 高级轮询:利用多跑 run_interval 的并行线程与全局变量数据同步
ptrade 怎么在 initialize 里设置多个 run_interval,并保证它们之间数据不冲突?

解决方案

在 ptrade 量化交易平台中,run_interval 接口允许用户以设定的时间间隔(单位为秒)周期性地运行指定的自定义函数。当在 initialize 中多次调用 run_interval 设定不同的任务时,系统会以多个线程并行运行这些函数。

由于多线程并行运行时会共享全局变量 g,如果多个线程同时对同一个全局变量进行读写操作,可能会导致数据竞争(Data Race)或冲突。以下是实现多线程安全轮询的策略与完整代码示例。

一、 核心设计原则

  1. 任务职责分离(Single Responsibility)
    尽量让不同的 run_interval 线程负责不同的标的或不同的逻辑。例如,线程 A 负责获取行情快照,线程 B 负责执行下单交易,避免它们同时修改同一个状态变量。
  2. 私有/独立变量隔离
    在全局变量 g 中,为不同的线程分配独立的子字典或子变量。例如 g.thread_A_datag.thread_B_data,从物理上隔离数据写入。
  3. 使用线程安全的数据结构
    Python 的内置 queue.Queue 是线程安全的,非常适合用于多线程之间传递交易信号或订单任务。

二、 完整策略源码示例

以下策略展示了如何在 initialize 中注册两个不同的 run_interval 任务:

  • 任务 1 (query_market_task):每 3 秒轮询一次行情快照,并将数据写入独立的全局变量中。
  • 任务 2 (trade_execution_task):每 5 秒轮询一次,读取行情数据并根据逻辑执行下单,避免了读写冲突。
import queue

def initialize(context):
    # 1. 设置操作的股票池
    g.security = '600570.SS'  # 恒生电子
    set_universe(g.security)
    
    # 2. 初始化线程隔离的全局数据结构
    g.market_data = {}       # 专门用于存放行情快照的字典
    g.signal_queue = queue.Queue()  # 线程安全的信号队列,用于传递交易指令
    g.trade_executed = False # 交易状态标识
    
    # 3. 注册多个 run_interval 任务(并行运行)
    # 任务 A:每 3 秒获取一次行情快照(高频查询)
    run_interval(context, query_market_task, seconds=3)
    
    # 任务 B:每 5 秒检查一次信号并执行交易(低频执行)
    run_interval(context, trade_execution_task, seconds=5)

def query_market_task(context):
    """
    任务 A:独立负责行情查询,只写 g.market_data,不参与交易
    """
    try:
        snapshot = get_snapshot(g.security)
        if snapshot and g.security in snapshot:
            # 安全写入独立变量
            g.market_data[g.security] = snapshot[g.security]
            log.info(f"[线程A] 成功更新 {g.security} 行情快照,最新价: {snapshot[g.security].get('last_px')}")
            
            # 简单策略逻辑:若最新价高于昨收价,向队列推送买入信号
            last_px = snapshot[g.security].get('last_px', 0)
            preclose_px = snapshot[g.security].get('preclose_px', 0)
            if last_px > preclose_px and not g.trade_executed:
                g.signal_queue.put(('BUY', last_px))
    except Exception as e:
        log.error(f"[线程A] 查询行情异常: {str(e)}")

def trade_execution_task(context):
    """
    任务 B:独立负责交易执行,只读 g.market_data,通过线程安全队列获取信号
    """
    # 检查是否有待处理的交易信号
    if not g.signal_queue.empty():
        try:
            signal_type, price = g.signal_queue.get_nowait()
            if signal_type == 'BUY' and not g.trade_executed:
                # 执行下单
                order_id = order(g.security, 100, limit_price=price)
                if order_id:
                    g.trade_executed = True
                    log.info(f"[线程B] 触发买入信号,以价格 {price} 委托买入 100 股,订单ID: {order_id}")
                g.signal_queue.task_done()
        except queue.Empty:
            pass
        except Exception as e:
            log.error(f"[线程B] 交易执行异常: {str(e)}")

def handle_data(context, data):
    # 基础结构必选函数,多线程策略中主要逻辑已在 run_interval 中实现
    pass

三、 避坑与注意事项

  1. 最小时间间隔限制
    run_intervalseconds 参数最小设定值为 3 秒。如果传入小于 3 的数值,ptrade 引擎会自动将其默认调整为 3 秒。
  2. 避免在 before_trading_start 中修改多线程共享变量
    由于服务器重启或隔日拉起时,before_trading_start 会被重新执行,请确保该函数中没有重置多线程正在依赖的动态变量,或者配合 set_parameters(server_restart_not_do_before="1") 限制其重复执行。
  3. 异常处理(Try-Except)
    在多线程并行环境中,任何一个线程因未捕获的异常崩溃都可能导致整个策略引擎终止。因此,在每个 run_interval 目标函数内部,务必使用 try...except 包裹核心逻辑。