こんにちは。Tomoyuki(@tomoyuki65)です。
これまでは同期的な処理を中心に学んできましたが、実務においては非同期的な処理が必要になる場面も少なくありません。
例えばユーザーがボタンをクリックした際に、大量データを集計したり、外部APIを呼び出したりする場合など、完了までに時間のかかる処理を実行するようなケースです。
こうした処理を同期的に実行すると、レスポンスを返すまでに30秒以上かかってしまうことがあり、ユーザーを長時間待たせることになるため、ユーザビリティの低下やタイムアウトにもつながりかねません。
そのため、このような時間のかかる処理を実行するためにも、非同期処理についてもしっかり学んでおくことが大切です。
そこで私も非同期処理について調べることにしましたが、日本の実務でよく使われるクラウドサービス「AWS」において非同期処理を実装するなら、Amazon SQSというサービスを利用することになります。
また、今回も以前のAWS Lambdaの記事と同様に、Floci(フローシー)を利用して試していきます。
ということで、この記事では、Go言語(Golang)とFloci(フローシー)を使い、AWSのAmazon SQSによる非同期的な処理を実装する方法について解説します。
Go言語(Golang)とFloci(フローシー)を使ったAWSのAmazon SQSによる非同期処理入門
AWSのAmazon SQS(Simple Queue Service)とは?
Amazon SQS(Simple Queue Service) は、AWSが提供するフルマネージド型(サーバー管理が不要)のメッセージキューサービスです。
システム間でメッセージを一時的に保管・受け渡しすることで、アプリケーション同士を疎結合にし、非同期処理や負荷分散を実現できます。
Floci(フローシー)や事前準備について
今回はFloci(フローシー)を利用して試しますが、そんなFlociや事前準備の部分については、以前のAWS Lambdaの記事と同様なので、以下の記事を参考にして下さい。
関連記事

