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

如何在Ruby on Rails中接收Kubernetes API的流式响应?

实现思路与示例代码

核心方案:用后台Worker处理长连接

Rails主线程不能被长连接阻塞,必须把K8s流式请求放到后台进程中。推荐用Sidekiq(或Rails内置Active Job配合Sidekiq/Delayed Job)实现Worker,这类异步任务框架天生适配长时间运行的任务场景。

步骤1:配置K8s API访问权限

确保Worker能正常访问K8s API:集群内部署时,直接用挂载的ServiceAccount token;本地测试可从kubeconfig文件提取认证信息。

步骤2:编写异步Worker

用Ruby的http.rb或faraday库处理流式请求,这两个库对响应流的支持更贴合Rails生态,比EventMachine更易维护。

示例:Sidekiq Worker + http.rb 处理流式请求

# app/workers/k8s_stream_worker.rb
require 'http'

class K8sStreamWorker
  include Sidekiq::Worker
  sidekiq_options queue: :k8s_stream, retry: false # 长时间任务关闭自动重试

  def perform(resource_type = "pods", namespace = "default")
    # 配置K8s API地址与认证信息
    k8s_api_url = "https://kubernetes.default.svc.cluster.local:443/api/v1/namespaces/#{namespace}/#{resource_type}?watch=true"
    token = File.read('/var/run/secrets/kubernetes.io/serviceaccount/token')
    ca_cert = '/var/run/secrets/kubernetes.io/serviceaccount/ca.crt'

    begin
      # 发起流式请求并逐块处理响应
      HTTP.get(k8s_api_url,
        headers: {
          'Authorization' => "Bearer #{token}",
          'Accept' => 'application/json'
        },
        ssl_ca_file: ca_cert
      ).each do |chunk|
        next if chunk.empty?
        # 解析K8s watch返回的单条事件(换行分隔的JSON格式)
        event = JSON.parse(chunk)
        puts "[K8s Stream] #{event['type']} - #{event['object']['metadata']['name']}"
      end
    rescue => e
      puts "[K8s Stream Error] #{e.message}"
      # 出错后重新入队,保证监听不中断
      K8sStreamWorker.perform_async(resource_type, namespace)
    end
  end
end

步骤3:控制器触发Worker

控制器仅负责触发Worker,不会阻塞主线程:

# app/controllers/k8s_controller.rb
class K8sController < ApplicationController
  def start_stream
    K8sStreamWorker.perform_async(params[:resource_type], params[:namespace])
    head :ok
  end
end

步骤4:Faraday替代方案

若习惯使用Faraday,可替换Worker内的请求逻辑:

# Worker内的请求部分替换为以下代码
require 'faraday'
require 'json'

conn = Faraday.new(url: k8s_api_url) do |f|
  f.headers['Authorization'] = "Bearer #{token}"
  f.headers['Accept'] = 'application/json'
  f.ssl.ca_file = ca_cert
  f.adapter :net_http # 选用支持流式响应的适配器
end

conn.get do |req|
  req.options.on_data = lambda do |chunk, overall_received_bytes|
    next if chunk.empty?
    event = JSON.parse(chunk)
    puts "[K8s Stream] #{event['type']} - #{event['object']['metadata']['name']}"
  end
end

关键注意事项

  • Worker队列隔离:给这类长连接Worker单独分配队列(如示例中的k8s_stream),避免阻塞其他短任务的执行。
  • 断连重连:K8s API可能主动断开连接,需在异常捕获逻辑中重新触发Worker,维持监听连续性。
  • 内存监控:长时间运行的Worker易出现内存泄漏,定期检查内存占用,必要时设置Worker自动重启机制。
  • 认证安全:集群内部署优先使用ServiceAccount,本地测试避免硬编码token,从kubeconfig动态读取。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 11:01:34