
Async Processing Best Practices
- 8 installs
- 3 repo stars
- Updated June 5, 2026
- xtone/ai_development_tools
Helps with ai & agent building tasks.
About
async-processing-best-practices is a Claude Code skill in the AI & Agent Building category.
- async-processing-best-practices
- AI & Agent Building
- AI-coding skill
Async Processing Best Practices by the numbers
- 8 all-time installs (skills.sh)
- +1 installs in the week ending Jul 27, 2026 (Skillselion tracking)
- Ranked #12,339 of 16,546 AI & Agent Building skills by installs in the Skillselion catalog
- Data as of Jul 27, 2026 (Skillselion catalog sync)
npx skills add https://github.com/xtone/ai_development_tools --skill async-processing-best-practicesAdd your badge
Show developers this skill is listed on Skillselion. Paste this into your README.
| Installs | 8 |
|---|---|
| repo stars | ★ 3 |
| Last updated | June 5, 2026 |
| Repository | xtone/ai_development_tools ↗ |
What it does
Helps with ai & agent building tasks.
Files
async-processing-best-practices
Webアプリケーションにおける非同期処理を実現するための実装ガイドです。
いつ利用するか
次の場合にこのガイドラインを参照してください。
- 特定の処理を非同期で実現するかどうかを設計する場合
- 非同期処理を行うインフラを選定・設計する場合
- 非同期処理が必要なバックエンド処理を実装する場合
- 非同期処理のUIをフロントエンドで実装する場合
リファレンス
非同期処理を利用するべきかどうかの判断
HTTPレスポンを返す際に処理が完了している必要がないものについては非同期処理で実現することを検討します。主に以下の観点で非同期処理の採用を考えます。
- ユーザー体験(UX)
- 処理に1秒以上かかり、UIから実行する際に画面が操作不能になるケースを避けたい場合
- 外部依存
- 外部APIやサービスなど、自分たちでコントロールできないものを呼び出す際に、相手側のシステムの障害や遅延に性能低下やユーザー体験の低下が発生することが懸念される場合
- 即時性の低さ
- ユーザーが処理の結果を即座に知らなくてもいい場合
- リソースの保護
- APIのコネクション数を長い処理によって占有して、接続リソースを不要に消費してしまうことが懸念される場合
非同期処理の実現手法
以下の要素から実現手法を選択します
- Webアプリケーションをホスティングするインフラ環境
- Webアプリケーション実装に利用するフレームワーク・ミドルウェア
- 利用するクラウドサービス(AWS/Google Cloud)
- 実現したいサービスの規模感や求められる堅牢性、予算など
主に以下のパターンによって構築します。
- Ruby on RailsのActiveJob + ActionCableによる実装
- SidekiqもしくはSolidQueue
- Railsを利用した非同期処理の実現手法については @references/01_infrastructure.md を参照してください
- TypeScript BullMQによる実装
- Redisベースの高性能ジョブキュー
- Node.js/TypeScriptでバックエンドを構築する場合に最適
- 詳細は @references/02_typescript_bullmq.md を参照してください
- AWS SQSによるタスクキューイングとAWS AppSync Eventsによるイベント送信、DynamoDBによる状態管理
- AWSのサーバーレスアーキテクチャで実現
- Lambda + SQS + AppSync Events + DynamoDBの組み合わせ
- 詳細は @references/03_aws_serverless.md を参照してください
- Google CloudのCloud Tasksを利用したタスク管理とFirestoreによる通知・状態管理
- Google Cloudのサーバーレスアーキテクチャで実現
- Cloud Run + Cloud Tasks + Firestoreの組み合わせ
- 詳細は @references/04_gcp_serverless.md を参照してください
実現手法の比較
| 観点 | Rails (Solid Queue) | Rails (Sidekiq) | BullMQ | AWS Serverless | GCP Serverless |
|---|---|---|---|---|---|
| 追加インフラ | なし | Redis | Redis | SQS, DynamoDB等 | Cloud Tasks, Firestore |
| スケーラビリティ | 中 | 高 | 高 | 非常に高 | 非常に高 |
| 運用負荷 | 低 | 中 | 中 | 低 | 低 |
| コスト(低負荷時) | 低 | 中 | 中 | 非常に低 | 非常に低 |
| リアルタイム通知 | ActionCable | ActionCable | 独自実装 | AppSync Events | Firestore |
| 最大処理時間 | 無制限 | 無制限 | 無制限 | 15分(Lambda) | 60分(Cloud Run) |
| 適したユースケース | 小〜中規模Rails | 中〜大規模Rails | Node.js全般 | AWS中心のサーバーレス | GCP中心のサーバーレス |
非同期処理のバックエンド実装
バックエンド実装時には、以下の観点について各実現手法ごとのベストプラクティスを参照してください。
共通の実装観点
- パラメータの受け取り方
- シリアライズ可能なデータのみを渡す
- パラメータサイズの制限を考慮(大きなデータは外部ストレージを使用)
- IDベースの参照を推奨
- バリデーションの実装
- エラー発生時の処理と通知
- リトライ可能なエラーと不可能なエラーの分類
- リトライ設定(回数、バックオフ戦略)
- Dead Letter Queue(DLQ)の処理
- エラー通知(構造化ログ、外部サービス連携)
- 進捗の状態管理
- 状態管理用データモデルの設計
- 進捗率の計算と更新
- リアルタイム通知(WebSocket、SSE、ポーリング)
- キャンセル処理
- その他
- 冪等性の確保
- タイムアウト対策
- ジョブの分割とチェーン
- テスト戦略
実現手法別リファレンス
| 実現手法 | リファレンス |
|---|---|
| Ruby on Rails (ActiveJob) | @references/05_backend_rails.md |
| TypeScript BullMQ | @references/06_backend_bullmq.md |
| AWS Serverless (Lambda + SQS) | @references/07_backend_aws.md |
| GCP Serverless (Cloud Run + Cloud Tasks) | @references/08_backend_gcp.md |
非同期処理のUIを実現するフロントエンド実装
フロントエンド実装時には、以下の観点について各フレームワークごとのベストプラクティスを参照してください。
共通の実装観点
- 非同期処理の呼び出し方
- フォーム送信による非同期タスクの開始
- カスタムフックやComposableによる呼び出し管理
- 二重送信の防止
- 状態管理と表示
- 進捗状況のリアルタイム表示
- ポーリングによる状態更新
- WebSocket/SSEによるプッシュ通知
- Optimistic Updates(楽観的更新)
- 接続復帰処理
- WebSocket/SSE接続の自動再接続
- ページ非表示時の接続管理
- Visibility APIを使用した状態同期
- オフライン/オンライン状態の監視
- エラーハンドリング
- Error Boundaryによるエラー捕捉
- グローバルエラーハンドリング
- リトライ機能(Exponential Backoff)
- ユーザーへのエラー通知
フレームワーク別リファレンス
| フレームワーク | リファレンス |
|---|---|
| Rails Hotwire (Turbo + Stimulus + ActionCable) | @references/09_frontend_hotwire.md |
| Next.js (React / App Router) | @references/10_frontend_nextjs.md |
| Vue.js 3 (Composition API) | @references/11_frontend_vuejs.md |
| React (Webpacker/Shakapacker) | @references/12_frontend_react.md |
本番環境へのデプロイ
本番環境にデプロイする際の設定やベストプラクティスは以下を参照してください。
- 本番環境デプロイガイド: @references/13_deployment.md
- Pumaプラグイン設定
- ECS/Fargate環境
- Kubernetes環境
- CI/CD設定
Solid Cable(ActionCableの代替)
Redis不要でActionCableを使用したい場合は、Solid Cableを利用できます。
- Solid Cableリファレンス: @references/14_solid_cable.md
- セットアップ方法
- ActionCableとの違い
- 設定オプション
トラブルシューティング
問題が発生した場合は以下を参照してください。
- トラブルシューティングガイド: @references/15_troubleshooting.md
- ジョブが実行されない場合
- ステータスが更新されない場合
- WebSocket接続の問題
- デバッグ方法
環境分離(DB共有時の注意)
ステージング環境と本番環境で同じデータベースを共有している場合、以下の設定が必要です。
Solid Queueのキュー名分離
# config/solid_queue.yml
staging:
workers:
- queues:
- <%= ENV.fetch('QUEUE_PREFIX', 'staging') %>_default
production:
workers:
- queues:
- <%= ENV.fetch('QUEUE_PREFIX', 'production') %>_default# app/jobs/application_job.rb
class ApplicationJob < ActiveJob::Base
queue_as do
prefix = ENV.fetch('QUEUE_PREFIX', Rails.env)
"#{prefix}_default"
end
endSolid Cableのチャンネル分離
# config/cable.yml
staging:
adapter: solid_cable
channel_prefix: <%= ENV.fetch('CABLE_PREFIX', 'staging') %>
production:
adapter: solid_cable
channel_prefix: <%= ENV.fetch('CABLE_PREFIX', 'production') %>詳細は @references/13_deployment.md の「環境分離設定」を参照してください。
Ruby on Railsによる実装
非同期処理をRuby on Railsで利用する場合は以下のようなフレームワークおよびインフラ構成を採用します。
ActiveJob + Solid Queue
アプリケーションのデータベースとして既にSqlite3、MySQLもしくはPostgreSQLが採用されていて、ユーザーのリクエスト数がそこまで多くなくデータベースのリソースに余裕がある状況なら、追加のリソースが必要ないSolidQueueを利用して実現します。 メリットとしてキューの管理やワーカーの実行のために新しくインフラを用意する必要がなく、既に利用しているRDBやpumaをそのまま利用することが出来ます。 RDBは既存のものを利用しつつ、ワーカーの実行のみWebと切り離すことも可能です。
ActiveJob + Sidekiq
データベースのリソースに余裕がなかったり、アクセス数が多くて既存のインフラと分離して非同期処理を実現したい場合はSidekiqを利用して実現します。 キュー用のDBとしてはRedisを利用し、Sidekiq用にワーカープロセスを実行する環境を既存のRailsアプリケーション実行環境とは別に用意する必要があります。
Solid QueueとSidekiqの採用基準
Solid Queue を選ぶべきケース
- Rails 8以上で新規開発を始める
- インフラ構成をシンプルに保ちたい(Redisを管理したくない)
- ジョブの実行をDBトランザクションと同期させたい(例:ユーザー作成に失敗したら、歓迎メール送信ジョブも自動でキャンセルしたい)
- 数百万件/日 程度の一般的な負荷である
Sidekiq を選ぶべきケース
- 既に Sidekiq Pro/Enterprise を契約している、またはその機能が必須である
- 既にキャッシュ用途などで Redis を運用しており、導入コストが低い
- 秒間数千件以上のジョブが走る、極めて高いスループットが求められる
- Sidekiq のリッチな Web UI や、長年のコミュニティ知見に頼りたい
Solid Queueの実行環境の採用基準
Webと統合する(pumaプラグインを利用)べきケース
- インフラのランニングコストや管理コストを抑えたい
- ジョブの負荷が低く、CPUを長時間占有しないジョブが中心
- メール送信、簡易的なデータ更新など。重い画像処理や機械学習の推論、大量のデータ処理などをこの方式で行うとWebサイトのレスポンスが極端に遅くなったりタイムアウトするリスクがあります。
Webと分離するべきケース
- ジョブの負荷が高かったり、処理時間が長いことが見込まれている
- Web側のリクエスト数が多く、Web側のサービスに影響が出ることを避けたい
- Web側とジョブ側でCPUやメモリの消費量が大きく異なることが想定される場合
インフラ構成
AWSおよびGoogle Cloudにホスティングする場合は以下のような構成で実現します。
- AWS
- ワーカー: Fargateを利用
- Web用とジョブ用のDocker imageはWeb用とWorker用で同一のものを利用
- Solid Queueをpumaプラグインを利用する場合はWeb用のサービスにジョブを実行させる
- ジョブ用のワーカーを別に設定する場合は、ジョブ用のタスク定義を設定し異なるサービスとして実行させる
- キューDB
- Solid Queueの場合はWeb用のRDB(RDS/Aurora)を利用する
- Sidekiqの場合はElasticacheを利用してRedisもしくはValkeyを利用する
- 既にRedisを利用している場合は相乗りを検討し、新規に作る場合はValkeyの採用を検討する
- Google Cloud
- ワーカー: Cloud Runを利用
- Web用とジョブ用のDocker imageはWeb用とWorker用で同一のものを利用
- Solid Queueをpumaプラグインを利用する場合はWeb用のサービスにジョブを実行させる
- min-instances: 1以上に設定しないとジョブが実行されない点に注意
- ジョブ用のワーカーを別に設定する場合、Web用とは別のサービスを実行させる
- この場合もmin-instances: 1以上に設定する必要がある
- ワーカー用のサービスが停止したときのことを考慮する
- ワーカー用のプロセスにSIGTERMが送られた際にGraceful Shutdownが正しく動くようにする必要がある
- キューDB
- Solid Queueの場合はWeb用のCloud SQLを利用する
- Sidekiqの場合はCloud Memorystore for Redisを利用する
Dockerfile
Rails 8のベストプラクティスに基づいた、軽量かつセキュアなDockerfileのサンプルです。事前にインストールするライブラリ(libpq-dev等)はGemfileに応じて調整してください。
# syntax = docker/dockerfile:1
ARG RUBY_VERSION=3.3.0
FROM ruby:$RUBY_VERSION-slim AS base
WORKDIR /rails
ENV RAILS_ENV="production" \
BUNDLE_DEPLOYMENT="1" \
BUNDLE_PATH="/usr/local/bundle" \
BUNDLE_WITHOUT="development test"
# --- Build stage ---
FROM base AS build
RUN apt-get update -qq && \
apt-get install --no-install-recommends -y build-essential git libpq-dev pkg-config
COPY Gemfile Gemfile.lock ./
RUN bundle install && \
rm -rf ~/.bundle/ "${BUNDLE_PATH}"/ruby/*/cache "${BUNDLE_PATH}"/ruby/*/bundler/gems/*/.git
COPY . .
RUN bundle exec bootsnap precompile --gemfile app/ lib/
# Asset precompilation (SECRET_KEY_BASEはダミーでOK)
RUN SECRET_KEY_BASE_DUMMY=1 ./bin/rails assets:precompile
# --- Final stage ---
FROM base
RUN apt-get update -qq && \
apt-get install --no-install-recommends -y curl libpq5 && \
rm -rf /var/lib/apt/lists /var/cache/apt/archives
COPY --from=build /usr/local/bundle /usr/local/bundle
COPY --from=build /rails /rails
# 非特権ユーザーで実行
RUN useradd rails --create-home --shell /bin/bash && \
chown -R rails:rails db log storage tmp
USER rails:rails
ENTRYPOINT ["/rails/bin/docker-entrypoint"]
EXPOSE 3000
CMD ["./bin/rails", "server"]AWS: Terraform (ECS Fargate)
WebとWorkerを同一イメージ、別サービスで定義する構成のサンプルです。
# ECS Task Definition (Shared)
resource "aws_ecs_task_definition" "app" {
family = "rails-app"
network_mode = "awsvpc"
requires_compatibilities = ["FARGATE"]
cpu = "512"
memory = "1024"
execution_role_arn = aws_iam_role.ecs_task_execution_role.arn
task_role_arn = aws_iam_role.ecs_task_role.arn
container_definitions = jsonencode([{
name = "rails-container"
image = "${aws_ecr_repository.app.repository_url}:latest"
portMappings = [{ containerPort = 3000 }]
environment = [
{ name = "DATABASE_URL", value = "postgres://..." },
{ name = "RAILS_ENV", value = "production" },
{ name = "SOLID_QUEUE_IN_PUMA", value = "true" } # 統合型の場合
]
logConfiguration = {
logDriver = "awslogs"
options = {
"awslogs-group" = "/ecs/rails-app"
"awslogs-region" = "ap-northeast-1"
"awslogs-stream-prefix" = "ecs"
}
}
}])
}
# ECS Service (Web)
resource "aws_ecs_service" "web" {
name = "web-service"
cluster = aws_ecs_cluster.main.id
task_definition = aws_ecs_task_definition.app.arn
desired_count = 2
launch_type = "FARGATE"
network_configuration {
subnets = aws_subnet.private[*].id
security_groups = [aws_security_group.ecs.id]
}
load_balancer {
target_group_arn = aws_lb_target_group.app.arn
container_name = "rails-container"
container_port = 3000
}
}Google Cloud: Terraform (Cloud Run v2)
Direct VPC Egressを利用したモダンなCloud Run構成です。
# Cloud Run Service (Unified Web & Worker)
resource "google_cloud_run_v2_service" "app" {
name = "rails-app"
location = "asia-northeast1"
template {
containers {
image = "asia-northeast1-docker.pkg.dev/project/repo/image:latest"
env {
name = "DATABASE_URL"
value = "postgres://user:pass@10.x.x.x:5432/dbname"
}
env {
name = "SOLID_QUEUE_IN_PUMA"
value = "true"
}
ports {
container_port = 3000
}
resources {
limits = {
cpu = "1"
memory = "1Gi"
}
}
}
# Direct VPC Egress 設定
vpc_access {
network_interfaces {
network = "default"
subnets = "default"
}
egress = "ALL_TRAFFIC"
}
scaling {
min_instance_count = 1 # Solid Queueのポーリングを維持するため
max_instance_count = 10
}
}
}フロントエンドへの状態通知(ActionCable)
非同期ジョブの進捗や完了状態をフロントエンドにリアルタイムで通知するには、ActionCableを使用します。
ActionCable Channel
# app/channels/job_status_channel.rb
class JobStatusChannel < ApplicationCable::Channel
def subscribed
job_id = params[:job_id]
stream_from "job_status_#{job_id}"
end
def unsubscribed
# クリーンアップ処理(必要に応じて)
end
endジョブからの通知
# app/jobs/example_job.rb
class ExampleJob < ApplicationJob
queue_as :default
def perform(job_record_id)
job_record = JobRecord.find(job_record_id)
# 処理開始を通知
broadcast_status(job_record, 'processing', 0)
# 実際の処理
result = process_task(job_record) do |progress|
broadcast_status(job_record, 'processing', progress)
end
# 完了を通知
job_record.update!(status: 'completed', result: result)
broadcast_status(job_record, 'completed', 100, result)
rescue StandardError => e
job_record.update!(status: 'failed', error_message: e.message)
broadcast_status(job_record, 'failed', nil, nil, e.message)
raise
end
private
def broadcast_status(job_record, status, progress, result = nil, error = nil)
ActionCable.server.broadcast(
"job_status_#{job_record.id}",
{
job_id: job_record.id,
status: status,
progress: progress,
result: result,
error: error,
updated_at: Time.current.iso8601
}
)
end
def process_task(job_record)
# 進捗を報告しながら処理を実行
total_steps = 10
total_steps.times do |i|
# 実際の処理...
sleep 1
yield ((i + 1) * 100 / total_steps) if block_given?
end
{ message: 'Task completed successfully' }
end
endフロントエンド(JavaScript)
// app/javascript/channels/job_status_channel.js
import consumer from "./consumer"
export function subscribeToJobStatus(jobId, callbacks) {
return consumer.subscriptions.create(
{ channel: "JobStatusChannel", job_id: jobId },
{
received(data) {
switch (data.status) {
case 'processing':
callbacks.onProgress?.(data.progress)
break
case 'completed':
callbacks.onComplete?.(data.result)
break
case 'failed':
callbacks.onError?.(data.error)
break
}
},
connected() {
callbacks.onConnected?.()
},
disconnected() {
callbacks.onDisconnected?.()
}
}
)
}React Hook(Hotwire/Stimulus以外の場合)
// frontend/src/hooks/useJobStatus.ts
import { useEffect, useState, useCallback } from 'react';
import { createConsumer } from '@rails/actioncable';
interface JobStatus {
status: 'pending' | 'processing' | 'completed' | 'failed';
progress?: number;
result?: Record<string, unknown>;
error?: string;
}
const consumer = createConsumer('/cable');
export function useJobStatus(jobId: string | null) {
const [status, setStatus] = useState<JobStatus | null>(null);
useEffect(() => {
if (!jobId) return;
const subscription = consumer.subscriptions.create(
{ channel: 'JobStatusChannel', job_id: jobId },
{
received(data: JobStatus) {
setStatus(data);
},
}
);
return () => {
subscription.unsubscribe();
};
}, [jobId]);
return status;
}Hotwire/Turbo Streams(Rails 7+推奨)
# app/jobs/example_job.rb(Turbo Streams版)
class ExampleJob < ApplicationJob
def perform(job_record_id)
job_record = JobRecord.find(job_record_id)
job_record.update!(status: 'processing')
broadcast_turbo_stream(job_record)
# 処理実行...
job_record.update!(status: 'completed', result: result)
broadcast_turbo_stream(job_record)
end
private
def broadcast_turbo_stream(job_record)
Turbo::StreamsChannel.broadcast_replace_to(
"job_#{job_record.id}",
target: "job_status_#{job_record.id}",
partial: "job_records/status",
locals: { job_record: job_record }
)
end
end<!-- app/views/job_records/_status.html.erb -->
<div id="job_status_<%= job_record.id %>">
<p>Status: <%= job_record.status %></p>
<% if job_record.processing? %>
<progress value="<%= job_record.progress %>" max="100"></progress>
<% end %>
<% if job_record.completed? %>
<p>Result: <%= job_record.result %></p>
<% end %>
<% if job_record.failed? %>
<p class="error">Error: <%= job_record.error_message %></p>
<% end %>
</div><!-- View側でのsubscribe -->
<%= turbo_stream_from "job_#{@job_record.id}" %>
<%= render 'job_records/status', job_record: @job_record %>インフラ設定の注意点
AWS (ALB + ECS)
ActionCableのWebSocket接続を維持するため、ALBのスティッキーセッションとアイドルタイムアウトを設定します。
resource "aws_lb_target_group" "app" {
# ... 既存の設定 ...
stickiness {
type = "lb_cookie"
cookie_duration = 86400
enabled = true
}
}
resource "aws_lb_listener" "https" {
# ... 既存の設定 ...
# WebSocket用にアイドルタイムアウトを延長
idle_timeout = 3600
}Google Cloud (Cloud Run)
Cloud RunではWebSocket接続は最大60分まで維持されます。長時間のジョブには再接続ロジックを実装してください。
resource "google_cloud_run_v2_service" "app" {
template {
# ... 既存の設定 ...
# セッションアフィニティを有効化
session_affinity = true
}
}外部フロントエンド(Next.js等)との連携
RailsをAPIサーバーとして使用し、フロントエンドをNext.js等で別途構築する場合のベストプラクティスです。
通知方法の選択
| 方法 | 特徴 | 推奨ケース |
|---|---|---|
| ActionCable(スタンドアロン) | Rails標準、双方向通信 | リアルタイム性が重要、双方向通信が必要 |
| Server-Sent Events(SSE) | シンプル、軽量、HTTP標準 | 単方向通知で十分、シンプルな実装を好む |
| ポーリング | 最もシンプル、インフラ制約なし | 更新頻度が低い、WebSocket非対応環境 |
CORS設定
# config/initializers/cors.rb
Rails.application.config.middleware.insert_before 0, Rack::Cors do
allow do
origins ENV.fetch('FRONTEND_URL', 'http://localhost:3000')
resource '*',
headers: :any,
methods: [:get, :post, :put, :patch, :delete, :options, :head],
credentials: true
# ActionCable用のWebSocket接続を許可
resource '/cable',
headers: :any,
methods: [:get, :post, :options],
credentials: true
end
end認証(JWT)
外部フロントエンドとの連携では、セッションベースではなくJWT認証を使用することが一般的です。
# Gemfile
gem 'jwt'# app/services/jwt_service.rb
class JwtService
SECRET_KEY = Rails.application.credentials.secret_key_base
def self.encode(payload, exp = 24.hours.from_now)
payload[:exp] = exp.to_i
JWT.encode(payload, SECRET_KEY, 'HS256')
end
def self.decode(token)
decoded = JWT.decode(token, SECRET_KEY, true, algorithm: 'HS256')
HashWithIndifferentAccess.new(decoded.first)
rescue JWT::DecodeError, JWT::ExpiredSignature
nil
end
endActionCable + JWT認証
# app/channels/application_cable/connection.rb
module ApplicationCable
class Connection < ActionCable::Connection::Base
identified_by :current_user
def connect
self.current_user = find_verified_user
end
private
def find_verified_user
# クエリパラメータからトークンを取得
token = request.params[:token]
return reject_unauthorized_connection unless token
payload = JwtService.decode(token)
return reject_unauthorized_connection unless payload
user = User.find_by(id: payload[:user_id])
return reject_unauthorized_connection unless user
user
end
end
end# config/initializers/action_cable.rb
Rails.application.config.action_cable.allowed_request_origins = [
ENV.fetch('FRONTEND_URL', 'http://localhost:3000'),
/https?:\/\/localhost:\d+/
]
# 本番環境ではURLを明示的に指定
if Rails.env.production?
Rails.application.config.action_cable.allowed_request_origins = [
ENV.fetch('FRONTEND_URL')
]
endNext.js フロントエンド実装(ActionCable)
// lib/actionCable.ts
import { createConsumer, Consumer, Subscription } from '@rails/actioncable';
let consumer: Consumer | null = null;
export const getConsumer = (token: string): Consumer => {
if (!consumer) {
const wsUrl = `${process.env.NEXT_PUBLIC_WS_URL}/cable?token=${token}`;
consumer = createConsumer(wsUrl);
}
return consumer;
};
export const disconnectConsumer = (): void => {
if (consumer) {
consumer.disconnect();
consumer = null;
}
};// hooks/useJobStatus.ts
import { useEffect, useState, useRef } from 'react';
import { getConsumer } from '../lib/actionCable';
import { Subscription } from '@rails/actioncable';
import { useAuth } from './useAuth'; // JWT認証フック
interface JobStatus {
status: 'pending' | 'processing' | 'completed' | 'failed';
progress?: number;
result?: Record<string, unknown>;
error?: string;
}
export const useJobStatus = (jobId: string | null) => {
const [status, setStatus] = useState<JobStatus | null>(null);
const [connected, setConnected] = useState(false);
const subscriptionRef = useRef<Subscription | null>(null);
const { token } = useAuth();
useEffect(() => {
if (!jobId || !token) return;
const consumer = getConsumer(token);
subscriptionRef.current = consumer.subscriptions.create(
{ channel: 'JobStatusChannel', job_id: jobId },
{
connected() {
setConnected(true);
},
disconnected() {
setConnected(false);
},
received(data: JobStatus) {
setStatus(data);
},
}
);
return () => {
subscriptionRef.current?.unsubscribe();
subscriptionRef.current = null;
};
}, [jobId, token]);
return { status, connected };
};Server-Sent Events(SSE)による実装
ActionCableより軽量なSSEを使用する場合の実装です。
# app/controllers/api/v1/job_streams_controller.rb
module Api
module V1
class JobStreamsController < ApplicationController
include ActionController::Live
before_action :authenticate_user!
before_action :set_job
def show
response.headers['Content-Type'] = 'text/event-stream'
response.headers['Cache-Control'] = 'no-cache'
response.headers['X-Accel-Buffering'] = 'no' # nginx用
sse = SSE.new(response.stream, retry: 3000, event: 'job-status')
begin
# 初期状態を送信
sse.write(job_status_data)
# Redis Pub/Subで更新を購読
redis = Redis.new(url: ENV['REDIS_URL'])
redis.subscribe("job_status_#{@job.id}") do |on|
on.message do |_channel, message|
data = JSON.parse(message)
sse.write(data)
# 完了または失敗で終了
break if %w[completed failed].include?(data['status'])
end
end
rescue ActionController::Live::ClientDisconnected
# クライアント切断
ensure
redis&.close
sse.close
end
end
private
def set_job
@job = current_user.jobs.find(params[:id])
end
def job_status_data
{
job_id: @job.id,
status: @job.status,
progress: @job.progress,
result: @job.result,
error: @job.error_message
}
end
end
end
end# SSEヘルパークラス
# lib/sse.rb
class SSE
def initialize(stream, options = {})
@stream = stream
@options = options
end
def write(data, options = {})
options = @options.merge(options)
options.each do |key, value|
@stream.write("#{key}: #{value}\n")
end
@stream.write("data: #{data.to_json}\n\n")
end
def close
@stream.close
end
end# ジョブから通知を送信
# app/jobs/example_job.rb
class ExampleJob < ApplicationJob
def perform(job_record_id)
job_record = JobRecord.find(job_record_id)
redis = Redis.new(url: ENV['REDIS_URL'])
begin
publish_status(redis, job_record, 'processing', 0)
result = process_task(job_record) do |progress|
publish_status(redis, job_record, 'processing', progress)
end
job_record.update!(status: 'completed', result: result)
publish_status(redis, job_record, 'completed', 100, result)
rescue StandardError => e
job_record.update!(status: 'failed', error_message: e.message)
publish_status(redis, job_record, 'failed', nil, nil, e.message)
raise
ensure
redis.close
end
end
private
def publish_status(redis, job_record, status, progress, result = nil, error = nil)
data = {
job_id: job_record.id,
status: status,
progress: progress,
result: result,
error: error,
updated_at: Time.current.iso8601
}
redis.publish("job_status_#{job_record.id}", data.to_json)
# ActionCableにも通知(両方使う場合)
ActionCable.server.broadcast("job_status_#{job_record.id}", data)
end
endNext.js フロントエンド実装(SSE)
// hooks/useJobStatusSSE.ts
import { useState, useEffect, useRef } from 'react';
import { useAuth } from './useAuth';
interface JobStatus {
status: 'pending' | 'processing' | 'completed' | 'failed';
progress?: number;
result?: Record<string, unknown>;
error?: string;
}
export const useJobStatusSSE = (jobId: string | null) => {
const [status, setStatus] = useState<JobStatus | null>(null);
const [connected, setConnected] = useState(false);
const [error, setError] = useState<string | null>(null);
const eventSourceRef = useRef<EventSource | null>(null);
const { token } = useAuth();
useEffect(() => {
if (!jobId || !token) return;
const url = `${process.env.NEXT_PUBLIC_API_URL}/api/v1/jobs/${jobId}/stream`;
// SSEはヘッダーを送れないため、クエリパラメータでトークンを渡す
const eventSource = new EventSource(`${url}?token=${token}`, {
withCredentials: true,
});
eventSourceRef.current = eventSource;
eventSource.onopen = () => {
setConnected(true);
setError(null);
};
eventSource.addEventListener('job-status', (event) => {
const data = JSON.parse(event.data);
setStatus({
status: data.status,
progress: data.progress,
result: data.result,
error: data.error,
});
// 完了/失敗で接続を閉じる
if (['completed', 'failed'].includes(data.status)) {
eventSource.close();
setConnected(false);
}
});
eventSource.onerror = () => {
setError('接続エラーが発生しました');
setConnected(false);
eventSource.close();
};
return () => {
eventSource.close();
eventSourceRef.current = null;
};
}, [jobId, token]);
return { status, connected, error };
};ポーリングによる実装
WebSocketやSSEが使用できない環境向けのフォールバック実装です。
# app/controllers/api/v1/jobs_controller.rb
module Api
module V1
class JobsController < ApplicationController
before_action :authenticate_user!
def show
job = current_user.jobs.find(params[:id])
render json: {
id: job.id,
status: job.status,
progress: job.progress,
result: job.result,
error: job.error_message,
updated_at: job.updated_at.iso8601
}
end
end
end
end// hooks/useJobStatusPolling.ts
import { useState, useEffect, useCallback, useRef } from 'react';
import { useAuth } from './useAuth';
interface JobStatus {
status: 'pending' | 'processing' | 'completed' | 'failed';
progress?: number;
result?: Record<string, unknown>;
error?: string;
}
interface UseJobStatusPollingOptions {
interval?: number; // ポーリング間隔(ミリ秒)
enabled?: boolean;
}
export const useJobStatusPolling = (
jobId: string | null,
options: UseJobStatusPollingOptions = {}
) => {
const { interval = 2000, enabled = true } = options;
const [status, setStatus] = useState<JobStatus | null>(null);
const [loading, setLoading] = useState(false);
const [error, setError] = useState<string | null>(null);
const { token } = useAuth();
const intervalRef = useRef<NodeJS.Timeout | null>(null);
const fetchStatus = useCallback(async () => {
if (!jobId || !token) return;
try {
setLoading(true);
const response = await fetch(
`${process.env.NEXT_PUBLIC_API_URL}/api/v1/jobs/${jobId}`,
{
headers: {
Authorization: `Bearer ${token}`,
'Content-Type': 'application/json',
},
}
);
if (!response.ok) {
throw new Error('Failed to fetch job status');
}
const data = await response.json();
setStatus({
status: data.status,
progress: data.progress,
result: data.result,
error: data.error,
});
setError(null);
// 完了/失敗でポーリング停止
if (['completed', 'failed'].includes(data.status)) {
if (intervalRef.current) {
clearInterval(intervalRef.current);
intervalRef.current = null;
}
}
} catch (err) {
setError(err instanceof Error ? err.message : 'Unknown error');
} finally {
setLoading(false);
}
}, [jobId, token]);
useEffect(() => {
if (!jobId || !token || !enabled) return;
// 初回fetch
fetchStatus();
// ポーリング開始
intervalRef.current = setInterval(fetchStatus, interval);
return () => {
if (intervalRef.current) {
clearInterval(intervalRef.current);
intervalRef.current = null;
}
};
}, [jobId, token, enabled, interval, fetchStatus]);
return { status, loading, error, refetch: fetchStatus };
};統合フック(フォールバック付き)
WebSocket → SSE → ポーリングの順でフォールバックする統合フックです。
// hooks/useJobStatusWithFallback.ts
import { useState, useEffect } from 'react';
import { useJobStatus } from './useJobStatus'; // ActionCable
import { useJobStatusSSE } from './useJobStatusSSE';
import { useJobStatusPolling } from './useJobStatusPolling';
type TransportType = 'websocket' | 'sse' | 'polling';
interface JobStatus {
status: 'pending' | 'processing' | 'completed' | 'failed';
progress?: number;
result?: Record<string, unknown>;
error?: string;
}
export const useJobStatusWithFallback = (jobId: string | null) => {
const [transport, setTransport] = useState<TransportType>('websocket');
// WebSocket(ActionCable)
const {
status: wsStatus,
connected: wsConnected,
} = useJobStatus(transport === 'websocket' ? jobId : null);
// SSE
const {
status: sseStatus,
connected: sseConnected,
error: sseError,
} = useJobStatusSSE(transport === 'sse' ? jobId : null);
// ポーリング
const {
status: pollingStatus,
} = useJobStatusPolling(transport === 'polling' ? jobId : null);
// フォールバックロジック
useEffect(() => {
if (transport === 'websocket' && !wsConnected) {
// WebSocket接続失敗 → SSEにフォールバック
const timeout = setTimeout(() => {
if (!wsConnected) {
console.log('Falling back to SSE');
setTransport('sse');
}
}, 5000);
return () => clearTimeout(timeout);
}
}, [transport, wsConnected]);
useEffect(() => {
if (transport === 'sse' && sseError) {
// SSE接続失敗 → ポーリングにフォールバック
console.log('Falling back to polling');
setTransport('polling');
}
}, [transport, sseError]);
// 現在のトランスポートに応じたステータスを返す
const status = transport === 'websocket'
? wsStatus
: transport === 'sse'
? sseStatus
: pollingStatus;
return { status, transport };
};API設計のベストプラクティス
# config/routes.rb
Rails.application.routes.draw do
namespace :api do
namespace :v1 do
resources :jobs, only: [:create, :show] do
member do
get :stream # SSE用
end
end
end
end
# ActionCable
mount ActionCable.server => '/cable'
end# app/controllers/api/v1/jobs_controller.rb
module Api
module V1
class JobsController < ApplicationController
before_action :authenticate_user!
# POST /api/v1/jobs
def create
job = current_user.jobs.create!(
job_type: params[:type],
payload: params[:payload],
status: 'pending'
)
# ジョブをキューに投入
ProcessJobWorker.perform_async(job.id)
render json: {
id: job.id,
status: job.status,
# リアルタイム更新用のエンドポイント情報
_links: {
self: api_v1_job_url(job),
stream: stream_api_v1_job_url(job),
websocket: "#{websocket_url}/cable?token=#{current_token}"
}
}, status: :accepted
end
end
end
end環境変数設定例
# .env.production (Rails)
FRONTEND_URL=https://app.example.com
REDIS_URL=redis://localhost:6379/1
# .env.production (Next.js)
NEXT_PUBLIC_API_URL=https://api.example.com
NEXT_PUBLIC_WS_URL=wss://api.example.comTypeScript BullMQによる実装
非同期処理をTypeScript/Node.jsで実装する場合は、BullMQを採用します。BullMQはRedisをバックエンドとした高性能なジョブキューライブラリで、堅牢なリトライ機能、優先度付きキュー、レート制限など豊富な機能を提供します。
BullMQの特徴
- 高いパフォーマンス: Redisベースで秒間数千ジョブの処理が可能
- 型安全: TypeScriptファーストで設計されており、ジョブのペイロードやリターン値の型定義が可能
- 豊富な機能: 遅延ジョブ、繰り返しジョブ、優先度付きキュー、レート制限
- 可観測性: Bull Boardなどの管理UIが利用可能
- スケーラビリティ: ワーカーを水平スケールして負荷分散が可能
採用基準
BullMQ を選ぶべきケース
- Node.js/TypeScriptでバックエンドを構築している
- Redisを既に利用している、または導入に抵抗がない
- 複雑なジョブフロー(依存関係、子ジョブ)が必要
- 秒間数百〜数千ジョブの高スループットが求められる
- ジョブの優先度制御やレート制限が必要
他の選択肢を検討すべきケース
- サーバーレス環境でコールドスタートを避けたい場合 → AWS SQS + Lambda
- Redisの運用を避けたい場合 → Cloud Tasks/SQS
- シンプルなバッチ処理のみの場合 → シンプルなcronジョブ
アーキテクチャ
基本構成
┌─────────────┐ ┌─────────────┐ ┌─────────────┐
│ API Server │────▶│ Redis │◀────│ Worker │
│ (Producer) │ │ (Queue) │ │ (Consumer) │
└─────────────┘ └─────────────┘ └─────────────┘
│ │
│ ┌─────────────┐ │
└────────▶│ Database │◀─────────────┘
│ (状態管理) │
└─────────────┘コンポーネント
- Producer: APIサーバーがジョブをキューに投入
- Queue: Redisがジョブの永続化と配信を担当
- Worker: 別プロセス/コンテナでジョブを処理
- Database: ジョブの状態やビジネスデータを永続化
インフラ構成
AWS構成
- API Server: ECS Fargate または App Runner
- Worker: ECS Fargate(別タスク定義)
- Redis: ElastiCache for Redis または MemoryDB for Redis
- Database: RDS PostgreSQL または Aurora
Google Cloud構成
- API Server: Cloud Run
- Worker: Cloud Run Jobs または Compute Engine
- Redis: Memorystore for Redis
- Database: Cloud SQL
実装サンプル
パッケージインストール
# ioredis v5+はTypeScript型定義を内蔵しているため、@types/ioredisは不要
npm install bullmq ioredisキューの定義
// src/queues/types.ts
export interface EmailJobData {
to: string;
subject: string;
body: string;
templateId?: string;
}
export interface EmailJobResult {
messageId: string;
sentAt: Date;
}
export const EMAIL_QUEUE_NAME = 'email-queue';// src/queues/connection.ts
import { Redis } from 'ioredis';
const redisConfig = {
host: process.env.REDIS_HOST || 'localhost',
port: parseInt(process.env.REDIS_PORT || '6379'),
maxRetriesPerRequest: null, // BullMQの要件
};
export const createRedisConnection = () => new Redis(redisConfig);Producer(ジョブ投入側)
// src/queues/emailQueue.ts
import { Queue } from 'bullmq';
import { createRedisConnection } from './connection';
import { EmailJobData, EmailJobResult, EMAIL_QUEUE_NAME } from './types';
export const emailQueue = new Queue<EmailJobData, EmailJobResult>(
EMAIL_QUEUE_NAME,
{
connection: createRedisConnection(),
defaultJobOptions: {
attempts: 3,
backoff: {
type: 'exponential',
delay: 1000,
},
removeOnComplete: {
age: 24 * 3600, // 24時間後に完了ジョブを削除
count: 1000, // 最大1000件保持
},
removeOnFail: {
age: 7 * 24 * 3600, // 7日後に失敗ジョブを削除
},
},
}
);
// ジョブ投入ヘルパー
export const enqueueEmail = async (data: EmailJobData, options?: {
delay?: number;
priority?: number;
}) => {
const job = await emailQueue.add('send-email', data, {
delay: options?.delay,
priority: options?.priority,
});
return job.id;
};Worker(ジョブ処理側)
// src/workers/emailWorker.ts
import { Worker, Job } from 'bullmq';
import { createRedisConnection } from '../queues/connection';
import { EmailJobData, EmailJobResult, EMAIL_QUEUE_NAME } from '../queues/types';
import { sendEmail } from '../services/emailService';
import { logger } from '../utils/logger';
const processEmail = async (
job: Job<EmailJobData, EmailJobResult>
): Promise<EmailJobResult> => {
const { to, subject, body, templateId } = job.data;
logger.info(`Processing email job ${job.id}`, { to, subject });
// 進捗更新
await job.updateProgress(10);
try {
const result = await sendEmail({ to, subject, body, templateId });
await job.updateProgress(100);
return {
messageId: result.messageId,
sentAt: new Date(),
};
} catch (error) {
logger.error(`Failed to send email`, { jobId: job.id, error });
throw error; // リトライのために再スロー
}
};
export const createEmailWorker = () => {
const worker = new Worker<EmailJobData, EmailJobResult>(
EMAIL_QUEUE_NAME,
processEmail,
{
connection: createRedisConnection(),
concurrency: 10, // 同時処理数
limiter: {
max: 100, // 最大100ジョブ
duration: 1000, // 1秒あたり
},
}
);
worker.on('completed', (job, result) => {
logger.info(`Job ${job.id} completed`, { result });
});
worker.on('failed', (job, error) => {
logger.error(`Job ${job?.id} failed`, { error: error.message });
});
worker.on('error', (error) => {
logger.error('Worker error', { error });
});
return worker;
};ワーカーの起動スクリプト
// src/worker.ts
import { createEmailWorker } from './workers/emailWorker';
import { logger } from './utils/logger';
const workers = [
createEmailWorker(),
];
// Graceful shutdown
const shutdown = async () => {
logger.info('Shutting down workers...');
await Promise.all(workers.map(w => w.close()));
process.exit(0);
};
process.on('SIGTERM', shutdown);
process.on('SIGINT', shutdown);
logger.info('Workers started');APIエンドポイント例
// src/routes/email.ts
import { Router } from 'express';
import { enqueueEmail } from '../queues/emailQueue';
import { emailQueue } from '../queues/emailQueue';
const router = Router();
// メール送信リクエスト
router.post('/send', async (req, res) => {
const { to, subject, body } = req.body;
const jobId = await enqueueEmail({ to, subject, body });
res.status(202).json({
status: 'queued',
jobId,
});
});
// ジョブ状態確認
router.get('/status/:jobId', async (req, res) => {
const job = await emailQueue.getJob(req.params.jobId);
if (!job) {
return res.status(404).json({ error: 'Job not found' });
}
const state = await job.getState();
const progress = job.progress;
res.json({
jobId: job.id,
state,
progress,
data: job.data,
result: job.returnvalue,
failedReason: job.failedReason,
});
});
export default router;Dockerfile
# syntax = docker/dockerfile:1
ARG NODE_VERSION=20
FROM node:${NODE_VERSION}-slim AS base
WORKDIR /app
ENV NODE_ENV="production"
# --- Build stage ---
FROM base AS build
RUN apt-get update -qq && \
apt-get install --no-install-recommends -y build-essential python3
COPY package*.json ./
RUN npm ci --include=dev
COPY . .
RUN npm run build && \
npm prune --production
# --- Final stage ---
FROM base
RUN apt-get update -qq && \
apt-get install --no-install-recommends -y curl && \
rm -rf /var/lib/apt/lists/*
COPY --from=build /app/node_modules /app/node_modules
COPY --from=build /app/dist /app/dist
COPY --from=build /app/package.json /app/package.json
# 非特権ユーザーで実行
RUN useradd -m appuser
USER appuser
# ヘルスチェック
HEALTHCHECK --interval=30s --timeout=3s \
CMD curl -f http://localhost:3000/health || exit 1
EXPOSE 3000
CMD ["node", "dist/server.js"]ワーカー用起動コマンド
# ワーカーとして起動する場合
CMD ["node", "dist/worker.js"]AWS: Terraform (ECS Fargate)
# Variables
variable "environment" {
default = "production"
}
variable "redis_node_type" {
default = "cache.t3.micro"
}
# ElastiCache Redis
resource "aws_elasticache_subnet_group" "redis" {
name = "${var.environment}-redis-subnet"
subnet_ids = aws_subnet.private[*].id
}
resource "aws_elasticache_replication_group" "redis" {
replication_group_id = "${var.environment}-bullmq-redis"
description = "Redis for BullMQ"
node_type = var.redis_node_type
num_cache_clusters = 2
port = 6379
parameter_group_name = "default.redis7"
automatic_failover_enabled = true
subnet_group_name = aws_elasticache_subnet_group.redis.name
security_group_ids = [aws_security_group.redis.id]
at_rest_encryption_enabled = true
transit_encryption_enabled = true
}
# Security Group for Redis
resource "aws_security_group" "redis" {
name = "${var.environment}-redis-sg"
description = "Security group for Redis"
vpc_id = aws_vpc.main.id
ingress {
from_port = 6379
to_port = 6379
protocol = "tcp"
security_groups = [aws_security_group.ecs.id]
}
}
# ECS Task Definition (API)
resource "aws_ecs_task_definition" "api" {
family = "${var.environment}-api"
network_mode = "awsvpc"
requires_compatibilities = ["FARGATE"]
cpu = "512"
memory = "1024"
execution_role_arn = aws_iam_role.ecs_task_execution_role.arn
task_role_arn = aws_iam_role.ecs_task_role.arn
container_definitions = jsonencode([{
name = "api"
image = "${aws_ecr_repository.app.repository_url}:latest"
portMappings = [{ containerPort = 3000 }]
environment = [
{ name = "NODE_ENV", value = "production" },
{ name = "REDIS_HOST", value = aws_elasticache_replication_group.redis.primary_endpoint_address },
{ name = "REDIS_PORT", value = "6379" },
{ name = "DATABASE_URL", value = "postgres://..." }
]
command = ["node", "dist/server.js"]
logConfiguration = {
logDriver = "awslogs"
options = {
"awslogs-group" = "/ecs/${var.environment}/api"
"awslogs-region" = "ap-northeast-1"
"awslogs-stream-prefix" = "ecs"
}
}
healthCheck = {
command = ["CMD-SHELL", "curl -f http://localhost:3000/health || exit 1"]
interval = 30
timeout = 5
retries = 3
startPeriod = 60
}
}])
}
# ECS Task Definition (Worker)
resource "aws_ecs_task_definition" "worker" {
family = "${var.environment}-worker"
network_mode = "awsvpc"
requires_compatibilities = ["FARGATE"]
cpu = "256"
memory = "512"
execution_role_arn = aws_iam_role.ecs_task_execution_role.arn
task_role_arn = aws_iam_role.ecs_task_role.arn
container_definitions = jsonencode([{
name = "worker"
image = "${aws_ecr_repository.app.repository_url}:latest"
environment = [
{ name = "NODE_ENV", value = "production" },
{ name = "REDIS_HOST", value = aws_elasticache_replication_group.redis.primary_endpoint_address },
{ name = "REDIS_PORT", value = "6379" },
{ name = "DATABASE_URL", value = "postgres://..." }
]
command = ["node", "dist/worker.js"]
logConfiguration = {
logDriver = "awslogs"
options = {
"awslogs-group" = "/ecs/${var.environment}/worker"
"awslogs-region" = "ap-northeast-1"
"awslogs-stream-prefix" = "ecs"
}
}
}])
}
# ECS Service (API)
resource "aws_ecs_service" "api" {
name = "api-service"
cluster = aws_ecs_cluster.main.id
task_definition = aws_ecs_task_definition.api.arn
desired_count = 2
launch_type = "FARGATE"
network_configuration {
subnets = aws_subnet.private[*].id
security_groups = [aws_security_group.ecs.id]
}
load_balancer {
target_group_arn = aws_lb_target_group.api.arn
container_name = "api"
container_port = 3000
}
}
# ECS Service (Worker)
resource "aws_ecs_service" "worker" {
name = "worker-service"
cluster = aws_ecs_cluster.main.id
task_definition = aws_ecs_task_definition.worker.arn
desired_count = 2
launch_type = "FARGATE"
network_configuration {
subnets = aws_subnet.private[*].id
security_groups = [aws_security_group.ecs.id]
}
}
# Auto Scaling for Worker
resource "aws_appautoscaling_target" "worker" {
max_capacity = 10
min_capacity = 1
resource_id = "service/${aws_ecs_cluster.main.name}/${aws_ecs_service.worker.name}"
scalable_dimension = "ecs:service:DesiredCount"
service_namespace = "ecs"
}
resource "aws_appautoscaling_policy" "worker_cpu" {
name = "worker-cpu-scaling"
policy_type = "TargetTrackingScaling"
resource_id = aws_appautoscaling_target.worker.resource_id
scalable_dimension = aws_appautoscaling_target.worker.scalable_dimension
service_namespace = aws_appautoscaling_target.worker.service_namespace
target_tracking_scaling_policy_configuration {
target_value = 70.0
predefined_metric_specification {
predefined_metric_type = "ECSServiceAverageCPUUtilization"
}
scale_in_cooldown = 300
scale_out_cooldown = 60
}
}Google Cloud: Terraform (Cloud Run)
# Variables
variable "project_id" {
description = "GCP Project ID"
}
variable "region" {
default = "asia-northeast1"
}
# Memorystore Redis
resource "google_redis_instance" "bullmq" {
name = "bullmq-redis"
tier = "STANDARD_HA"
memory_size_gb = 1
region = var.region
authorized_network = google_compute_network.main.id
redis_version = "REDIS_7_0"
display_name = "BullMQ Redis Instance"
transit_encryption_mode = "SERVER_AUTHENTICATION"
}
# Cloud Run Service (API)
resource "google_cloud_run_v2_service" "api" {
name = "api"
location = var.region
template {
containers {
image = "${var.region}-docker.pkg.dev/${var.project_id}/app/api:latest"
env {
name = "NODE_ENV"
value = "production"
}
env {
name = "REDIS_HOST"
value = google_redis_instance.bullmq.host
}
env {
name = "REDIS_PORT"
value = tostring(google_redis_instance.bullmq.port)
}
ports {
container_port = 3000
}
resources {
limits = {
cpu = "1"
memory = "512Mi"
}
}
startup_probe {
http_get {
path = "/health"
port = 3000
}
initial_delay_seconds = 10
period_seconds = 10
failure_threshold = 3
}
liveness_probe {
http_get {
path = "/health"
port = 3000
}
period_seconds = 30
}
}
vpc_access {
network_interfaces {
network = google_compute_network.main.name
subnetwork = google_compute_subnetwork.main.name
}
egress = "PRIVATE_RANGES_ONLY"
}
scaling {
min_instance_count = 1
max_instance_count = 10
}
}
traffic {
percent = 100
type = "TRAFFIC_TARGET_ALLOCATION_TYPE_LATEST"
}
}
# Cloud Run Service (Worker)
resource "google_cloud_run_v2_service" "worker" {
name = "worker"
location = var.region
template {
containers {
image = "${var.region}-docker.pkg.dev/${var.project_id}/app/api:latest"
command = ["node", "dist/worker.js"]
env {
name = "NODE_ENV"
value = "production"
}
env {
name = "REDIS_HOST"
value = google_redis_instance.bullmq.host
}
env {
name = "REDIS_PORT"
value = tostring(google_redis_instance.bullmq.port)
}
resources {
limits = {
cpu = "1"
memory = "512Mi"
}
cpu_idle = false # ワーカーは常時起動
}
}
vpc_access {
network_interfaces {
network = google_compute_network.main.name
subnetwork = google_compute_subnetwork.main.name
}
egress = "PRIVATE_RANGES_ONLY"
}
scaling {
min_instance_count = 1 # ワーカーは常時1台以上
max_instance_count = 5
}
}
# ワーカーは外部からのアクセス不要
ingress = "INGRESS_TRAFFIC_INTERNAL_ONLY"
}
# IAM for Cloud Run
resource "google_cloud_run_service_iam_member" "api_invoker" {
location = google_cloud_run_v2_service.api.location
service = google_cloud_run_v2_service.api.name
role = "roles/run.invoker"
member = "allUsers"
}Bull Board(管理UI)の追加
// src/admin/bullBoard.ts
import { createBullBoard } from '@bull-board/api';
import { BullMQAdapter } from '@bull-board/api/bullMQAdapter';
import { ExpressAdapter } from '@bull-board/express';
import { emailQueue } from '../queues/emailQueue';
export const setupBullBoard = (app: Express) => {
const serverAdapter = new ExpressAdapter();
serverAdapter.setBasePath('/admin/queues');
createBullBoard({
queues: [new BullMQAdapter(emailQueue)],
serverAdapter,
});
// Basic認証などで保護することを推奨
app.use('/admin/queues', serverAdapter.getRouter());
};フロントエンドへの状態通知(Socket.io)
非同期ジョブの進捗や完了状態をフロントエンドにリアルタイムで通知するには、Socket.ioを使用します。
パッケージインストール
npm install socket.io
npm install -D @types/socket.ioSocket.ioサーバーの設定
// src/socket/server.ts
import { Server as HttpServer } from 'http';
import { Server, Socket } from 'socket.io';
import { logger } from '../utils/logger';
let io: Server;
export const initSocketServer = (httpServer: HttpServer): Server => {
io = new Server(httpServer, {
cors: {
origin: process.env.FRONTEND_URL || 'http://localhost:3000',
methods: ['GET', 'POST'],
},
path: '/socket.io',
});
io.on('connection', (socket: Socket) => {
logger.info(`Client connected: ${socket.id}`);
// ジョブ状態のサブスクリプション
socket.on('subscribe:job', (jobId: string) => {
socket.join(`job:${jobId}`);
logger.info(`Socket ${socket.id} subscribed to job:${jobId}`);
});
socket.on('unsubscribe:job', (jobId: string) => {
socket.leave(`job:${jobId}`);
logger.info(`Socket ${socket.id} unsubscribed from job:${jobId}`);
});
socket.on('disconnect', () => {
logger.info(`Client disconnected: ${socket.id}`);
});
});
return io;
};
export const getSocketServer = (): Server => {
if (!io) {
throw new Error('Socket.io server not initialized');
}
return io;
};ジョブ状態の通知ヘルパー
// src/socket/jobNotifier.ts
import { getSocketServer } from './server';
export interface JobStatusEvent {
jobId: string;
status: 'pending' | 'processing' | 'completed' | 'failed';
progress?: number;
result?: Record<string, unknown>;
error?: string;
updatedAt: string;
}
export const notifyJobStatus = (event: JobStatusEvent): void => {
const io = getSocketServer();
io.to(`job:${event.jobId}`).emit('job:status', event);
};Worker での通知
// src/workers/emailWorker.ts
import { Worker, Job } from 'bullmq';
import { createRedisConnection } from '../queues/connection';
import { EmailJobData, EmailJobResult, EMAIL_QUEUE_NAME } from '../queues/types';
import { notifyJobStatus } from '../socket/jobNotifier';
import { sendEmail } from '../services/emailService';
import { logger } from '../utils/logger';
const processEmail = async (
job: Job<EmailJobData, EmailJobResult>
): Promise<EmailJobResult> => {
const { to, subject, body, templateId } = job.data;
logger.info(`Processing email job ${job.id}`, { to, subject });
// 処理開始を通知
notifyJobStatus({
jobId: job.id!,
status: 'processing',
progress: 0,
updatedAt: new Date().toISOString(),
});
try {
// 進捗更新
await job.updateProgress(10);
notifyJobStatus({
jobId: job.id!,
status: 'processing',
progress: 10,
updatedAt: new Date().toISOString(),
});
const result = await sendEmail({ to, subject, body, templateId });
await job.updateProgress(100);
const jobResult = {
messageId: result.messageId,
sentAt: new Date(),
};
// 完了を通知
notifyJobStatus({
jobId: job.id!,
status: 'completed',
progress: 100,
result: jobResult,
updatedAt: new Date().toISOString(),
});
return jobResult;
} catch (error) {
const errorMessage = error instanceof Error ? error.message : 'Unknown error';
logger.error(`Failed to send email`, { jobId: job.id, error });
// エラーを通知
notifyJobStatus({
jobId: job.id!,
status: 'failed',
error: errorMessage,
updatedAt: new Date().toISOString(),
});
throw error;
}
};
export const createEmailWorker = () => {
const worker = new Worker<EmailJobData, EmailJobResult>(
EMAIL_QUEUE_NAME,
processEmail,
{
connection: createRedisConnection(),
concurrency: 10,
}
);
return worker;
};Express サーバーへの統合
// src/server.ts
import express from 'express';
import { createServer } from 'http';
import { initSocketServer } from './socket/server';
const app = express();
const httpServer = createServer(app);
// Socket.io初期化
initSocketServer(httpServer);
// ... 既存のルート設定 ...
const PORT = process.env.PORT || 3000;
httpServer.listen(PORT, () => {
console.log(`Server listening on port ${PORT}`);
});フロントエンド(React)
// frontend/src/lib/socket.ts
import { io, Socket } from 'socket.io-client';
const SOCKET_URL = process.env.NEXT_PUBLIC_API_URL || 'http://localhost:3000';
let socket: Socket | null = null;
export const getSocket = (): Socket => {
if (!socket) {
socket = io(SOCKET_URL, {
path: '/socket.io',
transports: ['websocket', 'polling'],
});
}
return socket;
};
export const disconnectSocket = (): void => {
if (socket) {
socket.disconnect();
socket = null;
}
};React Hook
// frontend/src/hooks/useJobStatus.ts
import { useState, useEffect, useCallback } from 'react';
import { getSocket } from '../lib/socket';
interface JobStatus {
status: 'pending' | 'processing' | 'completed' | 'failed';
progress?: number;
result?: Record<string, unknown>;
error?: string;
}
export const useJobStatus = (jobId: string | null) => {
const [status, setStatus] = useState<JobStatus | null>(null);
useEffect(() => {
if (!jobId) return;
const socket = getSocket();
// ジョブにサブスクライブ
socket.emit('subscribe:job', jobId);
// 状態更新を受信
const handleStatus = (event: JobStatus & { jobId: string }) => {
if (event.jobId === jobId) {
setStatus({
status: event.status,
progress: event.progress,
result: event.result,
error: event.error,
});
}
};
socket.on('job:status', handleStatus);
return () => {
socket.emit('unsubscribe:job', jobId);
socket.off('job:status', handleStatus);
};
}, [jobId]);
return status;
};React Component
// frontend/src/components/JobProgress.tsx
import { useJobStatus } from '../hooks/useJobStatus';
interface JobProgressProps {
jobId: string;
onComplete?: (result: Record<string, unknown>) => void;
onError?: (error: string) => void;
}
export const JobProgress: React.FC<JobProgressProps> = ({
jobId,
onComplete,
onError,
}) => {
const status = useJobStatus(jobId);
useEffect(() => {
if (status?.status === 'completed' && status.result) {
onComplete?.(status.result);
}
if (status?.status === 'failed' && status.error) {
onError?.(status.error);
}
}, [status, onComplete, onError]);
if (!status) {
return <div>接続中...</div>;
}
return (
<div className="job-progress">
<div className="status">状態: {status.status}</div>
{status.progress !== undefined && (
<div className="progress-bar">
<div
className="progress-fill"
style={{ width: `${status.progress}%` }}
/>
<span>{status.progress}%</span>
</div>
)}
{status.status === 'failed' && (
<div className="error">エラー: {status.error}</div>
)}
{status.status === 'completed' && (
<div className="result">完了しました</div>
)}
</div>
);
};Server-Sent Events(SSE)による代替実装
Socket.ioより軽量なSSEを使用する場合の実装例です。
// src/routes/sse.ts
import { Router, Request, Response } from 'express';
import { emailQueue } from '../queues/emailQueue';
import { QueueEvents } from 'bullmq';
import { createRedisConnection } from '../queues/connection';
const router = Router();
// SSEエンドポイント
router.get('/jobs/:jobId/stream', async (req: Request, res: Response) => {
const { jobId } = req.params;
res.setHeader('Content-Type', 'text/event-stream');
res.setHeader('Cache-Control', 'no-cache');
res.setHeader('Connection', 'keep-alive');
res.setHeader('X-Accel-Buffering', 'no'); // nginx用
const queueEvents = new QueueEvents(emailQueue.name, {
connection: createRedisConnection(),
});
const sendEvent = (data: object) => {
res.write(`data: ${JSON.stringify(data)}\n\n`);
};
// 進捗イベント
queueEvents.on('progress', ({ jobId: eventJobId, data }) => {
if (eventJobId === jobId) {
sendEvent({ type: 'progress', progress: data });
}
});
// 完了イベント
queueEvents.on('completed', ({ jobId: eventJobId, returnvalue }) => {
if (eventJobId === jobId) {
sendEvent({ type: 'completed', result: returnvalue });
cleanup();
}
});
// 失敗イベント
queueEvents.on('failed', ({ jobId: eventJobId, failedReason }) => {
if (eventJobId === jobId) {
sendEvent({ type: 'failed', error: failedReason });
cleanup();
}
});
const cleanup = () => {
queueEvents.close();
res.end();
};
// クライアント切断時のクリーンアップ
req.on('close', cleanup);
});
export default router;// frontend/src/hooks/useJobStatusSSE.ts
import { useState, useEffect } from 'react';
interface JobStatus {
status: 'pending' | 'processing' | 'completed' | 'failed';
progress?: number;
result?: Record<string, unknown>;
error?: string;
}
export const useJobStatusSSE = (jobId: string | null) => {
const [status, setStatus] = useState<JobStatus>({ status: 'pending' });
useEffect(() => {
if (!jobId) return;
const eventSource = new EventSource(
`${process.env.NEXT_PUBLIC_API_URL}/jobs/${jobId}/stream`
);
eventSource.onmessage = (event) => {
const data = JSON.parse(event.data);
switch (data.type) {
case 'progress':
setStatus((prev) => ({
...prev,
status: 'processing',
progress: data.progress,
}));
break;
case 'completed':
setStatus({
status: 'completed',
progress: 100,
result: data.result,
});
eventSource.close();
break;
case 'failed':
setStatus({
status: 'failed',
error: data.error,
});
eventSource.close();
break;
}
};
eventSource.onerror = () => {
eventSource.close();
};
return () => {
eventSource.close();
};
}, [jobId]);
return status;
};インフラ設定の注意点
AWS (ALB + ECS)
WebSocket/SSE接続を維持するため、ALBのアイドルタイムアウトを設定します。
resource "aws_lb" "api" {
# ... 既存の設定 ...
# WebSocket/SSE用にアイドルタイムアウトを延長
idle_timeout = 3600
}
resource "aws_lb_target_group" "api" {
# ... 既存の設定 ...
# スティッキーセッション(Socket.io polling fallback用)
stickiness {
type = "lb_cookie"
cookie_duration = 86400
enabled = true
}
}Google Cloud (Cloud Run)
Cloud RunではWebSocket接続は最大60分まで維持されます。
resource "google_cloud_run_v2_service" "api" {
template {
# ... 既存の設定 ...
# セッションアフィニティを有効化
session_affinity = true
}
}監視とアラート
CloudWatch メトリクス(AWS)
resource "aws_cloudwatch_metric_alarm" "worker_cpu_high" {
alarm_name = "worker-cpu-high"
comparison_operator = "GreaterThanThreshold"
evaluation_periods = 2
metric_name = "CPUUtilization"
namespace = "AWS/ECS"
period = 300
statistic = "Average"
threshold = 80
alarm_description = "Worker CPU utilization is high"
alarm_actions = [aws_sns_topic.alerts.arn]
dimensions = {
ClusterName = aws_ecs_cluster.main.name
ServiceName = aws_ecs_service.worker.name
}
}アプリケーションメトリクス
// src/metrics/queueMetrics.ts
import { emailQueue } from '../queues/emailQueue';
export const getQueueMetrics = async () => {
const [waiting, active, completed, failed, delayed] = await Promise.all([
emailQueue.getWaitingCount(),
emailQueue.getActiveCount(),
emailQueue.getCompletedCount(),
emailQueue.getFailedCount(),
emailQueue.getDelayedCount(),
]);
return {
waiting,
active,
completed,
failed,
delayed,
};
};AWS SQS + AppSync Events + DynamoDBによる実装
AWSのマネージドサービスを組み合わせて、完全なサーバーレスアーキテクチャで非同期処理を実現します。SQSでタスクキューイング、Lambdaでジョブ実行、AppSync Eventsでリアルタイム通知、DynamoDBで状態管理を行います。
アーキテクチャの特徴
- フルマネージド: インフラ管理が不要、スケーリングも自動
- 従量課金: 実行した分だけの課金で、アイドル時のコストが低い
- 高い耐障害性: AWSのマネージドサービスによるSLA保証
- リアルタイム通知: AppSync EventsによるWebSocket通信でフロントエンドにプッシュ通知
採用基準
この構成を選ぶべきケース
- サーバーレスアーキテクチャを採用している
- インフラ管理の負担を最小化したい
- 負荷の変動が大きく、オートスケーリングが必要
- ジョブの実行頻度が不定期で、アイドル時のコストを抑えたい
- AWSに集約したインフラ構成を維持したい
他の選択肢を検討すべきケース
- 処理時間が15分を超える場合 → Step Functions + ECS/Fargate
- 低レイテンシが必要な場合 → BullMQ/Sidekiq
- 複雑なジョブ依存関係がある場合 → Step Functions
- マルチクラウド/ベンダーロックイン回避が必要 → BullMQ
アーキテクチャ
全体構成
┌──────────────┐ ┌──────────────┐ ┌──────────────┐
│ API Gateway │────▶│ Lambda │────▶│ SQS │
│ (REST/HTTP) │ │ (Producer) │ │ (Queue) │
└──────────────┘ └──────────────┘ └──────────────┘
│
▼
┌──────────────┐ ┌──────────────┐ ┌──────────────┐
│ Client │◀────│ AppSync │◀────│ Lambda │
│ (Frontend) │ WS │ Events │ │ (Worker) │
└──────────────┘ └──────────────┘ └──────────────┘
│
▼
┌──────────────┐
│ DynamoDB │
│ (状態管理) │
└──────────────┘コンポーネントの役割
- API Gateway: リクエストの受付、認証
- Producer Lambda: ジョブの登録、SQSへのメッセージ送信
- SQS: メッセージキュー、リトライ、DLQ
- Worker Lambda: SQSトリガーでジョブを処理
- DynamoDB: ジョブの状態管理、メタデータ保存
- AppSync Events: フロントエンドへのリアルタイム通知
実装サンプル
プロジェクト構成
src/
├── functions/
│ ├── enqueue/ # ジョブ投入Lambda
│ │ └── handler.ts
│ ├── worker/ # ワーカーLambda
│ │ └── handler.ts
│ └── status/ # 状態確認Lambda
│ └── handler.ts
├── lib/
│ ├── dynamodb.ts # DynamoDB操作
│ ├── sqs.ts # SQS操作
│ └── appsync.ts # AppSync操作
└── types/
└── job.ts # 型定義型定義
// src/types/job.ts
export interface Job {
jobId: string;
type: string;
status: 'pending' | 'processing' | 'completed' | 'failed';
payload: Record<string, unknown>;
result?: Record<string, unknown>;
error?: string;
createdAt: string;
updatedAt: string;
completedAt?: string;
}
export interface EnqueueRequest {
type: string;
payload: Record<string, unknown>;
}
export interface JobStatusEvent {
jobId: string;
status: Job['status'];
progress?: number;
result?: Record<string, unknown>;
error?: string;
}DynamoDB操作
// src/lib/dynamodb.ts
import { DynamoDBClient } from '@aws-sdk/client-dynamodb';
import {
DynamoDBDocumentClient,
PutCommand,
GetCommand,
UpdateCommand,
} from '@aws-sdk/lib-dynamodb';
import { Job } from '../types/job';
const client = new DynamoDBClient({});
const docClient = DynamoDBDocumentClient.from(client);
const TABLE_NAME = process.env.JOBS_TABLE_NAME!;
export const createJob = async (job: Job): Promise<void> => {
await docClient.send(
new PutCommand({
TableName: TABLE_NAME,
Item: job,
})
);
};
export const getJob = async (jobId: string): Promise<Job | null> => {
const result = await docClient.send(
new GetCommand({
TableName: TABLE_NAME,
Key: { jobId },
})
);
return (result.Item as Job) || null;
};
export const updateJobStatus = async (
jobId: string,
status: Job['status'],
updates?: Partial<Pick<Job, 'result' | 'error' | 'completedAt'>>
): Promise<void> => {
const updateExpressions: string[] = [
'#status = :status',
'updatedAt = :updatedAt',
];
const expressionAttributeNames: Record<string, string> = {
'#status': 'status',
};
const expressionAttributeValues: Record<string, unknown> = {
':status': status,
':updatedAt': new Date().toISOString(),
};
if (updates?.result) {
updateExpressions.push('#result = :result');
expressionAttributeNames['#result'] = 'result';
expressionAttributeValues[':result'] = updates.result;
}
if (updates?.error) {
updateExpressions.push('#error = :error');
expressionAttributeNames['#error'] = 'error';
expressionAttributeValues[':error'] = updates.error;
}
if (updates?.completedAt) {
updateExpressions.push('completedAt = :completedAt');
expressionAttributeValues[':completedAt'] = updates.completedAt;
}
await docClient.send(
new UpdateCommand({
TableName: TABLE_NAME,
Key: { jobId },
UpdateExpression: `SET ${updateExpressions.join(', ')}`,
ExpressionAttributeNames: expressionAttributeNames,
ExpressionAttributeValues: expressionAttributeValues,
})
);
};SQS操作
// src/lib/sqs.ts
import { SQSClient, SendMessageCommand } from '@aws-sdk/client-sqs';
const client = new SQSClient({});
const QUEUE_URL = process.env.QUEUE_URL!;
export interface QueueMessage {
jobId: string;
type: string;
payload: Record<string, unknown>;
}
export const sendMessage = async (message: QueueMessage): Promise<void> => {
await client.send(
new SendMessageCommand({
QueueUrl: QUEUE_URL,
MessageBody: JSON.stringify(message),
MessageAttributes: {
jobType: {
DataType: 'String',
StringValue: message.type,
},
},
})
);
};AppSync Events操作
// src/lib/appsync.ts
import { SignatureV4 } from '@smithy/signature-v4';
import { Sha256 } from '@aws-crypto/sha256-js';
import { defaultProvider } from '@aws-sdk/credential-provider-node';
import { HttpRequest } from '@smithy/protocol-http';
import { JobStatusEvent } from '../types/job';
const APPSYNC_ENDPOINT = process.env.APPSYNC_HTTP_ENDPOINT!;
const APPSYNC_REGION = process.env.AWS_REGION!;
export const publishJobStatus = async (event: JobStatusEvent): Promise<void> => {
const url = new URL(APPSYNC_ENDPOINT);
const request = new HttpRequest({
method: 'POST',
hostname: url.hostname,
path: url.pathname,
headers: {
'Content-Type': 'application/json',
host: url.hostname,
},
body: JSON.stringify({
channel: `job/${event.jobId}`,
events: [JSON.stringify(event)],
}),
});
const signer = new SignatureV4({
credentials: defaultProvider(),
region: APPSYNC_REGION,
service: 'appsync',
sha256: Sha256,
});
const signedRequest = await signer.sign(request);
const response = await fetch(APPSYNC_ENDPOINT, {
method: signedRequest.method,
headers: signedRequest.headers as HeadersInit,
body: signedRequest.body,
});
if (!response.ok) {
throw new Error(`Failed to publish event: ${response.statusText}`);
}
};Producer Lambda
// src/functions/enqueue/handler.ts
import { APIGatewayProxyHandler } from 'aws-lambda';
import { randomUUID } from 'crypto';
import { createJob } from '../../lib/dynamodb';
import { sendMessage } from '../../lib/sqs';
import { EnqueueRequest, Job } from '../../types/job';
export const handler: APIGatewayProxyHandler = async (event) => {
try {
const body: EnqueueRequest = JSON.parse(event.body || '{}');
if (!body.type || !body.payload) {
return {
statusCode: 400,
body: JSON.stringify({ error: 'type and payload are required' }),
};
}
const jobId = randomUUID();
const now = new Date().toISOString();
const job: Job = {
jobId,
type: body.type,
status: 'pending',
payload: body.payload,
createdAt: now,
updatedAt: now,
};
// DynamoDBにジョブを作成
await createJob(job);
// SQSにメッセージを送信
await sendMessage({
jobId,
type: body.type,
payload: body.payload,
});
return {
statusCode: 202,
headers: {
'Content-Type': 'application/json',
},
body: JSON.stringify({
jobId,
status: 'pending',
}),
};
} catch (error) {
console.error('Error enqueueing job:', error);
return {
statusCode: 500,
body: JSON.stringify({ error: 'Internal server error' }),
};
}
};Worker Lambda
// src/functions/worker/handler.ts
import { SQSHandler, SQSRecord } from 'aws-lambda';
import { updateJobStatus, getJob } from '../../lib/dynamodb';
import { publishJobStatus } from '../../lib/appsync';
import { QueueMessage } from '../../lib/sqs';
// ジョブタイプごとの処理関数
const jobProcessors: Record<string, (payload: Record<string, unknown>) => Promise<Record<string, unknown>>> = {
'send-email': async (payload) => {
// メール送信ロジック
const { to, subject, body } = payload as { to: string; subject: string; body: string };
// 実際のメール送信処理...
return { messageId: `msg-${Date.now()}`, sentAt: new Date().toISOString() };
},
'process-image': async (payload) => {
// 画像処理ロジック
const { imageUrl } = payload as { imageUrl: string };
// 実際の画像処理...
return { processedUrl: `${imageUrl}-processed` };
},
};
const processRecord = async (record: SQSRecord): Promise<void> => {
const message: QueueMessage = JSON.parse(record.body);
const { jobId, type, payload } = message;
console.log(`Processing job ${jobId} of type ${type}`);
try {
// 状態を processing に更新
await updateJobStatus(jobId, 'processing');
await publishJobStatus({ jobId, status: 'processing', progress: 0 });
// ジョブタイプに応じた処理を実行
const processor = jobProcessors[type];
if (!processor) {
throw new Error(`Unknown job type: ${type}`);
}
const result = await processor(payload);
// 状態を completed に更新
await updateJobStatus(jobId, 'completed', {
result,
completedAt: new Date().toISOString(),
});
await publishJobStatus({ jobId, status: 'completed', progress: 100, result });
console.log(`Job ${jobId} completed successfully`);
} catch (error) {
const errorMessage = error instanceof Error ? error.message : 'Unknown error';
console.error(`Job ${jobId} failed:`, errorMessage);
// 状態を failed に更新
await updateJobStatus(jobId, 'failed', { error: errorMessage });
await publishJobStatus({ jobId, status: 'failed', error: errorMessage });
// エラーを再スローしてSQSにリトライさせる
throw error;
}
};
export const handler: SQSHandler = async (event) => {
const results = await Promise.allSettled(
event.Records.map(processRecord)
);
// 失敗したレコードを報告(部分的なバッチ失敗)
const failedRecords = results
.map((result, index) => (result.status === 'rejected' ? event.Records[index] : null))
.filter((record): record is SQSRecord => record !== null);
if (failedRecords.length > 0) {
return {
batchItemFailures: failedRecords.map((record) => ({
itemIdentifier: record.messageId,
})),
};
}
};Status Lambda
// src/functions/status/handler.ts
import { APIGatewayProxyHandler } from 'aws-lambda';
import { getJob } from '../../lib/dynamodb';
export const handler: APIGatewayProxyHandler = async (event) => {
const jobId = event.pathParameters?.jobId;
if (!jobId) {
return {
statusCode: 400,
body: JSON.stringify({ error: 'jobId is required' }),
};
}
const job = await getJob(jobId);
if (!job) {
return {
statusCode: 404,
body: JSON.stringify({ error: 'Job not found' }),
};
}
return {
statusCode: 200,
headers: {
'Content-Type': 'application/json',
},
body: JSON.stringify(job),
};
};Terraform
# Variables
variable "environment" {
default = "production"
}
variable "region" {
default = "ap-northeast-1"
}
# DynamoDB Table
resource "aws_dynamodb_table" "jobs" {
name = "${var.environment}-async-jobs"
billing_mode = "PAY_PER_REQUEST"
hash_key = "jobId"
attribute {
name = "jobId"
type = "S"
}
ttl {
attribute_name = "ttl"
enabled = true
}
tags = {
Environment = var.environment
}
}
# SQS Queue
resource "aws_sqs_queue" "jobs" {
name = "${var.environment}-async-jobs"
visibility_timeout_seconds = 300 # Lambda timeout + buffer
message_retention_seconds = 1209600 # 14 days
receive_wait_time_seconds = 20 # Long polling
redrive_policy = jsonencode({
deadLetterTargetArn = aws_sqs_queue.jobs_dlq.arn
maxReceiveCount = 3
})
tags = {
Environment = var.environment
}
}
# Dead Letter Queue
resource "aws_sqs_queue" "jobs_dlq" {
name = "${var.environment}-async-jobs-dlq"
message_retention_seconds = 1209600
tags = {
Environment = var.environment
}
}
# AppSync Events API
# 注意: aws_appsync_api リソースは Event API 専用です(api_type パラメータは不要)
resource "aws_appsync_api" "events" {
name = "${var.environment}-job-events"
event_config {
auth_provider {
auth_type = "AWS_IAM"
}
connection_auth_mode {
auth_type = "AWS_IAM"
}
default_publish_auth_mode {
auth_type = "AWS_IAM"
}
default_subscribe_auth_mode {
auth_type = "AWS_IAM"
}
}
}
resource "aws_appsync_channel_namespace" "jobs" {
api_id = aws_appsync_api.events.id
name = "job"
publish_auth_modes {
auth_type = "AWS_IAM"
}
subscribe_auth_modes {
auth_type = "AWS_IAM"
}
}
# Lambda IAM Role
resource "aws_iam_role" "lambda" {
name = "${var.environment}-async-lambda-role"
assume_role_policy = jsonencode({
Version = "2012-10-17"
Statement = [{
Action = "sts:AssumeRole"
Effect = "Allow"
Principal = {
Service = "lambda.amazonaws.com"
}
}]
})
}
resource "aws_iam_role_policy" "lambda" {
name = "${var.environment}-async-lambda-policy"
role = aws_iam_role.lambda.id
policy = jsonencode({
Version = "2012-10-17"
Statement = [
{
Effect = "Allow"
Action = [
"dynamodb:GetItem",
"dynamodb:PutItem",
"dynamodb:UpdateItem",
"dynamodb:Query"
]
Resource = aws_dynamodb_table.jobs.arn
},
{
Effect = "Allow"
Action = [
"sqs:SendMessage",
"sqs:ReceiveMessage",
"sqs:DeleteMessage",
"sqs:GetQueueAttributes"
]
Resource = aws_sqs_queue.jobs.arn
},
{
Effect = "Allow"
Action = [
"appsync:EventPublish"
]
Resource = "${aws_appsync_api.events.arn}/*"
},
{
Effect = "Allow"
Action = [
"logs:CreateLogGroup",
"logs:CreateLogStream",
"logs:PutLogEvents"
]
Resource = "arn:aws:logs:*:*:*"
}
]
})
}
# Lambda Functions
resource "aws_lambda_function" "enqueue" {
function_name = "${var.environment}-async-enqueue"
runtime = "nodejs20.x"
handler = "handler.handler"
role = aws_iam_role.lambda.arn
timeout = 30
memory_size = 256
filename = "dist/enqueue.zip"
source_code_hash = filebase64sha256("dist/enqueue.zip")
environment {
variables = {
JOBS_TABLE_NAME = aws_dynamodb_table.jobs.name
QUEUE_URL = aws_sqs_queue.jobs.url
APPSYNC_HTTP_ENDPOINT = "https://${aws_appsync_api.events.id}.appsync-api.${var.region}.amazonaws.com/event"
}
}
}
resource "aws_lambda_function" "worker" {
function_name = "${var.environment}-async-worker"
runtime = "nodejs20.x"
handler = "handler.handler"
role = aws_iam_role.lambda.arn
timeout = 300 # 5 minutes max
memory_size = 512
filename = "dist/worker.zip"
source_code_hash = filebase64sha256("dist/worker.zip")
environment {
variables = {
JOBS_TABLE_NAME = aws_dynamodb_table.jobs.name
APPSYNC_HTTP_ENDPOINT = "https://${aws_appsync_api.events.id}.appsync-api.${var.region}.amazonaws.com/event"
}
}
}
resource "aws_lambda_function" "status" {
function_name = "${var.environment}-async-status"
runtime = "nodejs20.x"
handler = "handler.handler"
role = aws_iam_role.lambda.arn
timeout = 30
memory_size = 256
filename = "dist/status.zip"
source_code_hash = filebase64sha256("dist/status.zip")
environment {
variables = {
JOBS_TABLE_NAME = aws_dynamodb_table.jobs.name
}
}
}
# SQS Event Source Mapping
resource "aws_lambda_event_source_mapping" "worker_sqs" {
event_source_arn = aws_sqs_queue.jobs.arn
function_name = aws_lambda_function.worker.arn
batch_size = 10
maximum_batching_window_in_seconds = 5
function_response_types = ["ReportBatchItemFailures"]
}
# API Gateway
resource "aws_apigatewayv2_api" "api" {
name = "${var.environment}-async-api"
protocol_type = "HTTP"
}
resource "aws_apigatewayv2_stage" "api" {
api_id = aws_apigatewayv2_api.api.id
name = "$default"
auto_deploy = true
}
resource "aws_apigatewayv2_integration" "enqueue" {
api_id = aws_apigatewayv2_api.api.id
integration_type = "AWS_PROXY"
integration_uri = aws_lambda_function.enqueue.invoke_arn
payload_format_version = "2.0"
}
resource "aws_apigatewayv2_integration" "status" {
api_id = aws_apigatewayv2_api.api.id
integration_type = "AWS_PROXY"
integration_uri = aws_lambda_function.status.invoke_arn
payload_format_version = "2.0"
}
resource "aws_apigatewayv2_route" "enqueue" {
api_id = aws_apigatewayv2_api.api.id
route_key = "POST /jobs"
target = "integrations/${aws_apigatewayv2_integration.enqueue.id}"
}
resource "aws_apigatewayv2_route" "status" {
api_id = aws_apigatewayv2_api.api.id
route_key = "GET /jobs/{jobId}"
target = "integrations/${aws_apigatewayv2_integration.status.id}"
}
resource "aws_lambda_permission" "enqueue" {
action = "lambda:InvokeFunction"
function_name = aws_lambda_function.enqueue.function_name
principal = "apigateway.amazonaws.com"
source_arn = "${aws_apigatewayv2_api.api.execution_arn}/*/*"
}
resource "aws_lambda_permission" "status" {
action = "lambda:InvokeFunction"
function_name = aws_lambda_function.status.function_name
principal = "apigateway.amazonaws.com"
source_arn = "${aws_apigatewayv2_api.api.execution_arn}/*/*"
}
# Outputs
output "api_endpoint" {
value = aws_apigatewayv2_api.api.api_endpoint
}
output "appsync_realtime_endpoint" {
value = aws_appsync_api.events.realtime_uris["REALTIME"]
}
output "appsync_http_endpoint" {
value = "https://${aws_appsync_api.events.id}.appsync-api.${var.region}.amazonaws.com/event"
}フロントエンド連携(AppSync Events)
// frontend/src/lib/jobSubscription.ts
import { Amplify } from 'aws-amplify';
import { events } from 'aws-amplify/data';
// Amplify設定
Amplify.configure({
API: {
Events: {
endpoint: process.env.NEXT_PUBLIC_APPSYNC_ENDPOINT!,
region: process.env.NEXT_PUBLIC_AWS_REGION!,
defaultAuthMode: 'iam',
},
},
});
interface JobStatusEvent {
jobId: string;
status: 'pending' | 'processing' | 'completed' | 'failed';
progress?: number;
result?: Record<string, unknown>;
error?: string;
}
export const subscribeToJobStatus = (
jobId: string,
onStatusChange: (event: JobStatusEvent) => void
) => {
const channel = events.connect(`job/${jobId}`);
channel.subscribe({
next: (event) => {
const data = JSON.parse(event.data as string) as JobStatusEvent;
onStatusChange(data);
},
error: (error) => {
console.error('Subscription error:', error);
},
});
return () => {
channel.close();
};
};React Hook
// frontend/src/hooks/useJobStatus.ts
import { useState, useEffect, useCallback } from 'react';
import { subscribeToJobStatus } from '../lib/jobSubscription';
interface JobStatus {
status: 'pending' | 'processing' | 'completed' | 'failed';
progress?: number;
result?: Record<string, unknown>;
error?: string;
}
export const useJobStatus = (jobId: string | null) => {
const [status, setStatus] = useState<JobStatus | null>(null);
useEffect(() => {
if (!jobId) return;
const unsubscribe = subscribeToJobStatus(jobId, (event) => {
setStatus({
status: event.status,
progress: event.progress,
result: event.result,
error: event.error,
});
});
return unsubscribe;
}, [jobId]);
return status;
};監視とアラート
# CloudWatch Alarms
resource "aws_cloudwatch_metric_alarm" "dlq_messages" {
alarm_name = "${var.environment}-async-dlq-messages"
comparison_operator = "GreaterThanThreshold"
evaluation_periods = 1
metric_name = "ApproximateNumberOfMessagesVisible"
namespace = "AWS/SQS"
period = 300
statistic = "Sum"
threshold = 0
alarm_description = "DLQ has messages - jobs are failing"
alarm_actions = [aws_sns_topic.alerts.arn]
dimensions = {
QueueName = aws_sqs_queue.jobs_dlq.name
}
}
resource "aws_cloudwatch_metric_alarm" "lambda_errors" {
alarm_name = "${var.environment}-worker-errors"
comparison_operator = "GreaterThanThreshold"
evaluation_periods = 2
metric_name = "Errors"
namespace = "AWS/Lambda"
period = 300
statistic = "Sum"
threshold = 5
alarm_description = "Worker Lambda is having errors"
alarm_actions = [aws_sns_topic.alerts.arn]
dimensions = {
FunctionName = aws_lambda_function.worker.function_name
}
}制限事項と注意点
Lambda制限
- 最大実行時間: 15分(長時間処理にはStep Functionsを検討)
- メモリ: 最大10GB
- 同時実行数: デフォルト1000(引き上げ可能)
SQS制限
- メッセージサイズ: 最大1MiB(2025年8月にAWSにより256KBから1MiBに拡張されました。大きなデータはS3に保存してURLを渡すことを推奨)
- メッセージ保持期間: 最大14日
DynamoDB制限
- アイテムサイズ: 最大400KB(属性名と値の合計)
- 書き込みスループット: オンデマンドモードでは自動スケーリング
コスト最適化
- Lambda: メモリと実行時間を最適化
- DynamoDB: TTLでデータを自動削除
- SQS: Long pollingで無駄なポーリングを削減
Google Cloud Tasks + Firestoreによる実装
Google Cloudのマネージドサービスを組み合わせて、サーバーレスアーキテクチャで非同期処理を実現します。Cloud Tasksでタスクキューイング、Cloud Runでジョブ実行、Firestoreでリアルタイム通知と状態管理を行います。
アーキテクチャの特徴
- フルマネージド: インフラ管理が不要
- リアルタイム同期: Firestoreのリアルタイムリスナーでフロントエンドに自動通知
- スケーラブル: Cloud RunとCloud Tasksの自動スケーリング
- HTTPベース: Cloud TasksはHTTPエンドポイントを直接呼び出すシンプルな設計
採用基準
この構成を選ぶべきケース
- Google Cloudをメインで利用している
- Firebase/Firestoreを既に利用している
- リアルタイム通知をシンプルに実現したい(Firestoreのリアルタイムリスナー活用)
- Cloud Run/Cloud Functionsでバックエンドを構築している
- HTTPベースのシンプルな設計を好む
他の選択肢を検討すべきケース
- 複雑なメッセージルーティングが必要 → Pub/Sub
- 超高スループット(数万件/秒)が必要 → Pub/Sub + Dataflow
- マルチクラウド環境 → BullMQ
- 詳細なジョブ管理UIが必要 → BullMQ + Bull Board
アーキテクチャ
全体構成
┌──────────────┐ ┌──────────────┐ ┌──────────────┐
│ Client │────▶│ Cloud Run │────▶│ Cloud Tasks │
│ (API Call) │ │ (Producer) │ │ (Queue) │
└──────────────┘ └──────────────┘ └──────────────┘
│ │
│ Realtime ▼
│ Listener ┌──────────────┐ ┌──────────────┐
└────────────────▶│ Firestore │◀│ Cloud Run │
│ (状態管理) │ │ (Worker) │
└──────────────┘ └──────────────┘コンポーネントの役割
- Cloud Run (Producer): APIリクエストを受け付け、Cloud Tasksにタスクを投入
- Cloud Tasks: タスクのキューイング、スケジューリング、リトライ管理
- Cloud Run (Worker): Cloud Tasksからのコールバックを受けてジョブを処理
- Firestore: ジョブの状態管理、リアルタイム通知の配信
実装サンプル
プロジェクト構成
src/
├── api/
│ ├── enqueue.ts # ジョブ投入エンドポイント
│ ├── worker.ts # ワーカーエンドポイント
│ └── status.ts # 状態確認エンドポイント
├── lib/
│ ├── firestore.ts # Firestore操作
│ ├── cloudTasks.ts # Cloud Tasks操作
│ └── auth.ts # 認証ヘルパー
└── types/
└── job.ts # 型定義型定義
// src/types/job.ts
export interface Job {
jobId: string;
type: string;
status: 'pending' | 'processing' | 'completed' | 'failed';
payload: Record<string, unknown>;
result?: Record<string, unknown>;
error?: string;
progress?: number;
createdAt: FirebaseFirestore.Timestamp;
updatedAt: FirebaseFirestore.Timestamp;
completedAt?: FirebaseFirestore.Timestamp;
}
export interface EnqueueRequest {
type: string;
payload: Record<string, unknown>;
scheduleTime?: string; // ISO8601形式の実行予定時刻
}
export interface TaskPayload {
jobId: string;
type: string;
payload: Record<string, unknown>;
}Firestore操作
// src/lib/firestore.ts
import { Firestore, FieldValue } from '@google-cloud/firestore';
import { Job } from '../types/job';
const firestore = new Firestore();
const JOBS_COLLECTION = 'async_jobs';
export const createJob = async (
job: Omit<Job, 'createdAt' | 'updatedAt'>
): Promise<void> => {
const now = FieldValue.serverTimestamp();
await firestore.collection(JOBS_COLLECTION).doc(job.jobId).set({
...job,
createdAt: now,
updatedAt: now,
});
};
export const getJob = async (jobId: string): Promise<Job | null> => {
const doc = await firestore.collection(JOBS_COLLECTION).doc(jobId).get();
if (!doc.exists) return null;
return { jobId: doc.id, ...doc.data() } as Job;
};
export const updateJobStatus = async (
jobId: string,
status: Job['status'],
updates?: Partial<Pick<Job, 'result' | 'error' | 'progress'>>
): Promise<void> => {
const data: Record<string, unknown> = {
status,
updatedAt: FieldValue.serverTimestamp(),
};
if (updates?.result !== undefined) data.result = updates.result;
if (updates?.error !== undefined) data.error = updates.error;
if (updates?.progress !== undefined) data.progress = updates.progress;
if (status === 'completed' || status === 'failed') {
data.completedAt = FieldValue.serverTimestamp();
}
await firestore.collection(JOBS_COLLECTION).doc(jobId).update(data);
};
// ジョブの自動削除(TTL)設定用
export const setJobTTL = async (
jobId: string,
expiresAt: Date
): Promise<void> => {
await firestore.collection(JOBS_COLLECTION).doc(jobId).update({
expiresAt: expiresAt,
});
};Cloud Tasks操作
// src/lib/cloudTasks.ts
import { CloudTasksClient } from '@google-cloud/tasks';
import { TaskPayload } from '../types/job';
const client = new CloudTasksClient();
const PROJECT_ID = process.env.GOOGLE_CLOUD_PROJECT!;
const LOCATION = process.env.CLOUD_TASKS_LOCATION || 'asia-northeast1';
const QUEUE_NAME = process.env.CLOUD_TASKS_QUEUE || 'async-jobs';
const WORKER_URL = process.env.WORKER_URL!;
const queuePath = client.queuePath(PROJECT_ID, LOCATION, QUEUE_NAME);
export const createTask = async (
payload: TaskPayload,
options?: {
scheduleTime?: Date;
dispatchDeadlineSeconds?: number;
}
): Promise<string> => {
const task: any = {
httpRequest: {
httpMethod: 'POST',
url: WORKER_URL,
headers: {
'Content-Type': 'application/json',
},
body: Buffer.from(JSON.stringify(payload)).toString('base64'),
// OIDCトークンでCloud Run認証
oidcToken: {
serviceAccountEmail: `${PROJECT_ID}@appspot.gserviceaccount.com`,
},
},
};
// 実行予定時刻の設定
if (options?.scheduleTime) {
task.scheduleTime = {
seconds: Math.floor(options.scheduleTime.getTime() / 1000),
};
}
// タイムアウト設定(デフォルト10分)
if (options?.dispatchDeadlineSeconds) {
task.dispatchDeadline = {
seconds: options.dispatchDeadlineSeconds,
};
}
const [response] = await client.createTask({
parent: queuePath,
task,
});
return response.name!;
};Producer API
// src/api/enqueue.ts
import { Request, Response } from 'express';
import { randomUUID } from 'crypto';
import { createJob } from '../lib/firestore';
import { createTask } from '../lib/cloudTasks';
import { EnqueueRequest } from '../types/job';
export const enqueueHandler = async (
req: Request,
res: Response
): Promise<void> => {
try {
const body: EnqueueRequest = req.body;
if (!body.type || !body.payload) {
res.status(400).json({ error: 'type and payload are required' });
return;
}
const jobId = randomUUID();
// Firestoreにジョブを作成
await createJob({
jobId,
type: body.type,
status: 'pending',
payload: body.payload,
progress: 0,
});
// Cloud Tasksにタスクを投入
const taskName = await createTask(
{
jobId,
type: body.type,
payload: body.payload,
},
{
scheduleTime: body.scheduleTime
? new Date(body.scheduleTime)
: undefined,
}
);
res.status(202).json({
jobId,
status: 'pending',
taskName,
});
} catch (error) {
console.error('Error enqueueing job:', error);
res.status(500).json({ error: 'Internal server error' });
}
};Worker API
// src/api/worker.ts
import { Request, Response } from 'express';
import { updateJobStatus } from '../lib/firestore';
import { TaskPayload } from '../types/job';
// ジョブタイプごとの処理関数
const jobProcessors: Record<
string,
(
jobId: string,
payload: Record<string, unknown>,
updateProgress: (progress: number) => Promise<void>
) => Promise<Record<string, unknown>>
> = {
'send-email': async (jobId, payload, updateProgress) => {
const { to, subject, body } = payload as {
to: string;
subject: string;
body: string;
};
await updateProgress(10);
// メール送信ロジック...
await updateProgress(100);
return { messageId: `msg-${Date.now()}`, sentAt: new Date().toISOString() };
},
'process-image': async (jobId, payload, updateProgress) => {
const { imageUrl } = payload as { imageUrl: string };
await updateProgress(10);
// 画像処理ロジック...
await updateProgress(50);
// さらに処理...
await updateProgress(100);
return { processedUrl: `${imageUrl}-processed` };
},
};
export const workerHandler = async (
req: Request,
res: Response
): Promise<void> => {
// Cloud Tasksからのリクエスト検証
const taskName = req.headers['x-cloudtasks-taskname'];
if (!taskName) {
res.status(403).json({ error: 'Forbidden' });
return;
}
const taskPayload: TaskPayload = req.body;
const { jobId, type, payload } = taskPayload;
console.log(`Processing job ${jobId} of type ${type}`);
try {
// 状態を processing に更新
await updateJobStatus(jobId, 'processing', { progress: 0 });
const processor = jobProcessors[type];
if (!processor) {
throw new Error(`Unknown job type: ${type}`);
}
// 進捗更新用のヘルパー
const updateProgress = async (progress: number) => {
await updateJobStatus(jobId, 'processing', { progress });
};
const result = await processor(jobId, payload, updateProgress);
// 状態を completed に更新
await updateJobStatus(jobId, 'completed', { result, progress: 100 });
console.log(`Job ${jobId} completed successfully`);
res.status(200).json({ success: true });
} catch (error) {
const errorMessage =
error instanceof Error ? error.message : 'Unknown error';
console.error(`Job ${jobId} failed:`, errorMessage);
// 状態を failed に更新
await updateJobStatus(jobId, 'failed', { error: errorMessage });
// 5xxを返すとCloud Tasksがリトライする
// リトライ不要な場合は2xxを返す
res.status(500).json({ error: errorMessage });
}
};Status API
// src/api/status.ts
import { Request, Response } from 'express';
import { getJob } from '../lib/firestore';
export const statusHandler = async (
req: Request,
res: Response
): Promise<void> => {
const jobId = req.params.jobId;
if (!jobId) {
res.status(400).json({ error: 'jobId is required' });
return;
}
const job = await getJob(jobId);
if (!job) {
res.status(404).json({ error: 'Job not found' });
return;
}
res.status(200).json(job);
};Express サーバー
// src/server.ts
import express from 'express';
import { enqueueHandler } from './api/enqueue';
import { workerHandler } from './api/worker';
import { statusHandler } from './api/status';
const app = express();
app.use(express.json());
// API Routes
app.post('/jobs', enqueueHandler);
app.post('/worker', workerHandler);
app.get('/jobs/:jobId', statusHandler);
// Health check
app.get('/health', (req, res) => {
res.status(200).json({ status: 'healthy' });
});
const PORT = process.env.PORT || 8080;
app.listen(PORT, () => {
console.log(`Server listening on port ${PORT}`);
});Dockerfile
# syntax = docker/dockerfile:1
ARG NODE_VERSION=20
FROM node:${NODE_VERSION}-slim AS base
WORKDIR /app
ENV NODE_ENV="production"
# --- Build stage ---
FROM base AS build
RUN apt-get update -qq && \
apt-get install --no-install-recommends -y build-essential python3
COPY package*.json ./
RUN npm ci --include=dev
COPY . .
RUN npm run build && \
npm prune --production
# --- Final stage ---
FROM base
RUN apt-get update -qq && \
apt-get install --no-install-recommends -y curl && \
rm -rf /var/lib/apt/lists/*
COPY --from=build /app/node_modules /app/node_modules
COPY --from=build /app/dist /app/dist
COPY --from=build /app/package.json /app/package.json
# 非特権ユーザーで実行
RUN useradd -m appuser
USER appuser
EXPOSE 8080
CMD ["node", "dist/server.js"]Terraform
# Variables
variable "project_id" {
description = "GCP Project ID"
}
variable "region" {
default = "asia-northeast1"
}
variable "environment" {
default = "production"
}
# Enable APIs
resource "google_project_service" "apis" {
for_each = toset([
"run.googleapis.com",
"cloudtasks.googleapis.com",
"firestore.googleapis.com",
"artifactregistry.googleapis.com",
])
project = var.project_id
service = each.value
disable_on_destroy = false
}
# Artifact Registry
resource "google_artifact_registry_repository" "app" {
location = var.region
repository_id = "app"
format = "DOCKER"
}
# Cloud Tasks Queue
resource "google_cloud_tasks_queue" "jobs" {
name = "${var.environment}-async-jobs"
location = var.region
rate_limits {
max_dispatches_per_second = 100
max_concurrent_dispatches = 10
}
retry_config {
max_attempts = 5
max_retry_duration = "3600s" # 1時間
min_backoff = "10s"
max_backoff = "300s"
max_doublings = 4
}
stackdriver_logging_config {
sampling_ratio = 1.0
}
}
# Service Account
resource "google_service_account" "app" {
account_id = "${var.environment}-async-app"
display_name = "Async Processing App Service Account"
}
# IAM bindings
resource "google_project_iam_member" "app_firestore" {
project = var.project_id
role = "roles/datastore.user"
member = "serviceAccount:${google_service_account.app.email}"
}
resource "google_project_iam_member" "app_tasks" {
project = var.project_id
role = "roles/cloudtasks.enqueuer"
member = "serviceAccount:${google_service_account.app.email}"
}
resource "google_project_iam_member" "app_run_invoker" {
project = var.project_id
role = "roles/run.invoker"
member = "serviceAccount:${google_service_account.app.email}"
}
# Cloud Run Service (API + Worker)
resource "google_cloud_run_v2_service" "app" {
name = "${var.environment}-async-app"
location = var.region
template {
service_account = google_service_account.app.email
containers {
image = "${var.region}-docker.pkg.dev/${var.project_id}/app/async-app:latest"
env {
name = "GOOGLE_CLOUD_PROJECT"
value = var.project_id
}
env {
name = "CLOUD_TASKS_LOCATION"
value = var.region
}
env {
name = "CLOUD_TASKS_QUEUE"
value = google_cloud_tasks_queue.jobs.name
}
env {
name = "WORKER_URL"
value = "https://${var.environment}-async-app-${random_id.suffix.hex}-an.a.run.app/worker"
}
ports {
container_port = 8080
}
resources {
limits = {
cpu = "1"
memory = "512Mi"
}
}
startup_probe {
http_get {
path = "/health"
port = 8080
}
initial_delay_seconds = 5
period_seconds = 10
failure_threshold = 3
}
liveness_probe {
http_get {
path = "/health"
port = 8080
}
period_seconds = 30
}
}
scaling {
min_instance_count = 1 # ワーカー処理のため常時1台
max_instance_count = 10
}
}
traffic {
percent = 100
type = "TRAFFIC_TARGET_ALLOCATION_TYPE_LATEST"
}
}
resource "random_id" "suffix" {
byte_length = 4
}
# Public access for API endpoints
resource "google_cloud_run_service_iam_member" "public" {
location = google_cloud_run_v2_service.app.location
service = google_cloud_run_v2_service.app.name
role = "roles/run.invoker"
member = "allUsers"
}
# Cloud Tasks can invoke worker endpoint
resource "google_cloud_run_service_iam_member" "tasks_invoker" {
location = google_cloud_run_v2_service.app.location
service = google_cloud_run_v2_service.app.name
role = "roles/run.invoker"
member = "serviceAccount:${google_service_account.app.email}"
}
# Firestore Database (Native mode)
resource "google_firestore_database" "default" {
project = var.project_id
name = "(default)"
location_id = var.region
type = "FIRESTORE_NATIVE"
concurrency_mode = "OPTIMISTIC"
app_engine_integration_mode = "DISABLED"
}
# Firestore Index for job queries
resource "google_firestore_index" "jobs_by_status" {
project = var.project_id
database = google_firestore_database.default.name
collection = "async_jobs"
fields {
field_path = "status"
order = "ASCENDING"
}
fields {
field_path = "createdAt"
order = "DESCENDING"
}
}
# Firestore TTL Policy (requires Firestore Native mode)
resource "google_firestore_field" "jobs_ttl" {
project = var.project_id
database = google_firestore_database.default.name
collection = "async_jobs"
field = "expiresAt"
ttl_config {}
}
# Outputs
output "service_url" {
value = google_cloud_run_v2_service.app.uri
}
output "tasks_queue" {
value = google_cloud_tasks_queue.jobs.name
}フロントエンド連携(Firestore リアルタイムリスナー)
// frontend/src/lib/firebase.ts
import { initializeApp } from 'firebase/app';
import { getFirestore } from 'firebase/firestore';
const firebaseConfig = {
apiKey: process.env.NEXT_PUBLIC_FIREBASE_API_KEY,
authDomain: process.env.NEXT_PUBLIC_FIREBASE_AUTH_DOMAIN,
projectId: process.env.NEXT_PUBLIC_FIREBASE_PROJECT_ID,
};
const app = initializeApp(firebaseConfig);
export const db = getFirestore(app);// frontend/src/hooks/useJobStatus.ts
import { useState, useEffect } from 'react';
import { doc, onSnapshot } from 'firebase/firestore';
import { db } from '../lib/firebase';
interface JobStatus {
status: 'pending' | 'processing' | 'completed' | 'failed';
progress?: number;
result?: Record<string, unknown>;
error?: string;
}
export const useJobStatus = (jobId: string | null) => {
const [status, setStatus] = useState<JobStatus | null>(null);
const [loading, setLoading] = useState(true);
useEffect(() => {
if (!jobId) {
setLoading(false);
return;
}
setLoading(true);
const unsubscribe = onSnapshot(
doc(db, 'async_jobs', jobId),
(snapshot) => {
if (snapshot.exists()) {
const data = snapshot.data();
setStatus({
status: data.status,
progress: data.progress,
result: data.result,
error: data.error,
});
} else {
setStatus(null);
}
setLoading(false);
},
(error) => {
console.error('Error listening to job status:', error);
setLoading(false);
}
);
return () => unsubscribe();
}, [jobId]);
return { status, loading };
};React Component
// frontend/src/components/JobProgress.tsx
import { useJobStatus } from '../hooks/useJobStatus';
interface JobProgressProps {
jobId: string;
onComplete?: (result: Record<string, unknown>) => void;
onError?: (error: string) => void;
}
export const JobProgress: React.FC<JobProgressProps> = ({
jobId,
onComplete,
onError,
}) => {
const { status, loading } = useJobStatus(jobId);
useEffect(() => {
if (status?.status === 'completed' && status.result) {
onComplete?.(status.result);
}
if (status?.status === 'failed' && status.error) {
onError?.(status.error);
}
}, [status, onComplete, onError]);
if (loading) {
return <div>読み込み中...</div>;
}
if (!status) {
return <div>ジョブが見つかりません</div>;
}
return (
<div className="job-progress">
<div className="status">状態: {status.status}</div>
{status.progress !== undefined && (
<div className="progress-bar">
<div
className="progress-fill"
style={{ width: `${status.progress}%` }}
/>
<span>{status.progress}%</span>
</div>
)}
{status.status === 'failed' && (
<div className="error">エラー: {status.error}</div>
)}
{status.status === 'completed' && (
<div className="result">完了しました</div>
)}
</div>
);
};監視とアラート
# Cloud Monitoring Alert Policy
resource "google_monitoring_alert_policy" "tasks_failed" {
display_name = "${var.environment} - Cloud Tasks Failed"
combiner = "OR"
conditions {
display_name = "High task failure rate"
condition_threshold {
filter = "resource.type=\"cloud_tasks_queue\" AND metric.type=\"cloudtasks.googleapis.com/queue/task_attempt_count\" AND metric.label.response_code!=\"200\""
duration = "300s"
comparison = "COMPARISON_GT"
threshold_value = 10
aggregations {
alignment_period = "60s"
per_series_aligner = "ALIGN_RATE"
}
}
}
notification_channels = [google_monitoring_notification_channel.email.name]
}
resource "google_monitoring_alert_policy" "run_errors" {
display_name = "${var.environment} - Cloud Run Errors"
combiner = "OR"
conditions {
display_name = "High error rate"
condition_threshold {
filter = "resource.type=\"cloud_run_revision\" AND metric.type=\"run.googleapis.com/request_count\" AND metric.label.response_code_class=\"5xx\""
duration = "300s"
comparison = "COMPARISON_GT"
threshold_value = 5
aggregations {
alignment_period = "60s"
per_series_aligner = "ALIGN_RATE"
}
}
}
notification_channels = [google_monitoring_notification_channel.email.name]
}
resource "google_monitoring_notification_channel" "email" {
display_name = "Email Notification"
type = "email"
labels = {
email_address = "alerts@example.com"
}
}セキュリティ設定
Firestore Security Rules
// firestore.rules
rules_version = '2';
service cloud.firestore {
match /databases/{database}/documents {
// async_jobs コレクション
match /async_jobs/{jobId} {
// 認証済みユーザーのみ読み取り可能
allow read: if request.auth != null;
// サーバーサイドからのみ書き込み可能(Service Account経由)
allow write: if false;
}
}
}制限事項と注意点
Cloud Tasks制限
- 最大タスクサイズ: 100KB
- 最大スケジュール期間: 30日先まで
- 最大リトライ回数: 設定可能(推奨5回程度)
- タイムアウト: 最大30分(Cloud Run宛ての場合)
Cloud Run制限
- 最大リクエストタイムアウト: 60分
- 最大インスタンス数: デフォルト1000(引き上げ可能)
- コールドスタート: min_instance_count=0だと発生
Firestore制限
- ドキュメントサイズ: 最大1MB
- 書き込みレート: 1ドキュメントあたり1回/秒(推奨)
- 同時リアルタイムリスナー: プロジェクトあたり最大500(課金有効時)
- 注意: 効率的なクエリとリスナーの設計で制限内に収める
コスト最適化
- Cloud Run: min_instance_countを調整
- Cloud Tasks: 不要なタスクを早期キャンセル
- Firestore: TTLで古いドキュメントを自動削除