ローカル開発環境の構築
次にローカル開発環境を構築するため、以下のコマンドを実行して各種ファイルを作成します。
$ mkdir go-sqs && cd go-sqs
$ mkdir -p deploy/docker/local/api/go && touch deploy/docker/local/api/go/Dockerfile
$ mkdir -p src/api && touch src/api/main.go
$ touch .env compose.yaml .gitignore
次に作成したファイルをそれぞれ以下のように記述します。
・「deploy/docker/local/api/go/Dockerfile」
FROM golang:1.26.5-alpine3.23
WORKDIR /go/src/api
COPY ./src/api .
# go.modがあれば依存関係をインストール
RUN if [ -f ./go.mod ]; then \
go install; \
fi
# 開発用のライブラリをインストール
RUN go install github.com/air-verse/air@v1.65.3
EXPOSE 8080
・「src/api/main.go」
package main
import (
"context"
"fmt"
"net/http"
"os"
"os/signal"
"syscall"
"time"
"github.com/labstack/echo/v5"
)
func main() {
// 全体の共通コンテキスト
ctx := context.Background()
// echoのルーター設定
e := echo.New()
e.GET("/", func(c *echo.Context) error {
res := map[string]string{
"message": "Hello World !!",
}
return c.JSON(http.StatusOK, res)
})
// ポート番号の取得
port := os.Getenv("GO_API_PORT")
if port == "" {
port = "8080"
}
// Echoサーバーの起動設定
sc := echo.StartConfig{
Address: fmt.Sprintf(":%s", port),
GracefulTimeout: 5 * time.Second, // 終了シグナル受信後、リクエスト完了を待つ最大時間
}
// OSの終了シグナル(SIGINT/SIGTERM)を受け取ったらサーバーを正常終了するためのコンテキストを生成
signalCtx, stopSignal := signal.NotifyContext(ctx, os.Interrupt, syscall.SIGTERM)
defer stopSignal()
// サーバー起動
e.Logger.Info("start go-api")
if err := sc.Start(signalCtx, e); err != nil {
e.Logger.Error("failed to start go-api server", "error", err)
}
}
・「.env」
ENV=local
GO_API_PORT=8080
GO_LAMBDA_PORT=8081
GO_SQS_WORKER_PORT=8082
TZ=Asia/Tokyo
AWS_ACCESS_KEY_ID=dummy
AWS_SECRET_ACCESS_KEY=dummy
AWS_REGION=ap-northeast-1
AWS_BASE_ENDPOINT=http://floci:4566
LAMBDA_QUEUE_URL=http://floci:4566/000000000000/lambda-queue
WORKER_QUEUE_URL=http://floci:4566/000000000000/worker-queue
SQS_MAX_NUMBER_OF_MESSAGES=1
SQS_MESSAGE_TIMEOUT_SECONDS=3
※後述で利用する環境変数の値も入れています。
・「compose.yaml」
services:
go-api:
container_name: go-api
build:
context: .
dockerfile: ./deploy/docker/local/api/go/Dockerfile
command: air -c .air.toml
volumes:
- ./src/api:/go/src/api
ports:
- "${GO_API_PORT}:${GO_API_PORT}"
env_file:
- .env
tty: true
stdin_open: true
floci:
container_name: floci
image: floci/floci:latest
ports:
- "4566:4566"
volumes:
- go_sqs_floci_data:/app/data
- /var/run/docker.sock:/var/run/docker.sock
environment:
TZ: Asia/Tokyo
FLOCI_STORAGE_MODE: hybrid
volumes:
go_sqs_floci_data:
・「.gitignore」
.DS_Store
.env
/src/*/tmp
次に以下のコマンドを実行し、Dockerコンテナのビルドを行います。
$ docker compose build --no-cache
次に以下のコマンドを実行し、Goの初期化をします。
$ docker compose run --rm go-api go mod init go-api
$ docker compose run --rm go-api go mod tidy
$ docker compose run --rm go-api air init
次に以下のコマンドを実行し、Dockerコンテナを起動します。
$ docker compose up -d
次に以下のコマンドを実行し、ステータスを確認します。
$ docker compose ps
実行後、以下のように各Dockerコンテナが起動していればOKです。

次にcurlコマンドを使って以下のコマンドを実行し、ルートのレスポンス結果を確認します。
$ curl -v http://localhost:8080
実行後、以下のように想定通りのレスポンス結果が返ってこればOKです。

次に以下のコマンドを実行し、一度Dockerコンテナを停止しておきます。
$ docker compose down -v
※ここではオプション「-v」でボリュームも削除しておきます。
AWSのAmazon SQSとLambdaを使った非同期処理の開発方法
次にSQSとLambdaを使った非同期処理を作ってみるため、以下のコマンドを実行して各種ファイルを作成します。
$ mkdir -p deploy/docker/local/lambda/go && touch deploy/docker/local/lambda/go/Dockerfile
$ mkdir -p src/lambda/event-lambda && touch src/lambda/event-lambda/main.go
$ touch Makefile
次に作成したファイルをそれぞれ以下のように記述します。
・「deploy/docker/local/lambda/go/Dockerfile」
FROM golang:1.26.5-alpine3.23
# インストール可能なパッケージ一覧の更新
RUN apk update && \
apk upgrade && \
# パッケージのインストール(--no-cacheでキャッシュ削除)
apk add --no-cache \
zip
WORKDIR /go/src/lambda
COPY ./src/lambda .
# go.modがあれば依存関係をインストール
RUN for dir in ./*; do \
if [ -d "$dir" ] && [ -f "$dir/go.mod" ]; then \
echo "Downloading dependencies in $dir"; \
(cd "$dir" && go mod download); \
fi; \
done
EXPOSE 8081
・「src/lambda/event-lambda/main.go」
package main
import (
"context"
"encoding/json"
"fmt"
"time"
"github.com/aws/aws-lambda-go/events"
"github.com/aws/aws-lambda-go/lambda"
)
// イベントの構造体を定義
type Event struct {
EventType string `json:"eventType"`
Message string `json:"message"`
}
func handler(ctx context.Context, event events.SQSEvent) error {
// イベントのレコードをループ処理
for _, record := range event.Records {
// レコードのBodyをEvent構造体にデコード
var e Event
if err := json.Unmarshal([]byte(record.Body), &e); err != nil {
return err
}
// イベント情報をログ出力
fmt.Printf(
"messageId=%s receiveCount=%s eventType=%s message=%s\n",
record.MessageId,
record.Attributes["ApproximateReceiveCount"],
e.EventType,
e.Message,
)
// DLQ動作確認用に意図的に失敗させる条件設定
switch e.EventType {
case "TEST_ERROR":
return fmt.Errorf(
"intentional test error: messageId=%s receiveCount=%s",
record.MessageId,
record.Attributes["ApproximateReceiveCount"],
)
case "TEST_TIMEOUT_ERROR":
// 5秒かかる処理をシミュレーション
select {
case <-time.After(5 * time.Second):
// 5秒経過(今回はここまで来ない)
return nil
case <-ctx.Done():
// 3秒でこちらに来る
return fmt.Errorf(
"processing timeout error: messageId=%s receiveCount=%s",
record.MessageId,
record.Attributes["ApproximateReceiveCount"],
)
}
}
}
return nil
}
func main() {
lambda.Start(handler)
}
※今回はサンプルとして、e.EventType が特定のタイプの場合にエラーとなるような条件を設定しています。
・「Makefile」
.PHONY: \
create-sqs-lambda redrive-lambda-dlq \
go-mod-tidy-event-lambda \
build-event-lambda deploy-local-event-lambda update-local-event-lambda \
create-esm-event-lambda
# ローカル環境用のawsコマンドの変数定義
AWS_LOCAL := aws --profile local --endpoint-url=http://localhost:4566
# =====================
# Lambda用のSQS定義
# =====================
# DLQと通常Queueの作成
create-sqs-lambda:
@set -e; \
DLQ_URL=$$($(AWS_LOCAL) sqs create-queue \
--queue-name lambda-queue-dlq \
--query QueueUrl \
--output text); \
DLQ_ARN=$$($(AWS_LOCAL) sqs get-queue-attributes \
--queue-url $$DLQ_URL \
--attribute-names QueueArn \
--query Attributes.QueueArn \
--output text); \
QUEUE_URL=$$($(AWS_LOCAL) sqs create-queue \
--queue-name lambda-queue \
--query QueueUrl \
--output text); \
$(AWS_LOCAL) sqs set-queue-attributes \
--queue-url $$QUEUE_URL \
--attributes '{"RedrivePolicy":"{\"deadLetterTargetArn\":\"'"$$DLQ_ARN"'\",\"maxReceiveCount\":\"3\"}","VisibilityTimeout":"18","MessageRetentionPeriod":"1209600","ReceiveMessageWaitTimeSeconds":"20","SqsManagedSseEnabled":"true"}'
# DLQから通常Queueへのリドライブ(DLQの全てを再実行)
redrive-lambda-dlq:
@set -e; \
DLQ_URL=$$($(AWS_LOCAL) sqs get-queue-url \
--queue-name lambda-queue-dlq \
--query QueueUrl \
--output text); \
DLQ_ARN=$$($(AWS_LOCAL) sqs get-queue-attributes \
--queue-url $$DLQ_URL \
--attribute-names QueueArn \
--query Attributes.QueueArn \
--output text); \
QUEUE_URL=$$($(AWS_LOCAL) sqs get-queue-url \
--queue-name lambda-queue \
--query QueueUrl \
--output text); \
QUEUE_ARN=$$($(AWS_LOCAL) sqs get-queue-attributes \
--queue-url $$QUEUE_URL \
--attribute-names QueueArn \
--query Attributes.QueueArn \
--output text); \
echo "Starting redrive from lambda-queue-dlq to lambda-queue..."; \
$(AWS_LOCAL) sqs start-message-move-task \
--source-arn $$DLQ_ARN \
--destination-arn $$QUEUE_ARN
# =====================
# event-lambda用の定義
# =====================
# go.mod更新
go-mod-tidy-event-lambda:
docker compose run --rm go-lambda go -C ./event-lambda mod tidy
# ビルド
build-event-lambda:
docker compose run --rm \
-e GOOS=linux \
-e GOARCH=amd64 \
-e CGO_ENABLED=0 \
go-lambda \
go -C ./event-lambda build -tags lambda.norpc -o bootstrap main.go
docker compose run --rm go-lambda \
zip -j event-lambda/function.zip event-lambda/bootstrap
# ローカル環境にデプロイ
deploy-local-event-lambda:
@set -e; \
QUEUE_URL=$$($(AWS_LOCAL) sqs get-queue-url \
--queue-name lambda-queue \
--query QueueUrl \
--output text); \
QUEUE_ARN=$$($(AWS_LOCAL) sqs get-queue-attributes \
--queue-url $$QUEUE_URL \
--attribute-names QueueArn \
--query "Attributes.QueueArn" \
--output text); \
ROLE_ARN=$$($(AWS_LOCAL) iam create-role \
--role-name event-lambda-role \
--assume-role-policy-document \
"{\"Version\":\"2012-10-17\",\"Statement\":[{\"Effect\":\"Allow\",\"Principal\":{\"Service\":\"lambda.amazonaws.com\"},\"Action\":\"sts:AssumeRole\"}]}" \
--query "Role.Arn" \
--output text); \
POLICY_ARN=$$($(AWS_LOCAL) iam create-policy \
--policy-name event-lambda-policy \
--policy-document \
'{\"Version\":\"2012-10-17\",\"Statement\":[{\"Effect\":\"Allow\",\"Action\":[\"sqs:ReceiveMessage\",\"sqs:DeleteMessage\",\"sqs:GetQueueAttributes\"],\"Resource\":\"'"$$QUEUE_ARN"'\"},{\"Effect\":\"Allow\",\"Action\":[\"logs:CreateLogGroup\",\"logs:CreateLogStream\",\"logs:PutLogEvents\"],\"Resource\":\"*\"}]}' \
--query "Policy.Arn" \
--output text); \
$(AWS_LOCAL) iam attach-role-policy \
--role-name event-lambda-role \
--policy-arn $$POLICY_ARN; \
sleep 1; \
$(AWS_LOCAL) lambda create-function \
--function-name event \
--runtime provided.al2023 \
--handler bootstrap \
--role $$ROLE_ARN \
--timeout 3 \
--environment 'Variables={TZ=Asia/Tokyo}' \
--zip-file fileb://src/lambda/event-lambda/function.zip;
# ローカル環境のlambdaを更新
update-local-event-lambda:
$(AWS_LOCAL) lambda update-function-code \
--function-name event \
--zip-file fileb://src/lambda/event-lambda/function.zip
# LambdaにSQSを紐付け
create-esm-event-lambda:
@set -e; \
QUEUE_URL=$$($(AWS_LOCAL) sqs get-queue-url \
--queue-name lambda-queue \
--query 'QueueUrl' \
--output text); \
QUEUE_ARN=$$($(AWS_LOCAL) sqs get-queue-attributes \
--queue-url "$$QUEUE_URL" \
--attribute-names QueueArn \
--query 'Attributes.QueueArn' \
--output text); \
$(AWS_LOCAL) lambda create-event-source-mapping \
--function-name event \
--event-source-arn "$$QUEUE_ARN" \
--batch-size 1 \
--enabled
※DLQと通常Queueの作成の「--attributes」、ローカル環境にデプロイの「--assume-role-policy-document」、「--policy-document」の設定については、本番環境を想定して設定していますが、実務では要件に合わせて適切な設定をして下さい。
・通常Queueのattributes設定一覧
| 属性名 | 設定値 | 説明 |
|---|---|---|
RedrivePolicy.deadLetterTargetArn |
$$DLQ_ARN |
処理に失敗したメッセージの移動先となるDLQのARNです。 |
RedrivePolicy.maxReceiveCount |
3 |
メッセージの最大受信回数です。3回受信されても正常に処理されなかった場合、DLQへ移動します。 |
VisibilityTimeout |
18 |
メッセージを受信した後、ほかのコンシューマーから見えなくなる時間です。単位は秒で、この設定では18秒です。 |
MessageRetentionPeriod |
1209600 |
メッセージをSQSに保持する期間です。単位は秒で、1209600秒は14日です。 |
ReceiveMessageWaitTimeSeconds |
20 |
メッセージ受信時のロングポーリング待機時間です。メッセージがない場合、最大20秒待機します。 |
SqsManagedSseEnabled |
true |
SQSが管理するサーバーサイド暗号化を有効にします。保存されるメッセージが自動的に暗号化されます。 |
※VisibilityTimeoutの設定秒数については、AWSの推奨式として「VisibilityTimeout ≥ Lambda Timeout × 6 + MaximumBatchingWindowInSeconds」というのがあります。今回の例ではLambdaのタイムアウト設定が3秒、MaximumBatchingWindowInSecondsはデフォルト設定(0秒)のため、「3 × 6 + 0 = 18」を設定しています。尚、MaximumBatchingWindowInSecondsは、「Lambdaのイベントソースマッピングで、メッセージをバッチとしてまとめるために待機する最大時間(秒)」を指します。
次にファイル「src/api/main.go」、「compose.yaml」をそれぞれ以下のように修正します。
・「src/api/main.go」
package main
import (
"context"
"encoding/json"
"fmt"
"net/http"
"os"
"os/signal"
"syscall"
"time"
"github.com/aws/aws-sdk-go-v2/aws"
awsconfig "github.com/aws/aws-sdk-go-v2/config"
"github.com/aws/aws-sdk-go-v2/service/sqs"
"github.com/labstack/echo/v5"
)
// SQSとLambda用のリクエストボディ定義
type SqsLambdaRequestBody struct {
EventType string `json:"eventType" validate:"required"`
Message string `json:"message" validate:"required"`
}
// SQSクライアント生成
func newSqsClient(ctx context.Context) (*sqs.Client, error) {
// AWSリージョンの取得
awsRegion := os.Getenv("AWS_REGION")
if awsRegion == "" {
awsRegion = "ap-northeast-1"
}
// SQSクライアント作成
cfg, err := awsconfig.LoadDefaultConfig(
ctx,
awsconfig.WithRegion(awsRegion),
)
if err != nil {
return nil, err
}
opts := []func(*sqs.Options){}
if awsBaseEndpoint := os.Getenv("AWS_BASE_ENDPOINT"); awsBaseEndpoint != "" {
opts = append(opts, func(o *sqs.Options) {
o.BaseEndpoint = aws.String(awsBaseEndpoint)
})
}
client := sqs.NewFromConfig(cfg, opts...)
return client, nil
}
func main() {
// 全体の共通コンテキスト
ctx := context.Background()
// SQSクライアント取得
sqsClient, err := newSqsClient(ctx)
if err != nil {
panic(fmt.Errorf("failed to create SQS client: %w", err))
}
// echoのルーター設定
e := echo.New()
e.GET("/", func(c *echo.Context) error {
res := map[string]string{
"message": "Hello World !!",
}
return c.JSON(http.StatusOK, res)
})
// グループ定義「/api/v1」
apiV1 := e.Group("/api/v1")
// POST /api/v1/sqs-lambda
apiV1.POST("/sqs-lambda", func(c *echo.Context) error {
ctx := c.Request().Context()
// リクエストボディのバリデーションチェック
var body SqsLambdaRequestBody
if err := c.Bind(&body); err != nil {
return echo.NewHTTPError(http.StatusBadRequest, "failed to bind request body")
}
if body.EventType == "" || body.Message == "" {
return echo.NewHTTPError(http.StatusBadRequest, "eventType and message are required")
}
// bodyをJSON形式のバイト列に変換
jsonBody, err := json.Marshal(body)
if err != nil {
return echo.NewHTTPError(http.StatusInternalServerError, "failed to marshal request body")
}
// Lambda用のQueue URLを取得
queueURL := os.Getenv("LAMBDA_QUEUE_URL")
if queueURL == "" {
queueURL = "http://floci:4566/000000000000/lambda-queue"
}
// SQS実行
output, err := sqsClient.SendMessage(ctx, &sqs.SendMessageInput{
QueueUrl: aws.String(queueURL),
MessageBody: aws.String(string(jsonBody)),
})
if err != nil {
return echo.NewHTTPError(http.StatusInternalServerError, "failed to send SQS message")
}
// レスポンス定義
res := map[string]string{
"messageId": aws.ToString(output.MessageId),
}
return c.JSON(http.StatusAccepted, res)
})
// ポート番号の取得
port := os.Getenv("GO_API_PORT")
if port == "" {
port = "8080"
}
// Echoサーバーの起動設定
sc := echo.StartConfig{
Address: fmt.Sprintf(":%s", port),
GracefulTimeout: 5 * time.Second, // 終了シグナル受信後、リクエスト完了を待つ最大時間
}
// OSの終了シグナル(SIGINT/SIGTERM)を受け取ったらサーバーを正常終了するためのコンテキストを生成
signalCtx, stopSignal := signal.NotifyContext(ctx, os.Interrupt, syscall.SIGTERM)
defer stopSignal()
// サーバー起動
e.Logger.Info("start go-api")
if err := sc.Start(signalCtx, e); err != nil {
e.Logger.Error("failed to start go-api server", "error", err)
}
}
・「compose.yaml」
services:
go-api:
container_name: go-api
build:
context: .
dockerfile: ./deploy/docker/local/api/go/Dockerfile
command: air -c .air.toml
volumes:
- ./src/api:/go/src/api
ports:
- "${GO_API_PORT}:${GO_API_PORT}"
env_file:
- .env
tty: true
stdin_open: true
go-lambda:
container_name: go-lambda
build:
context: .
dockerfile: ./deploy/docker/local/lambda/go/Dockerfile
command: >
sh -c '
if [ -n "$$LAMBDA_DIR" ]; then
exec go -C "./$$LAMBDA_DIR" run main.go
else
echo "LAMBDA_DIR is not set."
tail -f /dev/null
fi
'
ports:
- "${GO_LAMBDA_PORT}:${GO_LAMBDA_PORT}"
volumes:
- ./src/lambda:/go/src/lambda
env_file:
- .env
environment:
LAMBDA_DIR: ${LAMBDA_DIR:-}
floci:
container_name: floci
image: floci/floci:latest
ports:
- "4566:4566"
volumes:
- go_sqs_floci_data:/app/data
- /var/run/docker.sock:/var/run/docker.sock
environment:
TZ: Asia/Tokyo
FLOCI_STORAGE_MODE: hybrid
volumes:
go_sqs_floci_data:
次に以下のコマンドを実行し、Dockerコンテナのビルドを行います。
$ docker compose build --no-cache
次に以下のコマンドを実行し、go-apiコンテナの方のgo.modの更新を行います。
$ docker compose run --rm go-api go mod tidy
次に以下のコマンドを実行し、go-lambdaコンテナの方のGoの初期化をします。
$ docker compose run --rm go-lambda go -C ./event-lambda mod init event-lambda
$ make go-mod-tidy-event-lambda
次に以下のコマンドを実行し、Dockerコンテナを起動します。
$ docker compose up -d
次に以下のコマンドを実行し、ステータスを確認します。
$ docker compose ps
実行後、以下のように各Dockerコンテナが起動していればOKです。

次に以下のコマンドを実行し、flociコンテナにSQSを作成します。
$ make create-sqs-lambda
次に以下のコマンドを実行し、Goコードのビルドをします。
$ make build-event-lambda
ビルド実行後、ファイル「src/lambda/event-lambda/bootstrap」、「src/lambda/event-lambda/function.zip」が作成されればOKです。
次に以下のコマンドを実行し、flociコンテナにGoコードのLambdaをデプロイします。
$ make deploy-local-event-lambda
実行後、以下のようにエラーが出ずに完了すればOKです。

※ターミナルが占有されますが、キーボードの「q」を押すと抜けれます。
次に以下のコマンドを実行し、LambdaにSQSを紐付けます。
$ make create-esm-event-lambda
実行後、以下のようにエラーが出ずに完了すればOKです。

SQSを実行して試す
次にSQSを実行して試してみます。curlコマンドを使って以下のコマンドを実行します。
curl -X POST "http://localhost:8080/api/v1/sqs-lambda" \
-H "Content-Type: application/json" \
-d '{"eventType": "LogEvent", "message": "SQSからLambdaの非同期処理を起動!"}'
実行後、以下のように想定通りのレスポンス結果が返ってこればOKです。

次に以下のコマンドを実行し、flociコンテナのログを確認します。
$ docker compose logs floci
実行後、以下のように想定通りのログ出力がされていればOKです。

次に以下のコマンドを実行し、通常Queueのメッセージ件数を確認します。
aws-local sqs get-queue-attributes \
--queue-url http://localhost:4566/000000000000/lambda-queue \
--attribute-names ApproximateNumberOfMessages
実行後、ApproximateNumberOfMessagesが「0」ならOKです。

eventType「TEST_ERROR」のエラーパターンを試す
次にDLQの動作確認をするため、eventType「TEST_ERROR」のエラーパターンを試します。
curlコマンドを使って以下のコマンドを実行します。
curl -X POST "http://localhost:8080/api/v1/sqs-lambda" \
-H "Content-Type: application/json" \
-d '{"eventType": "TEST_ERROR", "message": "DLQの動作確認!"}'
実行後、以下のように想定通りのレスポンス結果が返ってこればOKです。

次に以下のコマンドを実行し、flociコンテナのログを確認します。
$ docker compose logs floci
実行後、以下のように3回エラーが発生し、最後にメッセージがDLQに移動させるメッセージが出力されていればOKです。

※SQS作成時の「--attributes」の設定で、maxReceiveCountを「3」にしているため、3回実行されるまでリトライされている。
次に以下のコマンドを実行し、通常Queueのメッセージ件数を確認します。
aws-local sqs get-queue-attributes \
--queue-url http://localhost:4566/000000000000/lambda-queue \
--attribute-names ApproximateNumberOfMessages
実行後、ApproximateNumberOfMessagesが「0」ならOKです。

次に以下のコマンドを実行し、DLQのメッセージ件数を確認します。
aws-local sqs get-queue-attributes \
--queue-url http://localhost:4566/000000000000/lambda-queue-dlq \
--attribute-names ApproximateNumberOfMessages
実行後、ApproximateNumberOfMessagesが「1」ならOKです。

次に以下のコマンドを実行し、DLQのメッセージ内容を確認します。
aws-local sqs receive-message \
--queue-url http://localhost:4566/000000000000/lambda-queue-dlq \
--max-number-of-messages 10 \
--attribute-names All \
--message-attribute-names All
※尚、特定のメッセージIDを直接指定して取得するようなことはできませんが、オプション「–query ‘Messages[?MessageId==`xxxxxxxx-xxxx-xxxx-xxxx-xxxxxxxxxxxx`]’」を使って絞り込みは可能です。
実行後、以下のようにDLQのメッセージ内容を確認できればOKです。

DLQを通常Queueへ戻して再実行したい場合について
DLQのメッセージは後から調査したり、必要に応じて再実行をすることになります。
上記で作成したMakefileに「redrive-lambda-dlq」を記述している通り、AWS CLIコマンド「aws sqs start-message-move-task」を使うことで、DLQにある全てのメッセージを通常Queueへ戻して再実行することは可能です。
ただし、特定のメッセージだけを再実行するようなコマンドはないため、そういったことをしたい場合は、以下のような流れで作業をして下さい。
2. メッセージの本文(必要に応じて属性も)を取得する。
3. 元のSQSキューに対して、同じ内容のメッセージを再送する。
・AWSコンソールから送信する
・または aws sqs send-message コマンドを実行する
4. 再送が成功し、不要になったDLQのメッセージを削除する(重複実行を避けるため)。
・AWSコンソールから送信する
・または aws sqs delete-message コマンドを実行する
aws-local sqs send-message \
--queue-url http://localhost:4566/000000000000/lambda-queue \
--message-body '{"eventType":"TEST_ERROR_RETRY","message":"DLQの動作確認!"}'
※今回の例ではエラー回避のため、eventTypeは「TEST_ERROR_RETRY」に修正して再実行させます。
実行後、以下のようにメッセージIDが返ってこればOKです。

次に以下のコマンドを実行し、flociコンテナのログを確認します。
$ docker compose logs floci
実行後、以下のように正常終了していることを確認します。

次に以下のコマンドを実行し、DLQのメッセージ内容にある「ReceiptHandle」を確認します。
aws-local sqs receive-message \
--queue-url http://localhost:4566/000000000000/lambda-queue-dlq \
--max-number-of-messages 10 \
--attribute-names All \
--message-attribute-names All \
--query 'Messages[?MessageId==`18db4a20-3ae0-473e-b930-7b3b2d6cfded`]'
実行後、以下のようにReceiptHandleの値を確認できればOKです。

※receive-messageを実行し、そのメッセージが受信されたタイミングで新しい ReceiptHandleが発行される仕組みのため、ReceiptHandleの値は固定値ではないので注意して下さい。
次にReceiptHandleの値を使った以下のコマンドを実行し、DLQの対象メッセージを削除します。
aws-local sqs delete-message \
--queue-url http://localhost:4566/000000000000/lambda-queue-dlq \
--receipt-handle "96e1fc0c-2d2e-45e8-8d8e-4437e4858768"
次に以下のコマンドを実行してDLQのメッセージ件数を確認し、先ほど1件あったものが0件になっていることを確認します。
aws-local sqs get-queue-attributes \
--queue-url http://localhost:4566/000000000000/lambda-queue-dlq \
--attribute-names ApproximateNumberOfMessages
実行後、ApproximateNumberOfMessagesが「0」ならOKです。

eventType「TEST_TIMEOUT_ERROR」のタイムアウトエラーパターンを試す
次にeventType「TEST_TIMEOUT_ERROR」のタイムアウトエラーパターンを試します。
curlコマンドを使って以下のコマンドを実行します。
curl -X POST "http://localhost:8080/api/v1/sqs-lambda" \
-H "Content-Type: application/json" \
-d '{"eventType": "TEST_TIMEOUT_ERROR", "message": "タイムアウトエラーの確認!"}'
実行後、以下のように想定通りのレスポンス結果が返ってこればOKです。

次に以下のコマンドを実行し、flociコンテナのログを確認します。
$ docker compose logs floci
実行後、以下のように3回エラーが発生し、最後にメッセージがDLQに移動させるメッセージが出力されていればOKです。

※Lambda作成時にオプション「--timeout 3」を付けているため、3秒経過でタイムアウトエラーとなっている。(リトライ回数は上記と同様にSQS作成時の設定)
次に以下のコマンドを実行し、DLQのメッセージ内容を確認します。
aws-local sqs receive-message \
--queue-url http://localhost:4566/000000000000/lambda-queue-dlq \
--max-number-of-messages 10 \
--attribute-names All \
--message-attribute-names All
実行後、以下のようにDLQのメッセージ内容を確認できればOKです。

次に以下のコマンドを実行し、一度Dockerコンテナを停止しておきます。
$ docker compose down
AWSのAmazon SQSとWorkerを使った非同期処理の開発方法
上記ではLambdaを利用した構成についてご紹介しましたが、Lambdaには実行時間やリソースに制約があるため、その点には注意が必要です。
もし長時間実行される処理や負荷の高い処理を行いたい場合は、Workerを利用する構成にする必要があります。
では次にSQSとWorkerを使った非同期処理を作ってみるため、以下のコマンドを実行して各種ファイルを作成します。
$ mkdir -p deploy/docker/local/worker/go && touch deploy/docker/local/worker/go/Dockerfile
$ mkdir -p src/worker && touch src/worker/main.go
次に作成したファイルをそれぞれ以下のように記述します。
・「deploy/docker/local/worker/go/Dockerfile」
FROM golang:1.26.5-alpine3.23
WORKDIR /go/src/worker
COPY ./src/worker .
# go.modがあれば依存関係をインストール
RUN if [ -f ./go.mod ]; then \
go install; \
fi
# 開発用のライブラリをインストール
RUN go install github.com/air-verse/air@v1.65.3
EXPOSE 8083
・「src/worker/main.go」
package main
import (
"context"
"encoding/json"
"fmt"
"log/slog"
"net/http"
"os"
"os/signal"
"strconv"
"syscall"
"time"
"github.com/aws/aws-sdk-go-v2/aws"
awsconfig "github.com/aws/aws-sdk-go-v2/config"
"github.com/aws/aws-sdk-go-v2/service/sqs"
"github.com/aws/aws-sdk-go-v2/service/sqs/types"
"github.com/labstack/echo/v5"
)
// イベントの構造体を定義
type Event struct {
EventType string `json:"eventType"`
Message string `json:"message"`
}
// SQSクライアント生成
func newSqsClient(ctx context.Context) (*sqs.Client, error) {
// AWSリージョンの取得
awsRegion := os.Getenv("AWS_REGION")
if awsRegion == "" {
awsRegion = "ap-northeast-1"
}
// SQSクライアント作成
cfg, err := awsconfig.LoadDefaultConfig(
ctx,
awsconfig.WithRegion(awsRegion),
)
if err != nil {
return nil, err
}
opts := []func(*sqs.Options){}
if awsBaseEndpoint := os.Getenv("AWS_BASE_ENDPOINT"); awsBaseEndpoint != "" {
opts = append(opts, func(o *sqs.Options) {
o.BaseEndpoint = aws.String(awsBaseEndpoint)
})
}
client := sqs.NewFromConfig(cfg, opts...)
return client, nil
}
// メッセージ処理
func processMessage(ctx context.Context, receiveCount string, msg types.Message) error {
// メッセージのBodyをEvent構造体にデコード
var e Event
if err := json.Unmarshal([]byte(aws.ToString(msg.Body)), &e); err != nil {
return fmt.Errorf("json unmarshal error: %w", err)
}
// イベント情報をログ出力
fmt.Printf(
"messageId=%s receiveCount=%s eventType=%s message=%s\n",
aws.ToString(msg.MessageId),
receiveCount,
e.EventType,
e.Message,
)
// DLQ動作確認用に意図的に失敗させる条件設定
switch e.EventType {
case "TEST_ERROR":
return fmt.Errorf(
"intentional test error: messageId=%s",
aws.ToString(msg.MessageId),
)
case "TEST_TIMEOUT_ERROR":
// 5秒かかる処理をシミュレーション
select {
case <-time.After(5 * time.Second):
// 5秒経過(今回はここまで来ない)
return nil
case <-ctx.Done():
// 3秒でこちらに来る
return fmt.Errorf(
"processing timeout: messageId=%s: %w",
aws.ToString(msg.MessageId),
ctx.Err(),
)
}
}
return nil
}
// SQSのポーリング
func pollSQS(ctx context.Context, env string, client *sqs.Client) {
// Worker用のQueue URLを取得
queueURL := os.Getenv("WORKER_QUEUE_URL")
if queueURL == "" {
queueURL = "http://floci:4566/000000000000/worker-queue"
}
// SQSに1回で受信する最大メッセージ数の取得
maxNumberOfMessagesStr := os.Getenv("SQS_MAX_NUMBER_OF_MESSAGES")
if maxNumberOfMessagesStr == "" {
maxNumberOfMessagesStr = "1"
}
// 数値変換
n, err := strconv.ParseInt(maxNumberOfMessagesStr, 10, 32)
if err != nil {
panic(fmt.Errorf(
"invalid SQS_MAX_NUMBER_OF_MESSAGES %q: %w",
maxNumberOfMessagesStr,
err,
))
}
maxNumberOfMessages := int32(n)
// 範囲チェック
if maxNumberOfMessages < 1 || maxNumberOfMessages > 10 {
panic("SQS_MAX_NUMBER_OF_MESSAGES must be between 1 and 10")
}
// タイムアウト設定の秒数を取得
timeoutStr := os.Getenv("SQS_MESSAGE_TIMEOUT_SECONDS")
if timeoutStr == "" {
timeoutStr = "3"
}
// 数値変換
timeoutSec, err := strconv.Atoi(timeoutStr)
if err != nil {
panic(fmt.Errorf(
"invalid SQS_MESSAGE_TIMEOUT_SECONDS %q: %w",
timeoutStr,
err,
))
}
// 範囲チェック
if timeoutSec <= 0 {
panic("SQS_MESSAGE_TIMEOUT_SECONDS must be greater than 0")
}
// ポーリング処理
for {
out, err := client.ReceiveMessage(ctx, &sqs.ReceiveMessageInput{
QueueUrl: aws.String(queueURL),
MaxNumberOfMessages: maxNumberOfMessages,
MessageSystemAttributeNames: []types.MessageSystemAttributeName{
types.MessageSystemAttributeNameApproximateReceiveCount,
},
})
if err != nil {
slog.Error("receive message error:", "error", err)
time.Sleep(3 * time.Second)
continue
}
if len(out.Messages) == 0 {
continue
}
for _, msg := range out.Messages {
// メッセージの受信回数を取得
receiveCount := msg.Attributes[string(types.MessageSystemAttributeNameApproximateReceiveCount)]
// メッセージ処理のタイムアウト設定
messageCtx, cancel := context.WithTimeout(ctx, time.Duration(timeoutSec)*time.Second)
if err := processMessage(messageCtx, receiveCount, msg); err != nil {
slog.Error("process message error:", "receiveCount", receiveCount, "error", err)
continue
}
// メッセージ処理のタイムアウト設定を1件処理するごとに解放
cancel()
// 処理成功後にメッセージを削除
_, err := client.DeleteMessage(ctx, &sqs.DeleteMessageInput{
QueueUrl: aws.String(queueURL),
ReceiptHandle: msg.ReceiptHandle,
})
if err != nil {
slog.Error("delete message error:", "receiveCount", receiveCount, "error", err)
}
}
}
}
func main() {
// ENVの取得
env := os.Getenv("ENV")
if env == "" {
env = "local"
}
// コンテキスト設定
ctx := context.Background()
// SQSクライアント取得
sqsClient, err := newSqsClient(ctx)
if err != nil {
panic(fmt.Errorf("failed to create SQS client: %w", err))
}
// echoのルーター設定
e := echo.New()
e.GET("/", func(c *echo.Context) error {
res := map[string]string{
"service": "go-sqs-worker",
"status": "ok",
"version": "1.0.0",
"environment": env,
}
return c.JSON(http.StatusOK, res)
})
// SQSのポーリング実行
go pollSQS(ctx, env, sqsClient)
// ポート番号の取得
port := os.Getenv("GO_SQS_WORKER_PORT")
if port == "" {
port = "8080"
}
// OSの終了シグナル(SIGINT/SIGTERM)を受け取ったらサーバーを正常終了するためのコンテキストを生成
signalCtx, stopSignal := signal.NotifyContext(ctx, os.Interrupt, syscall.SIGTERM)
defer stopSignal()
// Echoサーバーの起動設定
sc := echo.StartConfig{
Address: fmt.Sprintf(":%s", port),
GracefulTimeout: 5 * time.Second, // 終了シグナル受信後、リクエスト完了を待つ最大時間
}
// サーバー起動
e.Logger.Info("start go-sqs-worker")
if err := sc.Start(signalCtx, e); err != nil {
e.Logger.Error("failed to start go-sqs-worker server", "error", err)
}
}
※Workerの場合、対象のSQSをポーリングします。また、SQSに1回で受信する最大メッセージ数は環境変数「SQS_MAX_NUMBER_OF_MESSAGES」、タイムアウト設定の秒数は環境変数「SQS_MESSAGE_TIMEOUT_SECONDS」を利用して設定しています。
次にファイル「src/api/main.go」、「compose.yaml」、「Makefile」をそれぞれ以下のように修正します。
・「src/api/main.go」
package main
import (
"context"
"encoding/json"
"fmt"
"net/http"
"os"
"os/signal"
"syscall"
"time"
"github.com/aws/aws-sdk-go-v2/aws"
awsconfig "github.com/aws/aws-sdk-go-v2/config"
"github.com/aws/aws-sdk-go-v2/service/sqs"
"github.com/labstack/echo/v5"
)
// SQSとLambda用のリクエストボディ定義
type SqsLambdaRequestBody struct {
EventType string `json:"eventType" validate:"required"`
Message string `json:"message" validate:"required"`
}
// SQSとWorker用のリクエストボディ定義
type SqsWorkerRequestBody struct {
EventType string `json:"eventType" validate:"required"`
Message string `json:"message" validate:"required"`
}
// SQSクライアント生成
func newSqsClient(ctx context.Context) (*sqs.Client, error) {
// AWSリージョンの取得
awsRegion := os.Getenv("AWS_REGION")
if awsRegion == "" {
awsRegion = "ap-northeast-1"
}
// SQSクライアント作成
cfg, err := awsconfig.LoadDefaultConfig(
ctx,
awsconfig.WithRegion(awsRegion),
)
if err != nil {
return nil, err
}
opts := []func(*sqs.Options){}
if awsBaseEndpoint := os.Getenv("AWS_BASE_ENDPOINT"); awsBaseEndpoint != "" {
opts = append(opts, func(o *sqs.Options) {
o.BaseEndpoint = aws.String(awsBaseEndpoint)
})
}
client := sqs.NewFromConfig(cfg, opts...)
return client, nil
}
func main() {
// 全体の共通コンテキスト
ctx := context.Background()
// SQSクライアント取得
sqsClient, err := newSqsClient(ctx)
if err != nil {
panic(fmt.Errorf("failed to create SQS client: %w", err))
}
// echoのルーター設定
e := echo.New()
e.GET("/", func(c *echo.Context) error {
res := map[string]string{
"message": "Hello World !!",
}
return c.JSON(http.StatusOK, res)
})
// グループ定義「/api/v1」
apiV1 := e.Group("/api/v1")
// POST /api/v1/sqs-lambda
apiV1.POST("/sqs-lambda", func(c *echo.Context) error {
ctx := c.Request().Context()
// リクエストボディのバリデーションチェック
var body SqsLambdaRequestBody
if err := c.Bind(&body); err != nil {
return echo.NewHTTPError(http.StatusBadRequest, "failed to bind request body")
}
if body.EventType == "" || body.Message == "" {
return echo.NewHTTPError(http.StatusBadRequest, "eventType and message are required")
}
// bodyをJSON形式のバイト列に変換
jsonBody, err := json.Marshal(body)
if err != nil {
return echo.NewHTTPError(http.StatusInternalServerError, "failed to marshal request body")
}
// Lambda用のQueue URLを取得
queueURL := os.Getenv("LAMBDA_QUEUE_URL")
if queueURL == "" {
queueURL = "http://floci:4566/000000000000/lambda-queue"
}
// SQS実行
output, err := sqsClient.SendMessage(ctx, &sqs.SendMessageInput{
QueueUrl: aws.String(queueURL),
MessageBody: aws.String(string(jsonBody)),
})
if err != nil {
return echo.NewHTTPError(http.StatusInternalServerError, "failed to send SQS message")
}
// レスポンス定義
res := map[string]string{
"messageId": aws.ToString(output.MessageId),
}
return c.JSON(http.StatusAccepted, res)
})
// POST /api/v1/sqs-worker
apiV1.POST("/sqs-worker", func(c *echo.Context) error {
ctx := c.Request().Context()
// リクエストボディのバリデーションチェック
var body SqsWorkerRequestBody
if err := c.Bind(&body); err != nil {
return echo.NewHTTPError(http.StatusBadRequest, "failed to bind request body")
}
if body.EventType == "" || body.Message == "" {
return echo.NewHTTPError(http.StatusBadRequest, "eventType and message are required")
}
// bodyをJSON形式のバイト列に変換
jsonBody, err := json.Marshal(body)
if err != nil {
return echo.NewHTTPError(http.StatusInternalServerError, "failed to marshal request body")
}
// Lambda用のQueue URLを取得
queueURL := os.Getenv("WORKER_QUEUE_URL")
if queueURL == "" {
queueURL = "http://floci:4566/000000000000/worker-queue"
}
// SQS実行
output, err := sqsClient.SendMessage(ctx, &sqs.SendMessageInput{
QueueUrl: aws.String(queueURL),
MessageBody: aws.String(string(jsonBody)),
// MessageGroupId: aws.String("default"), // Fifoの場合に必須
})
if err != nil {
return echo.NewHTTPError(http.StatusInternalServerError, "failed to send SQS message")
}
// レスポンス定義
res := map[string]string{
"messageId": aws.ToString(output.MessageId),
}
return c.JSON(http.StatusAccepted, res)
})
// ポート番号の取得
port := os.Getenv("GO_API_PORT")
if port == "" {
port = "8080"
}
// Echoサーバーの起動設定
sc := echo.StartConfig{
Address: fmt.Sprintf(":%s", port),
GracefulTimeout: 5 * time.Second, // 終了シグナル受信後、リクエスト完了を待つ最大時間
}
// OSの終了シグナル(SIGINT/SIGTERM)を受け取ったらサーバーを正常終了するためのコンテキストを生成
signalCtx, stopSignal := signal.NotifyContext(ctx, os.Interrupt, syscall.SIGTERM)
defer stopSignal()
// サーバー起動
e.Logger.Info("start go-api")
if err := sc.Start(signalCtx, e); err != nil {
e.Logger.Error("failed to start go-api server", "error", err)
}
}
・「compose.yaml」
services:
go-api:
container_name: go-api
build:
context: .
dockerfile: ./deploy/docker/local/api/go/Dockerfile
command: air -c .air.toml
volumes:
- ./src/api:/go/src/api
ports:
- "${GO_API_PORT}:${GO_API_PORT}"
env_file:
- .env
tty: true
stdin_open: true
go-lambda:
container_name: go-lambda
build:
context: .
dockerfile: ./deploy/docker/local/lambda/go/Dockerfile
command: >
sh -c '
if [ -n "$$LAMBDA_DIR" ]; then
exec go -C "./$$LAMBDA_DIR" run main.go
else
echo "LAMBDA_DIR is not set."
tail -f /dev/null
fi
'
ports:
- "${GO_LAMBDA_PORT}:${GO_LAMBDA_PORT}"
volumes:
- ./src/lambda:/go/src/lambda
env_file:
- .env
environment:
LAMBDA_DIR: ${LAMBDA_DIR:-}
go-sqs-worker:
container_name: go-sqs-worker
build:
context: .
dockerfile: ./deploy/docker/local/worker/go/Dockerfile
command: air -c .air.toml
volumes:
- ./src/worker:/go/src/worker
ports:
- "${GO_SQS_WORKER_PORT}:${GO_SQS_WORKER_PORT}"
env_file:
- .env
tty: true
stdin_open: true
floci:
container_name: floci
image: floci/floci:latest
ports:
- "4566:4566"
volumes:
- go_sqs_floci_data:/app/data
- /var/run/docker.sock:/var/run/docker.sock
environment:
TZ: Asia/Tokyo
FLOCI_STORAGE_MODE: hybrid
volumes:
go_sqs_floci_data:
・「Makefile」
.PHONY: \
create-sqs-lambda redrive-lambda-dlq \
go-mod-tidy-event-lambda \
build-event-lambda deploy-local-event-lambda update-local-event-lambda \
create-esm-event-lambda \
create-sqs-worker redrive-worker-dlq create-sqs-worker-fifo redrive-worker-dlq-fifo
# ローカル環境用のawsコマンドの変数定義
AWS_LOCAL := aws --profile local --endpoint-url=http://localhost:4566
# =====================
# Lambda用のSQS定義
# =====================
# DLQと通常Queueの作成
create-sqs-lambda:
@set -e; \
DLQ_URL=$$($(AWS_LOCAL) sqs create-queue \
--queue-name lambda-queue-dlq \
--query QueueUrl \
--output text); \
DLQ_ARN=$$($(AWS_LOCAL) sqs get-queue-attributes \
--queue-url $$DLQ_URL \
--attribute-names QueueArn \
--query Attributes.QueueArn \
--output text); \
QUEUE_URL=$$($(AWS_LOCAL) sqs create-queue \
--queue-name lambda-queue \
--query QueueUrl \
--output text); \
$(AWS_LOCAL) sqs set-queue-attributes \
--queue-url $$QUEUE_URL \
--attributes '{"RedrivePolicy":"{\"deadLetterTargetArn\":\"'"$$DLQ_ARN"'\",\"maxReceiveCount\":\"3\"}","VisibilityTimeout":"18","MessageRetentionPeriod":"1209600","ReceiveMessageWaitTimeSeconds":"20","SqsManagedSseEnabled":"true"}'
# DLQから通常Queueへのリドライブ(DLQの全てを再実行)
redrive-lambda-dlq:
@set -e; \
DLQ_URL=$$($(AWS_LOCAL) sqs get-queue-url \
--queue-name lambda-queue-dlq \
--query QueueUrl \
--output text); \
DLQ_ARN=$$($(AWS_LOCAL) sqs get-queue-attributes \
--queue-url $$DLQ_URL \
--attribute-names QueueArn \
--query Attributes.QueueArn \
--output text); \
QUEUE_URL=$$($(AWS_LOCAL) sqs get-queue-url \
--queue-name lambda-queue \
--query QueueUrl \
--output text); \
QUEUE_ARN=$$($(AWS_LOCAL) sqs get-queue-attributes \
--queue-url $$QUEUE_URL \
--attribute-names QueueArn \
--query Attributes.QueueArn \
--output text); \
echo "Starting redrive from lambda-queue-dlq to lambda-queue..."; \
$(AWS_LOCAL) sqs start-message-move-task \
--source-arn $$DLQ_ARN \
--destination-arn $$QUEUE_ARN
# =====================
# event-lambda用の定義
# =====================
# go.mod更新
go-mod-tidy-event-lambda:
docker compose run --rm go-lambda go -C ./event-lambda mod tidy
# ビルド
build-event-lambda:
docker compose run --rm \
-e GOOS=linux \
-e GOARCH=amd64 \
-e CGO_ENABLED=0 \
go-lambda \
go -C ./event-lambda build -tags lambda.norpc -o bootstrap main.go
docker compose run --rm go-lambda \
zip -j event-lambda/function.zip event-lambda/bootstrap
# ローカル環境にデプロイ
deploy-local-event-lambda:
@set -e; \
QUEUE_URL=$$($(AWS_LOCAL) sqs get-queue-url \
--queue-name lambda-queue \
--query QueueUrl \
--output text); \
QUEUE_ARN=$$($(AWS_LOCAL) sqs get-queue-attributes \
--queue-url $$QUEUE_URL \
--attribute-names QueueArn \
--query "Attributes.QueueArn" \
--output text); \
ROLE_ARN=$$($(AWS_LOCAL) iam create-role \
--role-name event-lambda-role \
--assume-role-policy-document \
"{\"Version\":\"2012-10-17\",\"Statement\":[{\"Effect\":\"Allow\",\"Principal\":{\"Service\":\"lambda.amazonaws.com\"},\"Action\":\"sts:AssumeRole\"}]}" \
--query "Role.Arn" \
--output text); \
POLICY_ARN=$$($(AWS_LOCAL) iam create-policy \
--policy-name event-lambda-policy \
--policy-document \
'{\"Version\":\"2012-10-17\",\"Statement\":[{\"Effect\":\"Allow\",\"Action\":[\"sqs:ReceiveMessage\",\"sqs:DeleteMessage\",\"sqs:GetQueueAttributes\"],\"Resource\":\"'"$$QUEUE_ARN"'\"},{\"Effect\":\"Allow\",\"Action\":[\"logs:CreateLogGroup\",\"logs:CreateLogStream\",\"logs:PutLogEvents\"],\"Resource\":\"*\"}]}' \
--query "Policy.Arn" \
--output text); \
$(AWS_LOCAL) iam attach-role-policy \
--role-name event-lambda-role \
--policy-arn $$POLICY_ARN; \
sleep 1; \
$(AWS_LOCAL) lambda create-function \
--function-name event \
--runtime provided.al2023 \
--handler bootstrap \
--role $$ROLE_ARN \
--timeout 3 \
--environment 'Variables={TZ=Asia/Tokyo}' \
--zip-file fileb://src/lambda/event-lambda/function.zip;
# ローカル環境のlambdaを更新
update-local-event-lambda:
$(AWS_LOCAL) lambda update-function-code \
--function-name event \
--zip-file fileb://src/lambda/event-lambda/function.zip
# LambdaにSQSを紐付け
create-esm-event-lambda:
@set -e; \
QUEUE_URL=$$($(AWS_LOCAL) sqs get-queue-url \
--queue-name lambda-queue \
--query 'QueueUrl' \
--output text); \
QUEUE_ARN=$$($(AWS_LOCAL) sqs get-queue-attributes \
--queue-url "$$QUEUE_URL" \
--attribute-names QueueArn \
--query 'Attributes.QueueArn' \
--output text); \
$(AWS_LOCAL) lambda create-event-source-mapping \
--function-name event \
--event-source-arn "$$QUEUE_ARN" \
--batch-size 1 \
--enabled
# =====================
# Worker用のSQS定義
# =====================
# DLQと通常Queueの作成
create-sqs-worker:
@set -e; \
DLQ_URL=$$($(AWS_LOCAL) sqs create-queue \
--queue-name worker-queue-dlq \
--query QueueUrl \
--output text); \
DLQ_ARN=$$($(AWS_LOCAL) sqs get-queue-attributes \
--queue-url $$DLQ_URL \
--attribute-names QueueArn \
--query Attributes.QueueArn \
--output text); \
QUEUE_URL=$$($(AWS_LOCAL) sqs create-queue \
--queue-name worker-queue \
--query QueueUrl \
--output text); \
$(AWS_LOCAL) sqs set-queue-attributes \
--queue-url $$QUEUE_URL \
--attributes '{"RedrivePolicy":"{\"deadLetterTargetArn\":\"'"$$DLQ_ARN"'\",\"maxReceiveCount\":\"3\"}","VisibilityTimeout":"18","MessageRetentionPeriod":"1209600","ReceiveMessageWaitTimeSeconds":"20","SqsManagedSseEnabled":"true"}'
# DLQから通常Queueへのリドライブ(DLQの全てを再実行)
redrive-worker-dlq:
@set -e; \
DLQ_URL=$$($(AWS_LOCAL) sqs get-queue-url \
--queue-name worker-queue-dlq \
--query QueueUrl \
--output text); \
DLQ_ARN=$$($(AWS_LOCAL) sqs get-queue-attributes \
--queue-url $$DLQ_URL \
--attribute-names QueueArn \
--query Attributes.QueueArn \
--output text); \
QUEUE_URL=$$($(AWS_LOCAL) sqs get-queue-url \
--queue-name worker-queue \
--query QueueUrl \
--output text); \
QUEUE_ARN=$$($(AWS_LOCAL) sqs get-queue-attributes \
--queue-url $$QUEUE_URL \
--attribute-names QueueArn \
--query Attributes.QueueArn \
--output text); \
echo "Starting redrive from worker-queue-dlq to worker-queue..."; \
$(AWS_LOCAL) sqs start-message-move-task \
--source-arn $$DLQ_ARN \
--destination-arn $$QUEUE_ARN
# FIFOのDLQと通常Queueの作成
create-sqs-worker-fifo:
@set -e; \
DLQ_URL=$$($(AWS_LOCAL) sqs create-queue \
--queue-name worker-queue-dlq.fifo \
--query QueueUrl \
--output text); \
DLQ_ARN=$$($(AWS_LOCAL) sqs get-queue-attributes \
--queue-url $$DLQ_URL \
--attribute-names QueueArn \
--query Attributes.QueueArn \
--output text); \
QUEUE_URL=$$($(AWS_LOCAL) sqs create-queue \
--queue-name worker-queue.fifo \
--attributes FifoQueue=true,ContentBasedDeduplication=true \
--query "QueueUrl" \
--output text); \
$(AWS_LOCAL) sqs set-queue-attributes \
--queue-url $$QUEUE_URL \
--attributes '{"RedrivePolicy":"{\"deadLetterTargetArn\":\"'"$$DLQ_ARN"'\",\"maxReceiveCount\":\"3\"}","VisibilityTimeout":"18","MessageRetentionPeriod":"1209600","ReceiveMessageWaitTimeSeconds":"20","SqsManagedSseEnabled":"true"}'
# FIFOのDLQから通常Queueへのリドライブ(DLQの全メッセージを再実行)
redrive-worker-dlq-fifo:
@set -e; \
DLQ_URL=$$($(AWS_LOCAL) sqs get-queue-url \
--queue-name worker-queue-dlq.fifo \
--query QueueUrl \
--output text); \
DLQ_ARN=$$($(AWS_LOCAL) sqs get-queue-attributes \
--queue-url $$DLQ_URL \
--attribute-names QueueArn \
--query Attributes.QueueArn \
--output text); \
QUEUE_URL=$$($(AWS_LOCAL) sqs get-queue-url \
--queue-name worker-queue.fifo \
--query "QueueUrl" \
--output text); \
QUEUE_ARN=$$($(AWS_LOCAL) sqs get-queue-attributes \
--queue-url $$QUEUE_URL \
--attribute-names QueueArn \
--query Attributes.QueueArn \
--output text); \
echo "Starting redrive from worker-queue-dlq.fifo to worker-queue.fifo..."; \
$(AWS_LOCAL) sqs start-message-move-task \
--source-arn $$DLQ_ARN \
--destination-arn $$QUEUE_ARN
※DLQと通常Queueの作成の「--attributes」、ローカル環境にデプロイの「--assume-role-policy-document」、「--policy-document」の設定については、本番環境を想定して設定していますが、実務では要件に合わせて適切な設定をして下さい。また、FIFO Queueの場合は、sqs create-queueのオプション「--attributes」で「FifoQueue=true」が必須です。「ContentBasedDeduplication」は任意ですが、これはメッセージ本文をもとに重複メッセージを自動で判定するための設定で、有効にすると「MessageDeduplicationId」を明示的に指定しなくても、SQS がメッセージ本文から重複判定を行います。(重複削除が可能になりますが、要件次第なので注意!)
次に以下のコマンドを実行し、Dockerコンテナのビルドを行います。
$ docker compose build --no-cache
次に以下のコマンドを実行し、Goの初期化をします。
$ docker compose run --rm go-sqs-worker go mod init go-sqs-worker
$ docker compose run --rm go-sqs-worker go mod tidy
$ docker compose run --rm go-sqs-worker air init
次に以下のコマンドを実行し、Dockerコンテナを起動します。
$ docker compose up -d
次に以下のコマンドを実行し、ステータスを確認します。
$ docker compose ps
実行後、以下のように各Dockerコンテナが起動していればOKです。

次にcurlコマンドを使って以下のコマンドを実行し、ルートのレスポンス結果を確認します。
$ curl -v http://localhost:8082
実行後、以下のように想定通りのレスポンス結果が返ってこればOKです。

次に以下のコマンドを実行し、flociコンテナにWorker用のSQSを作成します。
$ make create-sqs-worker
SQSを実行して試す
次にSQSを実行して試してみます。curlコマンドを使って以下のコマンドを実行します。
curl -X POST "http://localhost:8080/api/v1/sqs-worker" \
-H "Content-Type: application/json" \
-d '{"eventType": "LogEvent", "message": "SQSからWorkerの非同期処理を起動!"}'
実行後、以下のように想定通りのレスポンス結果が返ってこればOKです。

次に以下のコマンドを実行し、go-sqs-workerコンテナのログを確認します。
$ docker compose logs go-sqs-worker
実行後、以下のように想定通りのログ出力がされていればOKです。

尚、上記のLambdaの方で試したeventType「TEST_ERROR」のエラーパターンや、eventType「TEST_TIMEOUT_ERROR」のタイムアウトエラーパターンは、同様の結果になるため、ここでは割愛します。
通常Queueのメッセージの処理順序について
上記では通常のQueueを使用して動作確認をしましたが、このタイプのQueueではメッセージの処理順序は保証されません。
そのため、送信した順番どおりにメッセージが処理されるとは限らない点に注意して下さい。
もしメッセージの処理順序を保証する必要がある場合は、「FIFO Queue」を利用して下さい。
FIFO Queueとは?
FIFO(First In, First Out)Queue は、「先に送信したメッセージを先に処理する」という順序を保証するキューです。
通常のQueueとは異なり、メッセージの処理順序を維持できるため、処理の順番が重要なシステムで利用されます。
FIFO Queueを試したい場合について
上記で修正したMakefileにFIFO Queueについても追加してあるため、一部設定を修正すれば試すことができます。
試したい場合は、まず以下のコマンドを実行し、一度Dockerコンテナを停止します。
$ docker compose down
次にファイル「.env」の「WORKER_QUEUE_URL」を以下のように修正します。
・・・
WORKER_QUEUE_URL=http://floci:4566/000000000000/worker-queue.fifo
・・・
次にファイル「src/api/main.go」のWorker用の処理のSQS実行の部分にある、MessageGroupIdの部分のコメントアウトを外して有効化します。
・・・
// POST /api/v1/sqs-worker
apiV1.POST("/sqs-worker", func(c *echo.Context) error {
ctx := c.Request().Context()
// リクエストボディのバリデーションチェック
var body SqsWorkerRequestBody
if err := c.Bind(&body); err != nil {
return echo.NewHTTPError(http.StatusBadRequest, "failed to bind request body")
}
if body.EventType == "" || body.Message == "" {
return echo.NewHTTPError(http.StatusBadRequest, "eventType and message are required")
}
// bodyをJSON形式のバイト列に変換
jsonBody, err := json.Marshal(body)
if err != nil {
return echo.NewHTTPError(http.StatusInternalServerError, "failed to marshal request body")
}
// Lambda用のQueue URLを取得
queueURL := os.Getenv("WORKER_QUEUE_URL")
if queueURL == "" {
queueURL = "http://floci:4566/000000000000/worker-queue"
}
// SQS実行
output, err := sqsClient.SendMessage(ctx, &sqs.SendMessageInput{
QueueUrl: aws.String(queueURL),
MessageBody: aws.String(string(jsonBody)),
MessageGroupId: aws.String("default"), // Fifoの場合に必須
})
if err != nil {
return echo.NewHTTPError(http.StatusInternalServerError, "failed to send SQS message")
}
// レスポンス定義
res := map[string]string{
"messageId": aws.ToString(output.MessageId),
}
return c.JSON(http.StatusAccepted, res)
})
・・・
※メッセージの順序は「MessageGroupId」ごとに保証されます。今回はメッセージグループが一つだけのため、「default」を指定しています。もし複数のメッセージグループを扱う場合は、条件やカテゴリごとに一意の「MessageGroupId」を指定することで、それぞれのグループ内でメッセージの順序が保証されます。
次に以下のコマンドを実行し、Dockerコンテナを再起動します。
$ docker compose up -d
次に以下のコマンドを実行し、flociコンテナにWorker用のFIFO Queueを作成します。
$ make create-sqs-worker-fifo
次にSQSを実行して試してみます。curlコマンドを使って以下のコマンドを実行します。
curl -X POST "http://localhost:8080/api/v1/sqs-worker" \
-H "Content-Type: application/json" \
-d '{"eventType": "LogEvent", "message": "FIFO Queueを試す!"}'
実行後、以下のように想定通りのレスポンス結果が返ってこればOKです。

次に以下のコマンドを実行し、go-sqs-workerコンテナのログを確認します。
$ docker compose logs go-sqs-worker
実行後、以下のように想定通りのログ出力がされていればOKです。

今回はFIFO Queueのメッセージ処理順序の確認は省略しますが、処理順序を確認したい場合は、シェルスクリプトなどでループ処理を行い、複数のメッセージを連続して送信することで動作を確認できます。
AWSのAmazon SNSやStep Functionsとの組み合わせについて
今回はAmazon SQSを単体で利用した非同期処理についてご紹介しましたが、実務では他のAWSサービスと組み合わせて利用するケースもあります。
例えば、同じメッセージを複数のシステムへ配信したい場合は「Amazon SNS(Simple Notification Service)」、複数の処理を決められた順序で実行したり、分岐やリトライを制御したりしたい場合は「Step Functions」と組み合わせることで、より柔軟なシステムを構築できます。
非同期処理の冪等性について
非同期処理を作成する際は、冪等性(同じ処理を何度実行しても、最終的な結果が変わらない性質)を意識して設計することも非常に重要です。
その理由として、非同期処理では何らかのエラーによって処理が途中で停止した場合に、必要に応じて同じ処理をリトライさせることがあり、冪等性が担保されていないと、データの二重登録などの問題が発生する可能性があるためです。
例えば、銀行系の外部APIを非同期処理から呼び出すケースでは、外部APIの応答が遅延したことで非同期処理側がタイムアウトエラーになったとしても、実際には外部API側では処理が正常に完了していることがあります。
このような場合、同じ処理をリトライすると二重に実行される可能性がありますが、外部APIが一意のIDなどによって冪等性を担保していれば、同じIDで何度リクエストしても同じ結果が返されるため、安全にリトライすることが可能です。
一方で、外部API側で冪等性が担保されていない場合は、非同期処理側でステータスや処理済みの状態を管理するなど、アプリケーション側で冪等性を担保する仕組みが必要になります。
そのため、非同期処理を設計する際は、外部APIを含めた処理全体で冪等性を担保できるように設計することが重要です。
最後に
今回はGo言語(Golang)とFloci(フローシー)を使ったAWSのAmazon SQSによる非同期処理について解説しました。
非同期処理は、実務においてメール送信やバッチ処理、そして外部サービスとの連携など、様々な場面で利用される重要な仕組みであるため、学んでおくことは大切です。
本記事では、通常Queue、DLQ、FIFO Queueなど、Amazon SQSの基本的な部分は一通りご紹介したので、これから学ぶ方はぜひ参考にしてみて下さい。


コメント