mybook

サービス間通信 — 同期と非同期

「検索サービスを分離したはいいけど、商品情報をどうやって渡すんですか?」

ケンジの質問はシンプルだが、核心を突いていた。サービス間通信の設計は、マイクロサービスアーキテクチャの中で最も重要な設計決定の一つだ。


同期通信 vs 非同期通信

通信パターンには大きく2種類ある。

Loading diagram...

いつ同期、いつ非同期?

同期通信が適切なケース:
  ✓ レスポンスをすぐに返す必要がある(ユーザーが待っている)
  ✓ エラー時に即座にユーザーにフィードバックしたい
  例: 商品詳細の表示、ユーザー認証、在庫確認

非同期通信が適切なケース:
  ✓ 処理に時間がかかる(メール送信、推薦計算)
  ✓ 処理の失敗がユーザー体験に直結しない
  ✓ 複数サービスへのファンアウトが必要
  例: 注文確認メール送信、在庫ログ記録、分析イベント

REST HTTP: シンプルな同期通信

最もよく使われる同期通信方式。

# app/services/product_service_client.rb
class ProductServiceClient
  BASE_URL = ENV.fetch('PRODUCT_SERVICE_URL', 'http://product-service')
 
  def self.find(product_id)
    new.find(product_id)
  end
 
  def find(product_id)
    response = connection.get("/api/v1/products/#{product_id}")
    handle_response(response)
  end
 
  def find_batch(product_ids)
    response = connection.get('/api/v1/products', ids: product_ids.join(','))
    handle_response(response)
  end
 
  private
 
  def connection
    @connection ||= Faraday.new(BASE_URL) do |f|
      f.request :retry, max: 3, interval: 0.5, backoff_factor: 2
      f.response :json
      f.options.timeout = 3        # 読み取りタイムアウト
      f.options.open_timeout = 1   # 接続タイムアウト
      f.adapter Faraday.default_adapter
    end
  end
 
  def handle_response(response)
    case response.status
    when 200
      response.body
    when 404
      nil
    when 503
      raise ServiceUnavailableError, 'ProductService is unavailable'
    else
      raise ServiceError, "ProductService returned #{response.status}"
    end
  end
end

サーキットブレーカーパターン

依存サービスの障害時に、自サービスを守る重要なパターン。

# Gemfile
gem 'stoplight'
 
# config/initializers/circuit_breakers.rb
Stoplight::Light.default_data_store = Stoplight::DataStore::Redis.new(
  Redis.new(url: ENV['REDIS_URL'])
)
 
# app/services/resilient_product_client.rb
class ResilientProductClient
  def find(product_id)
    Stoplight("product-service-find") {
      ProductServiceClient.new.find(product_id)
    }
    .with_fallback { |error|
      Rails.logger.warn("ProductService circuit open: #{error&.message}")
      # フォールバック: キャッシュから返す
      Rails.cache.read("product:#{product_id}")
    }
    .with_threshold(3)    # 3回失敗でオープン
    .with_cool_off_time(30)  # 30秒後に半オープン
    .run
  end
end

gRPC: 高性能な同期通信

サービス間の内部通信で、パフォーマンスが重要な場合はgRPCが有効。

# proto/product_service.proto
syntax = "proto3";
 
package shopnova.product.v1;
 
service ProductService {
  rpc GetProduct(GetProductRequest) returns (Product);
  rpc GetProducts(GetProductsRequest) returns (GetProductsResponse);
}
 
message GetProductRequest {
  string product_id = 1;
}
 
message Product {
  string id = 1;
  string name = 2;
  int64 price_cents = 3;
  int32 stock_quantity = 4;
}
 
message GetProductsRequest {
  repeated string product_ids = 1;
}
 
message GetProductsResponse {
  repeated Product products = 1;
}
# app/services/grpc_product_client.rb
class GrpcProductClient
  def initialize
    @stub = Shopnova::Product::V1::ProductService::Stub.new(
      ENV.fetch('PRODUCT_SERVICE_GRPC_ADDR', 'product-service:50051'),
      :this_channel_is_insecure,
      timeout: 2
    )
  end
 
  def find(product_id)
    request = Shopnova::Product::V1::GetProductRequest.new(product_id: product_id)
    @stub.get_product(request)
  rescue GRPC::Unavailable => e
    raise ServiceUnavailableError, e.message
  end
end

REST vs gRPC の使い分け

観点RESTgRPC
可読性高(JSON)低(バイナリ)
パフォーマンス普通高速(HTTP/2、バイナリ)
型安全性高(Protobuf)
ブラウザ対応直接可要ゲートウェイ
採用場面外部API、BFF内部サービス間高頻度通信

SQS: 非同期メッセージキュー

AWS SQS を使ったサービス間の非同期通信。

# Gemfile
gem 'aws-sdk-sqs'
gem 'shoryuken'  # SQS用のWorkerライブラリ
 
# app/services/event_publisher.rb
class EventPublisher
  SQS_CLIENT = Aws::SQS::Client.new(region: 'ap-northeast-1')
 
  QUEUE_URLS = {
    order_events: ENV.fetch('ORDER_EVENTS_QUEUE_URL'),
    inventory_updates: ENV.fetch('INVENTORY_UPDATES_QUEUE_URL')
  }.freeze
 
  def self.publish(queue:, event_type:, payload:)
    message = {
      event_type: event_type,
      event_id: SecureRandom.uuid,
      timestamp: Time.current.iso8601,
      payload: payload
    }
 
    SQS_CLIENT.send_message(
      queue_url: QUEUE_URLS.fetch(queue),
      message_body: message.to_json,
      message_attributes: {
        'event_type' => {
          string_value: event_type.to_s,
          data_type: 'String'
        }
      }
    )
  end
end
 
