mybook

イベント駆動パターン — S3・SQS・EventBridge トリガー

「ファイルをアップロードしたら自動的に処理が走る」

ダイチはこの仕組みに興奮していた。Railsでは Active Job + Sidekiq でバックグラウンド処理を実装していたが、Lambdaならインフラ管理なしに同じことができる。

ECサイトの課題は3つあった。

  1. 商品画像のリサイズ: 出品者がアップロードした画像を複数サイズに変換
  2. 注文確認メール: 注文確定時のメール送信をメイン処理から切り離す
  3. 在庫レポート: 毎朝9時に在庫データをCSVで生成してSlackに通知

これらがイベント駆動パターンの典型的なユースケースだ。

S3 トリガー — 画像リサイズ処理

アーキテクチャ

Loading diagram...

SAM テンプレート

Resources:
  OriginBucket:
    Type: AWS::S3::Bucket
    Properties:
      BucketName: !Sub "${AWS::StackName}-origin-images"
 
  ThumbBucket:
    Type: AWS::S3::Bucket
    Properties:
      BucketName: !Sub "${AWS::StackName}-thumb-images"
 
  ResizeFunction:
    Type: AWS::Serverless::Function
    Properties:
      CodeUri: src/image_resize/
      Handler: app.lambda_handler
      Runtime: ruby3.2
      Timeout: 30
      MemorySize: 1024  # 画像処理は多めのメモリを確保
      Policies:
        - S3ReadPolicy:
            BucketName: !Ref OriginBucket
        - S3WritePolicy:
            BucketName: !Ref ThumbBucket
      Events:
        ImageUpload:
          Type: S3
          Properties:
            Bucket: !Ref OriginBucket
            Events: s3:ObjectCreated:*
            Filter:
              S3Key:
                Rules:
                  - Name: suffix
                    Value: '.jpg'
                  - Name: suffix
                    Value: '.png'

画像リサイズの Lambda 関数

# src/image_resize/app.rb
# frozen_string_literal: true
 
require 'json'
require 'aws-sdk-s3'
require 'mini_magick'  # ImageMagick の Ruby バインディング
require 'logger'
 
$logger = Logger.new($stdout)
$s3 = Aws::S3::Client.new
 
SIZES = {
  thumbnail: { width: 150, height: 150 },
  medium: { width: 600, height: 600 },
  large: { width: 1200, height: 1200 }
}.freeze
 
def lambda_handler(event:, context:)
  # S3 イベントには複数のレコードが含まれる場合がある
  event['Records'].each do |record|
    process_image(record)
  end
 
  { statusCode: 200, body: 'OK' }
rescue StandardError => e
  $logger.error("Error: #{e.class} - #{e.message}\n#{e.backtrace.first(5).join("\n")}")
  raise  # Lambda に失敗を通知してリトライさせる
end
 
private
 
def process_image(record)
  bucket = record.dig('s3', 'bucket', 'name')
  key = CGI.unescape(record.dig('s3', 'object', 'key'))
 
  $logger.info("Processing: s3://#{bucket}/#{key}")
 
  # S3 からオリジナル画像をダウンロード
  response = $s3.get_object(bucket:, key:)
  original_data = response.body.read
 
  # /tmp に保存(Lambda の tmp は 10GB まで利用可能)
  original_path = "/tmp/original_#{Time.now.to_i}#{File.extname(key)}"
  File.write(original_path, original_data, mode: 'wb')
 
  # 各サイズにリサイズしてアップロード
  SIZES.each do |size_name, dimensions|
    resize_and_upload(
      original_path:,
      key:,
      size_name:,
      dimensions:
    )
  end
ensure
  File.delete(original_path) if original_path && File.exist?(original_path)
end
 
def resize_and_upload(original_path:, key:, size_name:, dimensions:)
  output_path = "/tmp/#{size_name}_#{Time.now.to_i}#{File.extname(key)}"
 
  image = MiniMagick::Image.open(original_path)
  image.resize "#{dimensions[:width]}x#{dimensions[:height]}>"  # > は縮小のみ
  image.strip  # EXIF データを削除
  image.write(output_path)
 
  # サムネイル用バケットにアップロード
  dest_key = "#{size_name}/#{key}"
  $s3.put_object(
    bucket: ENV['THUMB_BUCKET_NAME'],
    key: dest_key,
    body: File.read(output_path, mode: 'rb'),
    content_type: 'image/jpeg',
    cache_control: 'public, max-age=31536000'
  )
 
  $logger.info("Uploaded: #{size_name} -> s3://#{ENV['THUMB_BUCKET_NAME']}/#{dest_key}")
ensure
  File.delete(output_path) if output_path && File.exist?(output_path)
end

WARNING

Lambda の /tmp は最大10GBだが、ウォームスタート時に前回の実行で残ったファイルが残る場合がある。必ず ensure ブロックで一時ファイルを削除すること。また、S3から直接ストリームでメモリに読み込む場合はメモリサイズに注意。

