第 9 章 · 04 生产者消费者集群 本节摘要:当算力不均衡或想让一个信号源驱动多个跟单 bot 时,freqtrade 的生产者-消费者(producer-consumer)模式派上用场。生产者 bot 正常算指标、出信号,通过 WebSocket 消息总线把「已分析的数据帧(analyzeddf)」和「白名单」广播出去;消费者 bot 订阅这些消息,直接拿到算好的指标,不必自己重算。本节讲清这套机制:生产者侧只需正常配 ,消费者侧配 块连上去;策略里用 拿数据;还会讲一个经典用法——把 FreqAI 跑在强机器上,让树莓派等弱设备当消费者只做跟单,以及延迟、一致性、信号传递的注意事项。 内容来源:原项目文档 、 ,汉化并套用体系化模板。
本节摘要:当算力不均衡或想让一个信号源驱动多个跟单 bot 时,freqtrade 的生产者-消费者(producer-consumer)模式派上用场。生产者 bot 正常算指标、出信号,通过 WebSocket 消息总线把「已分析的数据帧(analyzed_df)」和「白名单」广播出去;消费者 bot 订阅这些消息,直接拿到算好的指标,不必自己重算。本节讲清这套机制:生产者侧只需正常配
api_server,消费者侧配external_message_consumer块连上去;策略里用self.dp.get_producer_df(pair)拿数据;还会讲一个经典用法——把 FreqAI 跑在强机器上,让树莓派等弱设备当消费者只做跟单,以及延迟、一致性、信号传递的注意事项。
内容来源:原项目文档
docs/producer-consumer.md、docs/rest-api.md,汉化并套用体系化模板。
⚠️ 风险提示:消费者依赖生产者的数据和信号,生产者宕机或信号错误会直接传导到所有消费者。务必监控生产者健康,消费者侧也要有信号缺失时的降级逻辑。
阅读完本节,你应当能够:
传统模式下每个 bot 独立算指标、独立出信号。问题有两个:
producer-consumer 模式解决这两点:
核心思想:生产者算好指标和信号,广播给消费者;消费者拿到「已分析的数据帧」,可以直接用,也可以用自己的逻辑重新解读。每个消费者是独立的 bot,有独立的账户、配置、交易,只是共用生产者的计算结果。
💡 典型场景:把 FreqAI 训练放在有 GPU 的强机器上当生产者,让树莓派等弱设备当消费者,只解读信号做跟单——弱设备根本跑不动 FreqAI,但跟单毫无压力。
生产者不需要特殊配置,只要正常开启 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": { "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")指定取哪个。
消费者策略不再自己算指标,而是从生产者拿:
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:
"pairlists": [ { "method": "ProducerPairList", "number_assets": 5, "producer_name": "default" } ]
它从指定生产者拉白名单,并校验这些交易对在当前交易所是否有效。
经典实战架构:FreqAI + 弱机集群
生产者用 GPU 跑 FreqAI,把预测值(&-s_close)广播出去。多个树莓派消费者各自解读——一个只跟多头信号、一个只跟空头信号、一个用不同阈值——形成「一信号多策略」的集群,而弱设备完全不用跑模型训练。
api_server + ws_token 即可,策略照常写。external_message_consumer 块连生产者;策略必须 process_only_new_candles = False,用 get_producer_df(pair) 取数据并合并(加后缀)。remove_entry_exit_signals=false(默认)可直传信号;true 则只借指标。secure: true 或 VPN。ProducerPairList 复用白名单。至此第 9 章结束。下一章进入运维与扩展——Telegram 控制、FreqUI 与 REST API、插件系统、SQL 速查、策略迁移与 FAQ。