# 注文サービスでの使用
class OrdersController < ApplicationController
  def create
    order = Order.create!(order_params)
 
    # 非同期でInventoryServiceに在庫減算を依頼
    EventPublisher.publish(
      queue: :order_events,
      event_type: 'order_created',
      payload: {
        order_id: order.id,
        items: order.items.map { |i|
          { product_id: i.product_id, quantity: i.quantity }
        }
      }
    )
 
    render json: order, status: :created
  end
end
# app/workers/inventory_update_worker.rb(在庫サービス側)
class InventoryUpdateWorker
  include Shoryuken::Worker
 
  shoryuken_options queue: ENV.fetch('ORDER_EVENTS_QUEUE_URL'),
                    auto_delete: false  # 処理完了後に手動削除
 
  def perform(sqs_msg, body)
    event = JSON.parse(body)
 
    case event['event_type']
    when 'order_created'
      handle_order_created(event['payload'])
    end
 
    # 成功したらメッセージを削除
    sqs_msg.delete
  rescue StandardError => e
    Rails.logger.error("InventoryUpdateWorker failed: #{e.message}")
    # 削除しない → SQSが自動的に再試行
    raise
  end
 
  private
 
  def handle_order_created(payload)
    payload['items'].each do |item|
      StockItem.find_by!(product_id: item['product_id'])
               .decrement!(:quantity, item['quantity'])
    end
  end
end

SNS: ファンアウトパターン

1つのイベントを複数サービスに届けたいとき、SNS → SQS のファンアウトが有効。

Loading diagram...
# app/services/sns_publisher.rb
class SnsPublisher
  SNS_CLIENT = Aws::SNS::Client.new(region: 'ap-northeast-1')
 
  def self.publish(topic_arn:, event_type:, payload:)
    message = {
      event_type: event_type,
      event_id: SecureRandom.uuid,
      timestamp: Time.current.iso8601,
      payload: payload
    }.to_json
 
    SNS_CLIENT.publish(
      topic_arn: topic_arn,
      message: message,
      message_attributes: {
        'event_type' => {
          data_type: 'String',
          string_value: event_type
        }
      }
    )
  end
end
 
# 注文確定後の通知
class Order < ApplicationRecord
  after_commit :publish_order_event, on: :create
 
  private
 
  def publish_order_event
    SnsPublisher.publish(
      topic_arn: ENV.fetch('ORDER_EVENTS_TOPIC_ARN'),
      event_type: 'order.created',
      payload: {
        order_id: id,
        user_id: user_id,
        total_amount: total_amount,
        items: order_items.map(&:to_event_payload)
      }
    )
  end
end

AWS CDK でインフラを定義

# cloudformation/sqs-sns.yml(簡略化)
Resources:
  OrderEventsTopic:
    Type: AWS::SNS::Topic
    Properties:
      TopicName: shopnova-order-events
 
  InventoryQueue:
    Type: AWS::SQS::Queue
    Properties:
      QueueName: shopnova-inventory-updates
      VisibilityTimeout: 60
      RedrivePolicy:
        deadLetterTargetArn: !GetAtt InventoryDeadLetterQueue.Arn
        maxReceiveCount: 3
 
  InventoryDeadLetterQueue:
    Type: AWS::SQS::Queue
    Properties:
      QueueName: shopnova-inventory-updates-dlq
      MessageRetentionPeriod: 1209600  # 14日間保持
 
  InventorySubscription:
    Type: AWS::SNS::Subscription
    Properties:
      TopicArn: !Ref OrderEventsTopic
      Protocol: sqs
      Endpoint: !GetAtt InventoryQueue.Arn
      FilterPolicy:
        event_type:
          - order.created
          - order.cancelled

冪等性の確保

非同期通信では、同じメッセージが複数回届くことがある(at-least-once delivery)。

# app/workers/idempotent_inventory_worker.rb
class IdempotentInventoryWorker
  include Shoryuken::Worker
 
  shoryuken_options queue: ENV.fetch('INVENTORY_QUEUE_URL'), auto_delete: false
 
  def perform(sqs_msg, body)
    event = JSON.parse(body)
    event_id = event['event_id']
 
    # 冪等性キー: 同じevent_idは1回だけ処理
    return if already_processed?(event_id)
 
    ActiveRecord::Base.transaction do
      process_event(event)
      mark_as_processed(event_id)
    end
 
    sqs_msg.delete
  end
 
  private
 
  def already_processed?(event_id)
    ProcessedEvent.exists?(event_id: event_id)
  end
 
  def mark_as_processed(event_id)
    ProcessedEvent.create!(
      event_id: event_id,
      processed_at: Time.current
    )
  end
end
# db/migrate/xxxx_create_processed_events.rb
class CreateProcessedEvents < ActiveRecord::Migration[7.2]
  def change
    create_table :processed_events do |t|
      t.string :event_id, null: false
      t.datetime :processed_at, null: false
      t.timestamps
    end
 
    add_index :processed_events, :event_id, unique: true
  end
end

まとめ: 通信パターンの選択基準

「整理すると、こうだね」とミサキはまとめた。

ユーザーが待つ → 同期通信(REST/gRPC)
  └ 外部向け・可読性重視 → REST
  └ 内部・高頻度・型安全 → gRPC

ユーザーが待たない → 非同期通信(SQS/SNS)
  └ 1対1のキュー → SQS
  └ 1対多のブロードキャスト → SNS → SQS
  └ 必ず冪等性を実装する

WARNING

非同期通信を使う場合、必ず冪等性(同じメッセージを複数回処理しても結果が同じ)を実装すること。SQSはat-least-onceデリバリーを保証するため、重複配信は必ず起きる。

次章では、これらのサービスへの「玄関口」となるAPIゲートウェイパターンを学ぶ。