SQS トリガー — 非同期処理キュー

注文確認メールのような「失敗したらリトライしてほしい」処理にはSQSが最適だ。

アーキテクチャ

Loading diagram...
# template.yaml
Resources:
  OrderNotificationQueue:
    Type: AWS::SQS::Queue
    Properties:
      VisibilityTimeout: 60  # Lambda タイムアウト x 6 が推奨
      RedrivePolicy:
        deadLetterTargetArn: !GetAtt OrderDLQ.Arn
        maxReceiveCount: 3  # 3回失敗でDLQへ
 
  OrderDLQ:
    Type: AWS::SQS::Queue
    Properties:
      MessageRetentionPeriod: 1209600  # 14日間保持
 
  EmailFunction:
    Type: AWS::Serverless::Function
    Properties:
      CodeUri: src/email/
      Handler: app.lambda_handler
      Timeout: 10
      Policies:
        - SQSPollerPolicy:
            QueueName: !GetAtt OrderNotificationQueue.QueueName
        - Statement:
          - Effect: Allow
            Action: ses:SendEmail
            Resource: "*"
      Events:
        SQSEvent:
          Type: SQS
          Properties:
            Queue: !GetAtt OrderNotificationQueue.Arn
            BatchSize: 10              # 一度に最大10件処理
            FunctionResponseTypes:
              - ReportBatchItemFailures  # 失敗したメッセージのみリトライ
# src/email/app.rb
# frozen_string_literal: true
 
require 'json'
require 'aws-sdk-ses'
require 'logger'
 
$logger = Logger.new($stdout)
$ses = Aws::SES::Client.new
 
def lambda_handler(event:, context:)
  failed_message_ids = []
 
  event['Records'].each do |record|
    message_id = record['messageId']
 
    begin
      body = JSON.parse(record['body'])
      send_order_confirmation(body)
      $logger.info("Email sent for order #{body['order_id']}")
    rescue StandardError => e
      $logger.error("Failed to process #{message_id}: #{e.message}")
      failed_message_ids << message_id
    end
  end
 
  # ReportBatchItemFailures を使って失敗したメッセージのみリトライ
  {
    batchItemFailures: failed_message_ids.map { |id| { itemIdentifier: id } }
  }
end
 
private
 
def send_order_confirmation(order_data)
  $ses.send_email(
    source: 'no-reply@ec-site.example.com',
    destination: { to_addresses: [order_data['customer_email']] },
    message: {
      subject: { data: "ご注文ありがとうございます(注文番号: #{order_data['order_id']})" },
      body: {
        html: {
          data: build_html_body(order_data),
          charset: 'UTF-8'
        }
      }
    }
  )
end
 
def build_html_body(order_data)
  <<~HTML
    <h1>ご注文確認</h1>
    <p>#{order_data['customer_name']} 様</p>
    <p>注文番号: <strong>#{order_data['order_id']}</strong></p>
    <table>
      <tr><th>商品名</th><th>数量</th><th>金額</th></tr>
      #{order_data['items'].map { |item|
        "<tr><td>#{item['name']}</td><td>#{item['quantity']}</td><td>¥#{item['price']}</td></tr>"
      }.join}
    </table>
    <p>合計: ¥#{order_data['total']}</p>
  HTML
end

注文処理側(送信者)のコード:

# 注文処理 Lambda(または Rails アプリ)から SQS に投入
require 'aws-sdk-sqs'
 
$sqs = Aws::SQS::Client.new
QUEUE_URL = ENV['ORDER_NOTIFICATION_QUEUE_URL']
 
def enqueue_order_notification(order)
  message = {
    order_id: order.id,
    customer_email: order.user.email,
    customer_name: order.user.name,
    items: order.items.map { |item|
      { name: item.product.name, quantity: item.quantity, price: item.price }
    },
    total: order.total_price
  }
 
  $sqs.send_message(
    queue_url: QUEUE_URL,
    message_body: JSON.generate(message),
    message_attributes: {
      'order_id' => {
        data_type: 'String',
        string_value: order.id.to_s
      }
    }
  )
end

EventBridge — スケジュール実行

毎朝9時の在庫レポート生成はEventBridgeのスケジュールで実現する。

Resources:
  InventoryReportFunction:
    Type: AWS::Serverless::Function
    Properties:
      CodeUri: src/reports/
      Handler: inventory.lambda_handler
      Timeout: 300  # 5分
      MemorySize: 512
      Environment:
        Variables:
          SLACK_WEBHOOK_URL: !Sub '{{resolve:ssm:/myapp/slack/webhook_url}}'
      Policies:
        - DynamoDBReadPolicy:
            TableName: !Ref ProductsTable
      Events:
        DailyInventoryReport:
          Type: Schedule
          Properties:
            Schedule: cron(0 0 * * ? *)  # UTC 0:00 = JST 9:00
            Name: daily-inventory-report
            Enabled: true
