如何在Rails中实现从外部服务器流式读取的常驻后台进程
下面是几个适配Rails架构的可行思路,无需直接迁移到Elixir:
1. 基于Rails环境的自定义常驻进程
直接编写Ruby脚本并加载Rails环境,这样可以直接操作模型,无需通过API中转。这种方式比Rake task更灵活规范,且能复用Rails的所有资源。
示例脚本(lib/streaming_price_worker.rb):
require_relative '../config/environment' class StreamingPriceWorker def run # 初始化券商流客户端 broker_client = BrokerStreamAPI.new(api_key: ENV['BROKER_API_KEY']) # 监听流数据并处理 broker_client.subscribe do |price_data| # 直接操作Rails模型更新数据 StockPrice.find_or_initialize_by(symbol: price_data[:symbol]) do |stock| stock.current_price = price_data[:price] stock.last_updated = Time.now stock.save! end end # 维持连接的阻塞调用(根据券商API实现调整) broker_client.keep_alive end end StreamingPriceWorker.new.run
启动方式可以用nohup bundle exec ruby lib/streaming_price_worker.rb &,或者通过systemd配置成系统服务,确保进程崩溃后自动重启。
2. 利用Sidekiq实现常驻任务
Sidekiq虽主打单次异步任务,但可以通过阻塞调用实现常驻流监听,同时利用Sidekiq自带的进程监控、重启机制。
示例Worker(app/workers/streaming_worker.rb):
class StreamingWorker include Sidekiq::Worker sidekiq_options queue: :streaming, retry: false, unique: true def perform broker_client = BrokerStreamAPI.new(api_key: ENV['BROKER_API_KEY']) broker_client.subscribe do |price_data| # 处理数据逻辑 StockPrice.update_price(price_data[:symbol], price_data[:price]) end broker_client.keep_alive # 阻塞维持连接 end end
启动时指定仅运行streaming队列,且通过Sidekiq的并发配置或唯一任务插件,确保只启动一个实例,避免重复连接券商流。
3. 独立流进程+Rails API推送
如果更倾向于解耦,可选用Python/Go等对长连接更友好的语言编写流处理进程,将清洗后的价格数据通过Rails的REST API(如POST /api/stock_prices)推送至Rails应用。
这种方案的优势是流处理与Rails业务完全分离,后续优化流处理逻辑无需改动Rails代码;缺点是需要维护两个独立服务,增加部署复杂度。
4. 结合Action Cable实现前端实时推送
如果需要将实时价格同步到前端,可在后台进程处理数据后,通过Rails的Action Cable广播数据:
ActionCable.server.broadcast "stock_updates:#{price_data[:symbol]}", price: price_data[:price], updated_at: Time.now.iso8601
前端通过WebSocket订阅对应频道,即可实现实时价格展示。
内容的提问来源于stack exchange,提问作者John Small

