You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Ruby(非Rails)中消费HTTP Server Sent Events?

Ruby非Rails环境下消费SSE的可行方案

我完全懂你找靠谱SSE客户端的糟心经历——em-eventsource毛病多,sse-client-ruby又太简陋没给够配置说明。下面给你几个实用的解决方案,不管是用成熟gem还是自己写轻量实现都能搞定:

一、推荐可靠的Ruby SSE Gem

1. ruby-sse

这个gem维护状态不错,API设计直观,支持SSE核心功能(事件监听、自动重连、自定义请求头),基本能覆盖大部分场景。

先安装:

gem install ruby-sse

简单使用示例:

require 'sse/client'

SSE::Client.new('https://your-sse-endpoint.com/events') do |client|
  # 监听所有消息
  client.on_event do |event|
    puts "Received event #{event.type}: #{event.data}"
  end

  # 监听特定类型的事件
  client.on_event('user_update') do |event|
    puts "User updated: #{event.data}"
  end

  # 处理连接断开
  client.on_error do |error|
    puts "Connection error: #{error.message}"
  end
end

2. faraday-sse

如果你项目里已经在用Faraday做HTTP请求,这个插件能无缝集成SSE支持,不用额外学习新的API。

安装:

gem install faraday-sse

使用示例:

require 'faraday'
require 'faraday/sse'

conn = Faraday.new(url: 'https://your-sse-endpoint.com') do |f|
  f.use Faraday::Response::SSE
end

conn.get('/events') do |req|
  req.headers['Accept'] = 'text/event-stream'
end do |stream|
  stream.on_event do |event|
    puts "[#{event.type}] #{event.data}"
  end
end

二、自制轻量SSE客户端(不依赖第三方Gem)

如果不想引入额外依赖,用Ruby标准库就能实现一个基础的SSE消费者,核心是解析SSE的文本格式。

下面是一个完整的示例,包含连接、消息解析、自动重连逻辑:

require 'net/http'
require 'uri'

def consume_sse(sse_url)
  uri = URI.parse(sse_url)
  http = Net::HTTP.new(uri.host, uri.port)
  http.use_ssl = uri.scheme == 'https'
  http.open_timeout = 10
  http.read_timeout = 0 # 长连接禁用超时

  # 初始化当前事件的临时存储
  current_event = {
    data: [],
    event_type: 'message',
    event_id: nil,
    retry_delay: 3000 # 默认重试间隔3秒
  }

  loop do
    begin
      request = Net::HTTP::Get.new(uri.request_uri)
      request['Accept'] = 'text/event-stream'
      request['Connection'] = 'keep-alive'
      request['Cache-Control'] = 'no-cache'
      # 如果有上次的事件ID,带上用于重连续传
      request['Last-Event-ID'] = current_event[:event_id] if current_event[:event_id]

      http.request(request) do |response|
        # 逐块读取响应流
        response.read_body do |chunk|
          chunk.split("\n").each do |line|
            next if line.strip.empty?

            # 解析SSE的各个字段
            case line
            when /^data:(.*)/
              current_event[:data] << $1.strip
            when /^event:(.*)/
              current_event[:event_type] = $1.strip
            when /^id:(.*)/
              current_event[:event_id] = $1.strip
            when /^retry:(.*)/
              current_event[:retry_delay] = $1.to_i
            end

            # 空行表示当前事件结束,处理并重置
            if line.empty?
              full_message = current_event[:data].join("\n")
              puts "[#{current_event[:event_type]}] ID: #{current_event[:event_id]} | #{full_message}"

              # 重置临时存储,准备下一个事件
              current_event[:data] = []
              current_event[:event_type] = 'message'
            end
          end
        end
      end
    rescue Errno::ECONNRESET, Errno::EPIPE, Net::OpenTimeout
      # 连接断开,按指定间隔重试
      puts "Connection dropped, retrying in #{current_event[:retry_delay]/1000}s..."
      sleep(current_event[:retry_delay] / 1000.0)
      retry
    end
  end
end

# 替换成你的SSE端点URL
consume_sse('https://your-sse-server.com/stream')

这个自制实现支持:

  • 解析SSE标准的data、event、id、retry字段
  • 自动处理连接断开后的重试
  • 带上Last-Event-ID实现断点续传

内容的提问来源于stack exchange,提问作者Sod Almighty

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.27 10:00:22