如何在Elixir中使用Cowboy实现TradingView长连接保活API端点
基于Cowboy/Plug实现TradingView Streaming端点方案
1. 核心实现逻辑
TradingView的流端点要求长连接保活、多连接数据一致、chunked编码传输,普通Plug短生命周期请求无法满足需求,需要通过持有连接的常驻进程+全局广播机制实现:
- 用
send_chunked初始化chunked响应,禁用所有代理层缓冲 - 所有连接进程订阅同一个全局广播主题,保证所有连接收到的实时数据完全一致
- 内置心跳机制避免中间代理断开空闲连接
- 只处理订阅后的实时消息,不返回历史快照
2. 完整代码实现
路由配置
defmodule MyApp.Router do use MyApp.Web, :router # 其余路由规则省略 get "/streaming", MyApp.TradingViewStreamController, :index end
流处理控制器
defmodule MyApp.TradingViewStreamController do use MyApp.Web, :controller alias Phoenix.PubSub # 心跳间隔25秒,避免中间代理断开空闲连接 @heartbeat_interval 25_000 # 全局统一价格流主题,所有连接都订阅该主题保证数据一致 @stream_topic "trading:realtime_prices" def index(conn, _params) do conn = conn |> put_resp_header("content-type", "text/plain") |> put_resp_header("transfer-encoding", "chunked") |> put_resp_header("connection", "keep-alive") |> put_resp_header("cache-control", "no-cache, no-store, must-revalidate") # 禁用Nginx等反向代理的缓冲,保证chunked实时传输 |> put_resp_header("x-accel-buffering", "no") # 初始化chunked响应,状态码200 |> send_chunked(200) # 仅订阅实时流,不请求历史快照,符合接口要求 PubSub.subscribe(MyApp.PubSub, @stream_topic) # 启动心跳定时器 Process.send_after(self(), :send_heartbeat, @heartbeat_interval) # 进入消息循环,持有连接直到断开 stream_loop(conn) halt(conn) end defp stream_loop(conn) do receive do # 处理实时价格消息 {:new_price, data} -> # 所有消息末尾必须加换行符,符合TradingView格式要求 case chunk(conn, "#{Jason.encode!(data)}\n") do {:ok, conn} -> stream_loop(conn) # 连接断开时取消订阅,退出循环 {:error, _reason} -> PubSub.unsubscribe(MyApp.PubSub, @stream_topic) end # 发送心跳包,避免连接被代理掐断 :send_heartbeat -> # 心跳用注释格式,不影响业务数据解析 case chunk(conn, "#\n") do {:ok, conn} -> Process.send_after(self(), :send_heartbeat, @heartbeat_interval) stream_loop(conn) {:error, _reason} -> PubSub.unsubscribe(MyApp.PubSub, @stream_topic) end end end end
3. 上游数据推送实现
价格生成服务只需要向统一主题广播消息即可,所有订阅的连接会收到完全一致的实时数据:
# 上游价格更新时调用,向所有流连接推送同一份数据 PubSub.broadcast(MyApp.PubSub, @stream_topic, {:new_price, %{ symbol: "BTC_USDT", price: "42000.12", timestamp: System.system_time(:millisecond) }})
4. 反向代理配置(如有Nginx)
需要在Nginx配置中添加以下规则,避免缓冲chunked数据:
location /streaming { proxy_pass http://你的Elixir服务地址; proxy_http_version 1.1; proxy_buffering off; proxy_cache off; proxy_set_header Connection ""; }
内容的提问来源于stack exchange,提问作者give_me_ur_btcns
相关产品推荐
相关产品推荐

