第 9 章 · 04 生产者消费者集群


文档摘要

第 9 章 · 04 生产者消费者集群 本节摘要:当算力不均衡或想让一个信号源驱动多个跟单 bot 时,freqtrade 的生产者-消费者(producer-consumer)模式派上用场。生产者 bot 正常算指标、出信号,通过 WebSocket 消息总线把「已分析的数据帧(analyzeddf)」和「白名单」广播出去;消费者 bot 订阅这些消息,直接拿到算好的指标,不必自己重算。本节讲清这套机制:生产者侧只需正常配 ,消费者侧配 块连上去;策略里用 拿数据;还会讲一个经典用法——把 FreqAI 跑在强机器上,让树莓派等弱设备当消费者只做跟单,以及延迟、一致性、信号传递的注意事项。 内容来源:原项目文档 、 ,汉化并套用体系化模板。

第 9 章 · 04 生产者消费者集群

本节摘要:当算力不均衡或想让一个信号源驱动多个跟单 bot 时,freqtrade 的生产者-消费者(producer-consumer)模式派上用场。生产者 bot 正常算指标、出信号,通过 WebSocket 消息总线把「已分析的数据帧(analyzed_df)」和「白名单」广播出去;消费者 bot 订阅这些消息,直接拿到算好的指标,不必自己重算。本节讲清这套机制:生产者侧只需正常配 api_server,消费者侧配 external_message_consumer 块连上去;策略里用 self.dp.get_producer_df(pair) 拿数据;还会讲一个经典用法——把 FreqAI 跑在强机器上,让树莓派等弱设备当消费者只做跟单,以及延迟、一致性、信号传递的注意事项。

内容来源:原项目文档 docs/producer-consumer.mddocs/rest-api.md,汉化并套用体系化模板。

⚠️ 风险提示:消费者依赖生产者的数据和信号,生产者宕机或信号错误会直接传导到所有消费者。务必监控生产者健康,消费者侧也要有信号缺失时的降级逻辑。

学习目标

阅读完本节,你应当能够:

  1. 说清 producerconsumer 的角色分工。
  2. 配置消费者的 external_message_consumer 块。
  3. 在策略里用 get_producer_df() 获取生产者数据。
  4. 理解 remove_entry_exit_signals 的作用。
  5. 部署 FreqAI + 弱机跟单 的经典架构。

一、为什么要 producer-consumer:算力分离与一信号多跟单

传统模式下每个 bot 独立算指标、独立出信号。问题有两个:

  • 算力浪费:多个 bot 跑同一套策略,重复算同样的指标。
  • 算力不均:FreqAI、复杂策略需要强算力,但你想在多台弱设备上跟单。

producer-consumer 模式解决这两点:

核心思想:生产者算好指标和信号,广播给消费者;消费者拿到「已分析的数据帧」,可以直接用,也可以用自己的逻辑重新解读。每个消费者是独立的 bot,有独立的账户、配置、交易,只是共用生产者的计算结果。

💡 典型场景:把 FreqAI 训练放在有 GPU 的强机器上当生产者,让树莓派等弱设备当消费者,只解读信号做跟单——弱设备根本跑不动 FreqAI,但跟单毫无压力。

二、生产者侧:配好 api_server 即可

生产者不需要特殊配置,只要正常开启 REST API(下一章详讲),因为它依赖的 WebSocket 端点在 api_server 里:

"api_server": { "enabled": true, "listen_ip_address": "127.0.0.1", "listen_port": 8080, "ws_token": "sercet_Ws_t0ken" }

生产者的策略照常写——正常在 populate_indicators 里算指标、在 populate_entry_trend 里出信号。这些算好的数据帧会自动通过 WebSocket 广播给订阅的消费者。

⚠️ ws_token 要随机且保密。任何拿到这个 token 的人都能订阅你的生产者数据(包括信号)。用 secrets.token_urlsafe(25) 生成。

三、消费者侧:external_message_consumer 块

消费者在配置里加 external_message_consumer 块,连接到生产者:

{ "external_message_consumer": { "enabled": true, "producers": [ { "name": "default", "host": "127.0.0.1", "port": 8080, "secure": false, "ws_token": "sercet_Ws_t0ken" } ], "wait_timeout": 300, "ping_timeout": 10, "sleep_time": 10, "remove_entry_exit_signals": false, "initial_candle_limit": 1500, "message_size_limit": 8 } }

关键参数:

参数 含义 默认
enabled 开启消费者模式 false
producers 生产者列表(可多个) 必填
producers.name 生产者名称(策略里引用用) 必填
producers.host / port 生产者地址 必填
producers.ws_token 与生产者的 ws_token 一致 必填
remove_entry_exit_signals 收到数据帧时清零信号列 false
initial_candle_limit 初始期望的 K 线数 1500
message_size_limit 单条消息大小上限(MB) 8

