前回は AppSync Events で試すイベント駆動構成 (その1) で AppSync Events を使用したイベント駆動的なWebアプリケーションの構成例などを整理していました。
今回は実際に構成した内容をコードベースで紹介できたらと思っています。

| コンポーネント | 種別 | 役割 |
|---|---|---|
| CloudFront + S3 | frontend | Vite + React SPA のホスティング |
| AppSync Events | frontend | WebSocket Subscribe エンドポイント (tickets/orders) |
| API Gateway + Lambda | backend | POST /v1/tickets/orders を受け付け SQS に Enqueue |
| SQS + DLQ | backend | メッセージバッファ ※3 回失敗で DLQ へ移動 |
| Event Lambda | backend | SQS トリガー、DynamoDB 書き込み + AppSync Publish |
| AppSync Events | backend | HTTP Publish エンドポイント ※API キー認証 |
| DynamoDB | backend | イベント永続化 ※PK: event_id (UUID v7) |
Lambda Web Adapter を使って通常の Go HTTP サーバーとして実装しています。
今回インフラはCDKで定義しているのですが、ARM64 用の Adapter Layer を付与するだけで、http.ListenAndServe(":8080", ...) がそのまま Lambda 上で動作するようになります。
ローカル開発時には同じコードが docker compose up で起動するため、開発と本番の差が少ないのが利点です。
func main() {
ctx := context.Background()
cfg, err := config.Load()
if err != nil { log.Fatalf("failed to load config: %v", err) }
p, err := producer.NewProducer(ctx, cfg.App.Env, cfg.Producer.QueueURL)
if err != nil { log.Fatalf("failed to create producer: %v", err) }
mux := http.NewServeMux()
mux.Handle("POST /v1/tickets/orders", api.NewHandler(p))
// Lambda Web Adapter が :8080 を受け取って Lambda ランタイムに橋渡しする
http.ListenAndServe(":"+cfg.App.Port, corsMiddleware(mux))
}
リクエスト受付・バリデーション・SQS エンキューの辺りは特別なことはしていません。
非同期的な構成としているため、クライアントからのリクエストの検証ができ次第、イベント処理用の Lambda を発火させるため SQS にキューを渡し、APIサーバ自体はリクエストを受け付けた後は即座に 202 を返却して処理を終了します。
type postTicketOrderRequest struct {
EventID string `json:"event_id"`
EventName string `json:"event_name"`
SeatType string `json:"seat_type"`
Quantity int `json:"quantity"`
Amount int `json:"amount"`
}
type sqsMessage struct {
Payload map[string]any `json:"payload"`
EventType string `json:"event_type"`
}
func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
var req postTicketOrderRequest
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
writeJSONError(w, "invalid JSON body", http.StatusBadRequest); return
}
if req.EventID == "" || req.EventName == "" || req.SeatType == "" {
writeJSONError(w, "event_id, event_name and seat_type are required", http.StatusBadRequest); return
}
if req.Quantity <= 0 || req.Amount <= 0 {
writeJSONError(w, "quantity and amount must be greater than 0", http.StatusBadRequest); return
}
msg := sqsMessage{
EventType: "created",
Payload: map[string]any{
"event_id": req.EventID, "event_name": req.EventName,
"seat_type": req.SeatType, "quantity": req.Quantity, "amount": req.Amount,
},
}
body, _ := json.Marshal(msg)
if err := h.producer.Send(r.Context(), string(body)); err != nil {
writeJSONError(w, "failed to send message", http.StatusInternalServerError); return
}
w.WriteHeader(http.StatusAccepted) // 202: 受け付け。処理は非同期で行う。
}
SQSイベントハンドラーとして動作する Lambda になるため、キューを受け取った後に以下の流れで処理します。
func (h *Handler) Handle(ctx context.Context, sqsEvent events.SQSEvent) (events.SQSEventResponse, error) {
log.Printf("event handler invoked: records=%d", len(sqsEvent.Records))
var resp events.SQSEventResponse
for _, record := range sqsEvent.Records {
if err := h.processRecord(ctx, record); err != nil {
log.Printf("failure: message_id=%s err=%v", record.MessageId, err)
resp.BatchItemFailures = append(resp.BatchItemFailures, events.SQSBatchItemFailure{
ItemIdentifier: record.MessageId,
})
}
}
return resp, nil
}
func (h *Handler) processRecord(ctx context.Context, record events.SQSMessage) error {
var msg eventMessage
if err := json.Unmarshal([]byte(record.Body), &msg); err != nil {
return err
}
// 処理遅延シミュレート
time.Sleep(3 * time.Second)
// DynamoDB に書き込む
if err := h.store.PutEvent(ctx, msg.EventType, msg.Payload); err != nil {
return err
}
// AppSync Events に Publish する
orderID, _ := uuid.NewV7()
publishPayload := map[string]any{
"order_id": "ord-" + orderID.String(),
"status": "confirmed",
}
return h.notifier.PublishEvent(ctx, msg.EventType, publishPayload)
}
AppSync Events への Publish に関しては、以下のようにしています。
type eventsRequest struct {
Channel string `json:"channel"`
Events []string `json:"events"` // 各イベントは JSON 文字列として格納する
}
func (n *notifier) PublishEvent(ctx context.Context, eventType string, payload map[string]any) error {
edJSON, _ := json.Marshal(eventData{EventType: eventType, Payload: payload})
body, _ := json.Marshal(eventsRequest{
Channel: n.channel, // "tickets/orders"
Events: []string{string(edJSON)}, // JSON を文字列として渡す (二重シリアライズ)
})
req, _ := http.NewRequestWithContext(ctx, http.MethodPost, n.endpoint+"/event", bytes.NewReader(body))
req.Header.Set("Content-Type", "application/json")
req.Header.Set("x-api-key", n.apiKey) // API キー認証 (本番構成では IAM SigV4 に変更する)
resp, err := n.client.Do(req)
if err != nil { return err }
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
return fmt.Errorf("appsync: unexpected status %d", resp.StatusCode)
}
return nil
}
POST リクエストで Publish する形になりますが、以下2点がポイントとなります。
認証方法はいくつか用意されているようなのですが、今回は動作検証が目的のため API Key を使用しています。
その場合の注意点として、デフォルトでは7日で失効する仕様のため本番運用時には注意が必要です。
From the creation time, the time after which the API key expires. The date is represented as seconds since the epoch, rounded down to the nearest hour. The default value for this parameter is 7 days from creation time. For more information, see ApiKey.
上記は GraphQL API 側のドキュメントで Events API 専用ページは現時点では存在しないようなのですが、未指定で作成すると7日後で失効になっていることが確認できました。※おそらく認証基盤は AppSync 共通なのかと
なお有効期限は最大で365日に延長可能なようです。こちらは Events API のドキュメントがありました。
API_KEY authorization - Configuring authorization and authentication to secure Event APIs
API keys are configurable for up to 365 days, and you can extend an existing expiration date for up to another 365 days from that day.
有効期限については、CI/CDなどの定期更新フローを組むことも可能ですが、IAM SigV4 による認証が本番運用時には適していそうに見えます。※別記事でこの辺りは検証できたらと思っています
AppSync Events API の Publish リクエストでは、events 配列の各要素を JSON オブジェクトではなく JSON 文字列 として渡す必要があります。
Publishing events via HTTP - AWS AppSync Events API
Each specified event in your publish request must be a stringified valid JSON value.
以下は AppSync が要求しているリクエスト形式:
{
"method": "POST",
"headers": {
"content-type": "application/json",
"x-api-key": "da2-your-api-key"
},
"body": {
"channel": "default/channel",
"events": ["{\"event_1\":\"data_1\"}", "{\"event_2\":\"data_2\"}"]
}
}
よって、直感的には events に JSON オブジェクトをそのまま入れたくなりますが、その形式だと受け付けられません。
// NG: オブジェクトをそのまま入れる
Events: []any{
map[string]any{"event_type": "created", "order_id": "ord-xxx"},
}
// → {"channel":"tickets/orders","events":[{"event_type":"created","order_id":"ord-xxx"}]}
// AppSync はこの形式を拒否する
正しくは、一度 json.Marshal でバイト列に変換し、それを string にキャストしてから配列に入れます。
// OK: JSON を文字列化してから入れる (二重シリアライズ)
edJSON, _ := json.Marshal(map[string]any{
"event_type": "created",
"order_id": "ord-xxx",
})
// edJSON = []byte(`{"event_type":"created","order_id":"ord-xxx"}`)
body, _ := json.Marshal(eventsRequest{
Channel: "tickets/orders",
Events: []string{string(edJSON)}, // string にキャストして渡す
})
// → {"channel":"tickets/orders","events":["{\"event_type\":\"created\",\"order_id\":\"ord-xxx\"}"]}
// events[0] がオブジェクト型ではなく、ダブルクォートで囲まれた文字列型になっている
string(edJSON) 自体はバイト列を文字列に変換しているだけですが、外側の json.Marshal がその文字列をクォートとエスケープ付きの JSON 文字列として出力するため、結果的に二重シリアライズになります。
公式ドキュメントを読まないと気づけないポイントでした。
イベント通知は Amplify v6 で Subscribe するようにしています。
Connect to AWS AppSync Events - Amplify Docs
import { Amplify } from "aws-amplify"
Amplify.configure({
API: {
Events: {
endpoint: import.meta.env.VITE_APPSYNC_HTTP_URL, // HTTP Publish URL
region: import.meta.env.VITE_AWS_REGION,
defaultAuthMode: "apiKey",
apiKey: import.meta.env.VITE_APPSYNC_API_KEY,
},
},
})
endpoint に HTTP Publish URL を渡すと Amplify が内部で WebSocket URL に変換して接続してくれていそうです。※WebSocket用のURL(wss://)を環境変数などで持つ必要がない
認証方式については、今回は動作検証が目的のため apiKey を設定しています。※本番相当の環境で動かす場合は Cognito User Pool + Lambda Authorizer の組み合わせの方が良いと思います
ほか、細かい内容ですが Subscribe カスタムフックについては以下のようにしています。
events は aws-amplify から import したものであり、戻り値は useEffect などで呼び出されることを想定しクリーンアップ関数のような形にしています。
function subscribeChannel(addEvent, setError): () => void {
let channel: EventsChannel | null = null
events
.connect("tickets/orders") // チャンネルへ接続
.then((ch) => {
channel = ch
ch.subscribe({
next: (data: ReceivedEvent) => {
// 受信データを EventItem に変換してストアに追加
addEvent(buildEventItem(data.event.event_type, data.event.payload))
},
error: (err) => {
setError("Subscription 接続に失敗しました")
},
})
})
return () => channel?.close() // アンマウント時に接続を切断
}
useEffect(() => {
return subscribeChannel(addEvent, setError) // アンマウント時に channel.close() が呼ばれる
}, [])
状態管理については以下のように定義しています。
イベントを受信した際には新着順で画面上に表示したいため、[item, ...state.events] で先頭追加しつつ、ページ離脱などを考慮して clearEvents によるリセット処理を設定しています。
export const useEventFeedStore = create<EventFeedState>((set) => ({
events: [],
addEvent: (item) => set((state) => ({ events: [item, ...state.events] })),
clearEvents: () => set({ events: [] }),
}))
今回は AppSync Events を使ったイベント駆動の通知基盤の実装をコードベースで紹介しました。
次回は CDK のインフラ定義と、実際に画面上で動かしてみた後にどのような挙動になるのかを見てみたいと思います。
AppSync Events を使用したイベント駆動アーキテクチャの実装をCDKで定義し、WebSocket通信によるリアルタイムイベント配信の構成と動作確認を紹介。
AppSync Events を使用したイベント駆動型Webアプリケーション構築について、GraphQL Subscriptionとの違い、AppSync Eventsの概要、HTTP PublishエンドポイントとWebSocket Subscribeエンドポイントの役割を解説した記事。
AWS Lambdaのエラー通知をSlackに送信する際の設計ポイントを解説。ログ永続化、Subscription Filter上限対策、コスト最適化、通知遅延を考慮し、CloudWatch Logs→Lambda→Slack通知とFirehose→S3保存の構成を採用した実装例を紹介。
AWS AIPの知見を活かし、論文PDFを構造化JSONに変換するパイプラインを構築。S3、Lambda、Step Functions、Textract、Bedrockを組み合わせ、2つの抽出経路で日本語・英語PDFを処理し、メタデータ付きの構造化データをRAGのソースとして活用する仕組みを実装。
論文PDFを構造化JSONに変換するパイプラインの経路A(Textract+Bedrock)について、各ステップの設計理由と実装を詳細に解説。出力上限問題への対処として、モデルに本文を書かせず位置情報のみ返させることで、生成時間を半減させた改善事例を紹介。