# src/reports/inventory.rb
# frozen_string_literal: true
 
require 'json'
require 'net/http'
require 'uri'
require 'aws-sdk-dynamodb'
require 'csv'
 
$dynamodb = Aws::DynamoDB::Resource.new
$table = $dynamodb.table(ENV['PRODUCTS_TABLE_NAME'])
 
def lambda_handler(event:, context:)
  puts "Generating inventory report at #{Time.now.iso8601}"
 
  low_stock_products = scan_low_stock_products
  csv_data = generate_csv(low_stock_products)
 
  # S3 にCSVを保存
  save_to_s3(csv_data)
 
  # Slack に通知
  notify_slack(low_stock_products.size, csv_data)
 
  puts "Report generated: #{low_stock_products.size} low-stock products"
  { statusCode: 200, body: 'Report generated' }
end
 
private
 
def scan_low_stock_products
  products = []
 
  $table.scan(
    filter_expression: 'stock < :threshold',
    expression_attribute_values: { ':threshold' => 10 }
  ).each_page do |page|
    products.concat(page.items)
  end
 
  products.sort_by { |p| p['stock'].to_i }
end
 
def generate_csv(products)
  CSV.generate(headers: true, encoding: 'UTF-8') do |csv|
    csv << %w[商品ID 商品名 カテゴリ 在庫数]
    products.each do |product|
      csv << [product['id'], product['name'], product['category'], product['stock']]
    end
  end
end
 
def save_to_s3(csv_data)
  s3 = Aws::S3::Client.new
  date = Time.now.strftime('%Y%m%d')
 
  s3.put_object(
    bucket: ENV['REPORTS_BUCKET'],
    key: "inventory/#{date}_inventory_report.csv",
    body: csv_data,
    content_type: 'text/csv; charset=utf-8'
  )
end
 
def notify_slack(count, csv_preview)
  webhook_url = ENV['SLACK_WEBHOOK_URL']
  return unless webhook_url
 
  message = {
    text: "在庫レポート (#{Time.now.strftime('%Y/%m/%d')})",
    attachments: [{
      color: count > 0 ? 'warning' : 'good',
      fields: [{
        title: '在庫不足商品数',
        value: "#{count} 件",
        short: true
      }],
      footer: "S3に詳細CSVを保存しました"
    }]
  }
 
  uri = URI.parse(webhook_url)
  http = Net::HTTP.new(uri.host, uri.port)
  http.use_ssl = true
  request = Net::HTTP::Post.new(uri.path, 'Content-Type' => 'application/json')
  request.body = JSON.generate(message)
  http.request(request)
end

EventBridge ルール — イベントバスによる疎結合化

Lambda を直接呼ぶのではなく、EventBridge のイベントバスを介して疎結合にする設計も強力だ。

Loading diagram...
# 注文処理 Lambda: EventBridge にイベントを発行
def publish_order_completed(order)
  events_client = Aws::EventBridge::Client.new
 
  events_client.put_events(
    entries: [{
      source: 'ec-site.orders',
      detail_type: 'OrderCompleted',
      detail: JSON.generate({
        order_id: order.id,
        customer_id: order.user_id,
        total: order.total_price,
        items: order.items.map { |i| { product_id: i.product_id, quantity: i.quantity } }
      }),
      event_bus_name: ENV['EVENT_BUS_NAME']
    }]
  )
end
# メール送信Lambdaのトリガー設定
EmailFunction:
  Type: AWS::Serverless::Function
  Properties:
    Events:
      OrderCompleted:
        Type: EventBridgeRule
        Properties:
          EventBusName: !Ref AppEventBus
          Pattern:
            source:
              - ec-site.orders
            detail-type:
              - OrderCompleted

INFO

EventBridgeを使うと、注文処理Lambdaはメール送信・ポイント付与・ログ記録などの処理を知らなくてよい。新しい処理を追加するときも注文処理Lambdaを修正せず、EventBridgeのルールを追加するだけでよい。これがイベント駆動アーキテクチャの真価だ。

エラーハンドリングとリトライ

Loading diagram...

各トリガーのリトライ動作:

トリガーリトライ回数失敗時の扱い
S32回設定した DLQ に送信
SQSキュー設定によるDead Letter Queue
EventBridge24時間以内に最大185回DLQ または破棄
非同期呼び出し2回DLQ に送信

ダイチはこれらのパターンを使って、ECサイトのバックグラウンド処理の多くをLambdaに移行した。Sidekiqのワーカーサーバー(EC2 t3.small × 2台、月約¥5,000)を廃止でき、さらにSQSのリトライ機能がSidekiqより確実であることに気づいた。

「Sidekiqのメモリリークで週1回再起動してたのがなくなった」とダイチはチームの日報に書いた。