サービス間通信 — 同期と非同期
「検索サービスを分離したはいいけど、商品情報をどうやって渡すんですか?」
ケンジの質問はシンプルだが、核心を突いていた。サービス間通信の設計は、マイクロサービスアーキテクチャの中で最も重要な設計決定の一つだ。
同期通信 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
endgRPC: 高性能な同期通信
サービス間の内部通信で、パフォーマンスが重要な場合は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
endREST vs gRPC の使い分け
| 観点 | REST | gRPC |
|---|---|---|
| 可読性 | 高(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
endSNS: ファンアウトパターン
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
endAWS 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ゲートウェイパターンを学ぶ。