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

DelayedJob中SQSManager本地队列消息丢失问题求助

问题描述

基于Ruby on Rails 4.2.11.36 LTS(Ruby 3.3)开发的Web项目,部署于Puma服务器,同时在多服务器运行DelayedJob 4.1.11 Worker。需求是将Puma处理的前端请求日志、Delayed Job运行日志异步发送至AWS SQS队列。

实现了SQSManager类:用线程安全的本地Queue暂存日志,单独启动线程每10秒批量发送消息。Puma中通过after_fork钩子启动线程后运行正常,但Delayed Job中出现异常——Job线程写入队列后,发送线程检测队列为空,日志显示消息写入后“消失”。

相关代码

SQSManager 实现

require 'aws-sdk-sqs'

class SQSManager
  SEND_INTERVAL = 10 # seconds
  MAX_BATCH_SIZE = 10000

  @@sqs_queue_url = APP_CONFIG[:aws_sqs_queue_url]
  @@local_queue = Queue.new # Threadsafe intermediate buffer
  @@sqs_client = Aws::SQS::Client.new(region: ENV['AWS_REGION']) if ENV['AWS_REGION']
  @@worker_thread = nil

  # Add a log to the local queue
  def self.add_log(log)
      Rails.logger.warn "SQSManager (#{$$}): adding log, current # elements #{@@local_queue}"
      @@local_queue << cal # 存在笔误:应为log而非cal
    end
  end

  def self.start
    @@worker_thread = Thread.new do
      Rails.logger.warn "SQSManager (#{$$}): Starting thread"
      loop do
        sleep(SEND_INTERVAL)
        send_messages
        if @@local_queue.empty?
          Rails.logger.warn "SQSManager (#{$$}): Stopping thread"
          break
        end
      end
    end
  end

  private

  def self.send_messages
    Rails.logger.warn "SQSManager (#{$$}): processing local queues"
    if @@local_queue.empty?
      Rails.logger.warn "SQSManager (#{$$}): Messages empty"
    else
      Rails.logger.warn "SQSManager (#{$$}): Sending messages"
      MAX_BATCH_SIZE.times do
        unless @@local_queue.empty?
          begin
            @@sqs_client.send_message({queue_url: sqs_queue_url, message_body: @@local_queue.pop})
          rescue Aws::SQS::Errors::ServiceError => e
            Rails.logger.error "SQSManager (#{$$}): Error sending messages to #{sqs_queue_url}: #{e.message}.\nBacktrace:\n\t#{e.backtrace.join("\n\t")}"
          end
        else
          break
        end
      end
    end
  end
end

Delayed Job 初始化代码

#config/initializer/delayed_job_worker.rb
require 'delayed_job'

module DelayedJobWorker

  def self.included(base)

    base.class_eval do
      alias_method :original_stop, :stop
      alias_method :original_start, :start

      def start
        Delayed::Worker.logger.warn "original_start worker!"
        SQSManager.start # SQSManager initialization before the Worker original start
        original_start
      end

      def stop
        SQSManager.shutdown
        original_stop
      end
    end
  end
end

Delayed::Worker.include(DelayedJobWorker)

异常日志

SQSManager (15928): Starting thread
SQSManager (15928): adding log, current # elements 0
SQSManager (15928): adding log, current # elements 1
SQSManager (15928): processing local queues
SQSManager (15928): Messages empty
SQSManager (15928): processing local queues
SQSManager (15928): Messages empty
SQSManager (15928): adding log, current # elements 0
SQSManager (15928): processing local queues
SQSManager (15928): Messages empty

问题根源分析
  1. 进程隔离导致类变量失效:Delayed Job Worker默认以多进程模式运行,父进程启动后会fork子进程处理Job。Ruby的类变量@@在fork后,子进程会复制父进程内存空间,但后续修改完全独立。原代码在父进程的start方法中启动SQS线程,子进程写入的是自己的@@local_queue,而父进程的发送线程访问的是父进程的空队列,导致消息“消失”。
  2. 代码笔误:add_log方法中错误地将log写成cal,导致写入无效对象(日志中计数增加是因为写入了cal变量,而非实际日志内容)。
  3. 线程提前退出逻辑:发送线程在send_messages后检测队列空就直接退出,后续写入的消息无人处理。

