mybook

非同期処理とメッセージキュー — 処理を後回しにする力

「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のリクエスト内でやってはいけない。」


非同期処理の考え方

Loading diagram...

同期処理(悪い): リクエスト → サーバーが5秒処理 → レスポンス(この間Pumaスレッドを1本占有)

非同期処理(良い):

  1. ユーザーがリクエスト
  2. サーバーがジョブをキューに積む(数ミリ秒)
  3. 即座にレスポンス「受け付けました」(202 Accepted)
  4. バックグラウンドのWorkerが5秒の処理を実行
  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
end

Jobの実装とコントローラー

# 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
end

INFO

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 AlertSNSTopic

WARNING

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
end

INFO

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
end

WARNING

非同期ジョブは必ず冪等にする

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: FARGATE

INFO

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呼び出し)