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