💡 多个生产者:producers 是数组,可以连多个生产者,策略里用 self.dp.get_producer_df(pair, producer_name="xxx") 指定取哪个。

四、消费者策略:get_producer_df 取数据

消费者策略不再自己算指标,而是从生产者拿:

class ConsumerStrategy(IStrategy): process_only_new_candles = False # 消费者必须设这个 _columns_to_expect = ['rsi_default', 'tema_default', 'bb_middleband_default'] def populate_indicators(self, dataframe, metadata): pair = metadata['pair'] producer_pairs = self.dp.get_producer_pairs() producer_dataframe, _ = self.dp.get_producer_df(pair) if not producer_dataframe.empty: # 把生产者的列合并进来,加后缀 "_default" merged_dataframe = merge_informative_pair( dataframe, producer_dataframe, self.timeframe, self.timeframe, append_timeframe=False, suffix="default") return merged_dataframe else: # 生产者没数据时降级 dataframe[self._columns_to_expect] = 0 return dataframe def populate_entry_trend(self, df, metadata): # 用带后缀的列(来自生产者)做信号 df.loc[ (qtpylib.crossed_above(df['rsi_default'], self.buy_rsi.value)) & (df['tema_default'] <= df['bb_middleband_default']) & (df['volume'] > 0), 'enter_long'] = 1 return df

要点:

  • process_only_new_candles = False消费者必须设的,否则可能错过生产者推送的数据。
  • get_producer_df(pair) 返回生产者已分析的数据帧和最后分析时间。
  • 合并时给生产者的列加后缀(如 _default),避免与本地产物冲突。
  • 要处理生产者无数据的情况(刚启动、断线),给列填默认值降级。

💡 信号直传:如果生产者已经出了 enter_long 等信号,且 remove_entry_exit_signals=false(默认),消费者可以直接用这些信号。设 remove_entry_exit_signals=true 则会清零信号列,强迫消费者用自己的逻辑重新判断——适合「只借指标不借信号」的场景。

五、延迟、一致性与故障处理

producer-consumer 是异步消息系统,有几个工程要点:

延迟:生产者算完到消费者收到有网络和序列化延迟。对分钟级以上的策略通常无感,但对秒级高频可能有问题。消费者侧 wait_timeout(默认 300 秒)控制多久没消息就重新 ping。

一致性:生产者和消费者的时间周期、白名单应保持一致,否则合并数据帧会错位。如果生产者白名单变了,消费者通过白名单消息自动同步(也可用 ProducerPairList 直接复用生产者白名单)。

故障处理:

  • 生产者宕机:消费者收不到消息,get_producer_df 返回空,策略要降级(填默认值或暂停)。
  • 网络断开:sleep_time(默认 10 秒)后自动重连。
  • 消息过大:message_size_limit(默认 8MB)会丢弃超限消息。

⚠️ 安全提醒:生产者和消费者之间的通信默认不走 HTTPS(secure: false)。跨公网部署时务必设 secure: true 或走 VPN/SSH 隧道,避免信号和指标被窃听。

六、ProducerPairList 与实战架构

消费者想直接复用生产者的白名单,用 ProducerPairList:

"pairlists": [ { "method": "ProducerPairList", "number_assets": 5, "producer_name": "default" } ]

它从指定生产者拉白名单,并校验这些交易对在当前交易所是否有效。

经典实战架构:FreqAI + 弱机集群

生产者用 GPU 跑 FreqAI,把预测值(&-s_close)广播出去。多个树莓派消费者各自解读——一个只跟多头信号、一个只跟空头信号、一个用不同阈值——形成「一信号多策略」的集群,而弱设备完全不用跑模型训练。

本节要点回顾

  1. 角色分工:Producer 算指标出信号并广播;Consumer 订阅拿已分析数据帧,独立交易。
  2. 生产者配置:正常开 api_server + ws_token 即可,策略照常写。
  3. 消费者配置:external_message_consumer 块连生产者;策略必须 process_only_new_candles = False,用 get_producer_df(pair) 取数据并合并(加后缀)。
  4. 信号传递:remove_entry_exit_signals=false(默认)可直传信号;true 则只借指标。
  5. 故障与安全:要写降级逻辑应对生产者宕机;跨公网用 secure: true 或 VPN。
  6. 经典架构:FreqAI 跑在 GPU 强机上,树莓派当消费者解读信号跟单;ProducerPairList 复用白名单。

至此第 9 章结束。下一章进入运维与扩展——Telegram 控制、FreqUI 与 REST API、插件系统、SQL 速查、策略迁移与 FAQ。


发布者: 作者: 灏天文库 转发
评论区 (0)
U