解决方案

1. 修正Delayed Job初始化逻辑

将SQSManager.start移至Worker的after_fork钩子,确保每个子进程都启动独立的SQS线程:

#config/initializer/delayed_job_worker.rb
require 'delayed_job'

module DelayedJobWorker
  def self.included(base)
    base.class_eval do
      # 用after_fork钩子替代原start方法,确保子进程初始化SQSManager
      after_fork do |worker|
        Delayed::Worker.logger.warn "Worker #{worker.id} forked, starting SQSManager"
        SQSManager.start
      end

      before_stop do |worker|
        Delayed::Worker.logger.warn "Worker #{worker.id} stopping, shutting down SQSManager"
        SQSManager.shutdown
      end
    end
  end
end

Delayed::Worker.include(DelayedJobWorker)

2. 修复SQSManager核心逻辑

修正笔误、调整线程生命周期控制、增加失败重试机制:

require 'aws-sdk-sqs'

class SQSManager
  SEND_INTERVAL = 10 # seconds
  MAX_BATCH_SIZE = 10000

  @@sqs_queue_url = APP_CONFIG[:aws_sqs_queue_url]
  @@local_queue = Queue.new # Threadsafe intermediate buffer
  @@sqs_client = Aws::SQS::Client.new(region: ENV['AWS_REGION']) if ENV['AWS_REGION']
  @@worker_thread = nil
  @@running = false

  # Add a log to the local queue
  def self.add_log(log)
    Rails.logger.warn "SQSManager (#{$$}): adding log, current # elements #{@@local_queue.size}"
    @@local_queue << log
  rescue => e
    Rails.logger.error "SQSManager (#{$$}): Failed to add log to local queue: #{e.message}"
  end

  def self.start
    return if @@running || @@worker_thread&.alive?

    @@running = true
    @@worker_thread = Thread.new do
      Rails.logger.warn "SQSManager (#{$$}): Starting thread"
      while @@running
        sleep(SEND_INTERVAL)
        send_messages
      end
      # 退出前处理剩余所有消息
      send_messages until @@local_queue.empty?
      Rails.logger.warn "SQSManager (#{$$}): Thread stopped"
    end
  end

  def self.shutdown
    @@running = false
    @@worker_thread&.join if @@worker_thread&.alive?
  end

  private

  def self.send_messages
    Rails.logger.warn "SQSManager (#{$$}): processing local queues, #{@@local_queue.size} messages pending"
    return if @@local_queue.empty? || !@@sqs_client

    messages_sent = 0
    MAX_BATCH_SIZE.times do
      break if @@local_queue.empty?

      begin
        message = @@local_queue.pop(true) # 非阻塞弹出,避免线程挂起
        @@sqs_client.send_message({queue_url: @@sqs_queue_url, message_body: message})
        messages_sent += 1
      rescue ThreadError # 队列空时抛出
        break
      rescue Aws::SQS::Errors::ServiceError => e
        Rails.logger.error "SQSManager (#{$$}): Error sending message to #{@@sqs_queue_url}: #{e.message}.\nBacktrace:\n\t#{e.backtrace.join("\n\t")}"
        # 发送失败的消息重新放回队列,避免丢失
        @@local_queue << message
        break # 遇到错误暂停发送,避免批量失败
      end
    end
    Rails.logger.warn "SQSManager (#{$$}): Sent #{messages_sent} messages"
  end
end

3. 验证Worker运行模式

确保Delayed Job Worker未禁用fork机制,启动Worker时查看日志,确认每个子进程都输出了Worker X forked, starting SQSManager的日志,验证SQS线程在子进程中正常启动。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 08:55:00