非同期処理とメッセージキュー — 処理を後回しにする力
「5秒待たせる」の問題
EchoTaskにレポート生成機能を追加した日、ハルトはすぐに後悔した。
機能リリースから30分後、Slackに通知が届いた。「レポートボタンを押しても何も起きない。壊れてるの?」「自分も同じ問題」「ずっとロードしてる」とメッセージが積み上がった。
# app/controllers/api/v1/reports_controller.rb
def create
report = ReportGenerator.new(current_user).generate
# ↑ これが5〜10秒かかる
render json: { download_url: report.url }
endハルトはすぐにAPMダッシュボードを開いた。レスポンスタイムが平均7.3秒。さらに悪いことに、複数のユーザーが同時にレポート生成を押すと、Pumaのスレッドがすべてふさがり、他のAPIも応答しなくなっていた。
Puma threads: 5本 / 同時レポート生成: 5本
→ 6人目のユーザーはすべてのAPIが「タイムアウト」
→ ログイン、タスク更新、コメント投稿……全部が止まる
ハルトは椅子にもたれて天井を見た。「重い処理をAPIのリクエスト内でやってはいけない。」そんなことは分かっていた。でも「とりあえず動かす」を優先した結果がこれだ。
「重い処理はAPIのリクエスト内でやってはいけない。」
非同期処理の考え方
同期処理(悪い): リクエスト → サーバーが5秒処理 → レスポンス(この間Pumaスレッドを1本占有)
非同期処理(良い):
- ユーザーがリクエスト
- サーバーがジョブをキューに積む(数ミリ秒)
- 即座にレスポンス「受け付けました」(202 Accepted)
- バックグラウンドのWorkerが5秒の処理を実行
- 完了したらメールで通知
INFO
HTTPステータス 202 Accepted の意味
「リクエストを受け付けたが、まだ処理は完了していない」を意味するステータスコード。非同期処理の開始時に使う。進捗確認用のURLやジョブIDを一緒に返すのがベストプラクティスだ。
Sidekiq — Redisを使ったバックグラウンドジョブ
SidekiqはRailsの非同期処理ライブラリの定番だ。RedisをバックエンドとしてJobをキューイングし、Workerプロセスが並列で処理する。
# Gemfile
gem 'sidekiq'
gem 'sidekiq-web' # 管理UI
gem 'sidekiq-cron' # スケジューリングコンカレンシーチューニング
Sidekiqの最も重要な設定が :concurrency(同時実行スレッド数)だ。ここを間違えると、データベース接続プールが枯渇してエラーが多発する。
# config/sidekiq.yml
:concurrency: 10 # Sidekiqのスレッド数
:queues:
- [critical, 3] # 優先度高(重み3)—— 支払い通知など
- [default, 2] # 通常(重み2)—— メール送信など
- [low, 1] # 低優先度(重み1)—— レポート生成など
:max_retries: 3
:timeout: 25 # Graceful shutdownのタイムアウト(秒)WARNING
スレッド数とDBコネクションプールの関係
Sidekiqのスレッド数(:concurrency)は必ずDBコネクションプール数以下に設定する。database.yml のデフォルト pool: 5 のままで concurrency を10にすると、5スレッドはDB取得待ちで詰まる。
鉄則: pool >= concurrency + 2(例: concurrency:10 → pool:12以上)
# config/initializers/sidekiq.rb
Sidekiq.configure_server do |config|
config.redis = { url: ENV['REDIS_URL'], size: 10 }
# DBコネクションプールをスレッド数に合わせて動的設定
Rails.application.config.after_initialize do
db_config = ActiveRecord::Base.configurations.find_db_config(Rails.env)
new_config = db_config.configuration_hash.merge(
pool: Sidekiq.options[:concurrency] + 2
)
ActiveRecord::Base.establish_connection(new_config)
end
endJobの実装とコントローラー
# app/jobs/report_generation_job.rb
class ReportGenerationJob < ApplicationJob
queue_as :low
retry_on StandardError, wait: :polynomially_longer, attempts: 3
def perform(user_id, report_type, options = {})
user = User.find(user_id)
ProgressTracker.update(job_id, 10, "データを収集中...")
report = ReportGenerator.new(user, report_type, options).generate
ProgressTracker.update(job_id, 100, "完了")
ReportMailer.completed(user, report).deliver_later
ActionCable.server.broadcast("user_#{user_id}_notifications", {
type: 'report_completed',
download_url: report.download_url
})
rescue => e
ProgressTracker.update(job_id, -1, "エラー: #{e.message}")
ReportMailer.failed(user, e.message).deliver_later
raise
end
end# コントローラー側(202 Acceptedで即返す)
class Api::V1::ReportsController < ApplicationController
def create
job = ReportGenerationJob.perform_later(
current_user.id, params[:type], params[:options]
)
render json: {
job_id: job.job_id,
status_url: api_v1_report_status_url(job.job_id),
message: "レポートの生成を開始しました。完了したらメールでお知らせします。"
}, status: :accepted
end
def status
render json: ProgressTracker.get(params[:job_id])
end
endジョブの進捗管理 — Redis HASHで状態保存
# app/services/progress_tracker.rb
class ProgressTracker
TTL = 24.hours
def self.update(job_id, percent, message)
key = "job_progress:#{job_id}"
Redis.current.hset(key,
"percent", percent.to_s,
"message", message,
"updated_at", Time.current.iso8601
)
Redis.current.expire(key, TTL.to_i)
end
def self.get(job_id)
data = Redis.current.hgetall("job_progress:#{job_id}")
return { status: "not_found" } if data.empty?
{
status: data["percent"].to_i == 100 ? "completed" :
data["percent"].to_i == -1 ? "failed" : "processing",
percent: data["percent"].to_i,
message: data["message"],
updated_at: data["updated_at"]
}
end
endジョブの優先度管理とスケジューリング(sidekiq-cron)
# config/schedule.yml
weekly_report:
cron: "0 9 * * 1" # 毎週月曜 9:00
class: "WeeklyReportJob"
queue: low
description: "週次レポートの生成・送信"
daily_cleanup:
cron: "0 3 * * *" # 毎日 3:00
class: "DailyCleanupJob"
queue: default
description: "期限切れセッション・一時ファイルの削除"# config/initializers/sidekiq.rb(スケジュール設定を追加)
Sidekiq.configure_server do |config|
schedule_file = "config/schedule.yml"
Sidekiq::Cron::Job.load_from_hash(YAML.load_file(schedule_file)) if File.exist?(schedule_file)
endバッチ処理 — 大量データをfind_in_batchesで分割
# 危険: 全データをメモリに載せる(50万件がメモリに!)
User.active.all.each { |u| ReportMailer.weekly(u).deliver_later }
# 安全: find_in_batches で分割処理
class WeeklyReportJob < ApplicationJob
queue_as :low
def perform
User.active.find_in_batches(batch_size: 1000) do |users|
users.each { |u| WeeklyReportMailJob.perform_later(u.id) }
sleep(0.1) # バッチ間でDB負荷を分散
end
end
endINFO
find_in_batches vs find_each
find_in_batches: バッチ単位(配列)でブロックに渡す。バッチ内を一括処理したいときfind_each: 1件ずつブロックに渡す。件ごとに独立した処理をするとき
どちらも LIMIT 1000 OFFSET n でDBから少しずつ取得するためメモリ消費が少ない。
SQS — AWSのマネージドキュー
Redisが落ちるとSidekiqのジョブが消える。AWSのSQSはマネージドサービスで耐障害性と可用性を保証する。
FIFOキューとスタンダードキューの選択基準
| 比較項目 | スタンダードキュー | FIFOキュー |
|---|---|---|
| スループット | ほぼ無制限 | 300 msg/秒(バッチは3000) |
| 順序保証 | なし(ベストエフォート) | 厳密な順序保証 |
| 重複配信 | あり(at-least-once) | なし(exactly-once) |
| 用途 | メール送信、画像変換など | 金融取引、注文処理など |
| コスト | 安い | やや高い |
Dead Letter Queue の設定とアラーム
# cloudformation/sqs.yml(抜粋)
DeadLetterQueue:
Type: AWS::SQS::Queue
Properties:
QueueName: echo-task-dlq
MessageRetentionPeriod: 1209600 # 14日間保持
DefaultQueue:
Type: AWS::SQS::Queue
Properties:
QueueName: echo-task-default
VisibilityTimeout: 300
ReceiveMessageWaitTimeSeconds: 20 # Long Polling
RedrivePolicy:
deadLetterTargetArn: !GetAtt DeadLetterQueue.Arn
maxReceiveCount: 3 # 3回失敗したらDLQへ
DLQAlarm:
Type: AWS::CloudWatch::Alarm
Properties:
AlarmName: echo-task-dlq-not-empty
MetricName: ApproximateNumberOfMessagesVisible
Namespace: AWS/SQS
Statistic: Sum
Period: 60
EvaluationPeriods: 1
Threshold: 1
ComparisonOperator: GreaterThanOrEqualToThreshold
Dimensions:
- Name: QueueName
Value: !GetAtt DeadLetterQueue.QueueName
AlarmActions:
- !Ref AlertSNSTopicWARNING
DLQは「調査待ちの問題一覧」
DLQにメッセージが入ったらCloudWatchアラームでSlackに通知する。DLQをただの「失敗の墓場」にしてはいけない。定期的にメッセージを確認し、バグ修正後はRedriveで再処理することで、1件も取りこぼさない運用を実現する。
Golang でのWorker実装 — aws-sdk-go-v2
EchoTaskのデータ分析基盤はGoで書かれている。SQSポーリングとメッセージ処理をGoで実装した。
// worker/sqs_worker.go
package worker
import (
"context"
"encoding/json"
"log/slog"
"time"
"github.com/aws/aws-sdk-go-v2/aws"
"github.com/aws/aws-sdk-go-v2/config"
"github.com/aws/aws-sdk-go-v2/service/sqs"
"github.com/aws/aws-sdk-go-v2/service/sqs/types"
)
type JobPayload struct {
UserID int64 `json:"user_id"`
ReportType string `json:"report_type"`
RequestedAt string `json:"requested_at"`
}
type SQSWorker struct {
client *sqs.Client
queueURL string
handler func(ctx context.Context, payload JobPayload) error
}
func NewSQSWorker(ctx context.Context, queueURL string, handler func(context.Context, JobPayload) error) (*SQSWorker, error) {
cfg, _ := config.LoadDefaultConfig(ctx, config.WithRegion("ap-northeast-1"))
return &SQSWorker{client: sqs.NewFromConfig(cfg), queueURL: queueURL, handler: handler}, nil
}
func (w *SQSWorker) Start(ctx context.Context, concurrency int) {
sem := make(chan struct{}, concurrency) // セマフォで並列数を制御
for {
select {
case <-ctx.Done():
slog.Info("Workerをシャットダウン中...")
return
default:
}
// Long Polling: メッセージが来るまで最大20秒待つ
output, err := w.client.ReceiveMessage(ctx, &sqs.ReceiveMessageInput{
QueueUrl: aws.String(w.queueURL),
MaxNumberOfMessages: 10,
WaitTimeSeconds: 20,
})
if err != nil {
slog.Error("SQSメッセージ受信エラー", "error", err)
time.Sleep(5 * time.Second)
continue
}
for _, msg := range output.Messages {
sem <- struct{}{}
go func(m types.Message) {
defer func() { <-sem }()
w.processMessage(ctx, m)
}(msg)
}
}
}
func (w *SQSWorker) processMessage(ctx context.Context, msg types.Message) {
var payload JobPayload
if err := json.Unmarshal([]byte(*msg.Body), &payload); err != nil {
slog.Error("JSONパースエラー", "error", err)
w.deleteMessage(ctx, msg) // パースエラーは修正不可 → DLQへ
return
}
if err := w.handler(ctx, payload); err != nil {
slog.Error("ジョブ処理エラー", "user_id", payload.UserID, "error", err)
return // 削除しない → VisibilityTimeout後に再試行 → maxReceiveCount超過でDLQ
}
w.deleteMessage(ctx, msg) // 成功したら削除
}
func (w *SQSWorker) deleteMessage(ctx context.Context, msg types.Message) {
w.client.DeleteMessage(ctx, &sqs.DeleteMessageInput{
QueueUrl: aws.String(w.queueURL),
ReceiptHandle: msg.ReceiptHandle,
})
}SNS + EventBridge — イベントルーティング
タスク完了時に「通知」「ログ」「課金」を直接呼び出しで繋ぐと密結合になる。EventBridgeでイベントの内容によってルーティング先を変える。
# cloudformation/eventbridge.yml(抜粋)
EchoTaskEventBus:
Type: AWS::Events::EventBus
Properties:
Name: echo-task-events
TaskCompletedToNotificationRule:
Type: AWS::Events::Rule
Properties:
EventBusName: !Ref EchoTaskEventBus
EventPattern:
source: ["echo-task.tasks"]
detail-type: ["TaskCompleted"]
Targets:
- Arn: !GetAtt NotificationQueue.Arn
Id: NotificationQueue
PlanChangedToBillingRule:
Type: AWS::Events::Rule
Properties:
EventBusName: !Ref EchoTaskEventBus
EventPattern:
source: ["echo-task.subscriptions"]
detail-type: ["PlanUpgraded", "PlanDowngraded", "SubscriptionCancelled"]
Targets:
- Arn: !GetAtt BillingQueue.Arn
Id: BillingQueue# Railsからのイベント発行
class Task < ApplicationRecord
after_update :publish_event, if: :saved_change_to_status?
private
def publish_event
return unless status == 'completed'
Aws::EventBridge::Client.new.put_events(
entries: [{
event_bus_name: ENV['EVENT_BUS_NAME'],
source: 'echo-task.tasks',
detail_type: 'TaskCompleted',
detail: { task_id: id, project_id: project_id, completed_at: updated_at.iso8601 }.to_json
}]
)
end
endINFO
SNS vs EventBridge の使い分け
- SNS: シンプルなfan-out(1対多通知)に最適。設定が簡単でコストが低い
- EventBridge: イベントの内容でルーティング先を変えたい場合に使う。外部SaaSとの連携やスケジューリングもできる
ジョブの冪等性
「同じジョブが2回実行されても問題ない」——これが冪等性だ。SQSはAt-least-once配信を保証するため、ネットワーク障害でジョブが重複実行されることがある。
# 危険: 冪等でないジョブ(2回実行で2通送信)
class SendWelcomeEmailJob < ApplicationJob
def perform(user_id)
UserMailer.welcome(User.find(user_id)).deliver_now
end
end
# 安全: DBフラグで冪等化
class SendWelcomeEmailJob < ApplicationJob
def perform(user_id)
user = User.find(user_id)
return if user.welcome_email_sent_at.present? # 既に送信済みならスキップ
UserMailer.welcome(user).deliver_now
user.update!(welcome_email_sent_at: Time.current)
end
end
# より堅牢: Redisのアトミック操作でロック
class ChargeSubscriptionJob < ApplicationJob
def perform(user_id, period)
key = "charge:#{user_id}:#{period}"
acquired = Redis.current.set(key, 1, nx: true, ex: 86400) # SET NX EX
return unless acquired # 別Workerが先に取得した場合はスキップ
ChargeService.new(user_id).charge_monthly
end
endWARNING
非同期ジョブは必ず冪等にする
SQSのAt-least-once配信、リトライ、デプロイ中の再起動——様々な理由でジョブは複数回実行される。特に課金・メール送信・外部API呼び出しは冪等性の確保を怠らないこと。
ECS上でWorkerを動かす
APIサーバーとWorkerは同じDockerイメージを使いながら、起動コマンドだけを変えてECSのサービスとして分離する。
# cloudformation/ecs-worker.yml(抜粋)
WorkerTaskDefinition:
Type: AWS::ECS::TaskDefinition
Properties:
Family: echo-task-worker
Cpu: '512'
Memory: '1024'
ContainerDefinitions:
- Name: sidekiq
Image: !Sub "${AWS::AccountId}.dkr.ecr.ap-northeast-1.amazonaws.com/echo-task:latest"
Command: ["bundle", "exec", "sidekiq", "-C", "config/sidekiq.yml"]
Environment:
- Name: REDIS_URL
Value: !Sub "redis://${ElastiCacheCluster.PrimaryEndPoint.Address}:6379"
Secrets:
- Name: DATABASE_URL
ValueFrom: !Sub "arn:aws:secretsmanager:ap-northeast-1:${AWS::AccountId}:secret:echo-task/db-url"
LogConfiguration:
LogDriver: awslogs
Options:
awslogs-group: "/ecs/echo-task-worker"
awslogs-region: ap-northeast-1
awslogs-stream-prefix: sidekiq
WorkerService:
Type: AWS::ECS::Service
Properties:
ServiceName: echo-task-worker
Cluster: !Ref ECSCluster
TaskDefinition: !Ref WorkerTaskDefinition
DesiredCount: 2 # 冗長化のため2台
LaunchType: FARGATEINFO
Workerのオートスケーリング
SQSのメッセージ数をCloudWatchカスタムメトリクスとして送り、ECS Application Auto Scalingのトリガーにできる。キューが溜まったら自動でWorkerを増やし、空になったら減らす——負荷に応じた自動運転でコストを最適化できる。
改善結果
3日間の実装を経て、EchoTaskの非同期化が完了した。ハルトはデプロイ後のダッシュボードを眺めた。
Before: 同期レポート生成
- APIレスポンスタイム: 5〜10秒
- 同時実行でPumaスレッドが枯渇 → 503エラー多発
After: Sidekiq + SQS非同期処理
- APIレスポンスタイム: 平均52ms(「受け付けました」)
- Workerが10並列でジョブを処理
- 3回失敗でDLQへ → CloudWatchアラームでSlack通知
- Redisで進捗管理 → フロントで進捗バー表示
Slack通知が届いた。「レポート生成が爆速になった!」「バックグラウンドで動いてる感じが分かって安心する」——ハルトは初めて「良いシステムを作った」という手応えを感じた。
EchoTaskは順調に成長し、ユーザー数が50万人を超えた。そこでハルトは根本的な疑問を持つようになった。
「機能が増えすぎて、コードが複雑になってきた。タスク管理、通知、課金、分析——これらを1つのRailsアプリに詰め込むのは限界かもしれない。」
次の章では、マイクロサービスという選択肢を学ぶ。
INFO
この章のキーポイント
- 重い処理はAPIレスポンス内でやらない → 202 Acceptedとジョブ IDを返す
- Sidekiqのスレッド数はDBコネクションプール以下に設定(pool >= concurrency + 2)
- sidekiq-cronでスケジューリング、find_in_batchesで大量データを分割処理
- Redis HASHで進捗管理し「今どこまで進んでいるか」をユーザーに伝える
- SQS スタンダードキューは高スループット・低コスト、FIFOキューは順序保証・重複排除
- DLQとCloudWatchアラームで「失敗の見える化」を怠らない
- EventBridgeでイベントをルーティングし、サービス間を疎結合に保つ
- Goでは aws-sdk-go-v2 + セマフォで並列数を制御し、Graceful Shutdownを実装する
- ジョブは必ず冪等に実装する(特に課金・メール送信・外部API呼び出し)