PR

Go言語(Golang)とFloci(フローシー)を使ったAWSのAmazon SQSによる非同期処理入門

3. 応用

こんにちは。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の記事と同様なので、以下の記事を参考にして下さい。

 

関連記事

Go言語(Golang)とFloci(フローシー)を使ったAWS Lambda開発入門
こんにちは。Tomoyuki(@tomoyuki65)です。イベント駆動の非同期処理や軽量なAPIを作りたい場合、AWS Lambdaを使うケースが多かったりするため、Go言語(Golang)で 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」の設定については、本番環境を想定して設定していますが、実務では要件に合わせて適切な設定をして下さい。

 

DLQ(Dead Letter Queue)とは、処理に何度も失敗したメッセージを一時的に退避させるためのキューです。通常の処理は止めず、エラーになったメッセージだけを分離して保存し、後から原因調査や再処理ができるようにします。

 

・通常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へ戻して再実行することは可能です。

ただし、特定のメッセージだけを再実行するようなコマンドはないため、そういったことをしたい場合は、以下のような流れで作業をして下さい。

1. DLQ内の対象メッセージを確認する。
2. メッセージの本文(必要に応じて属性も)を取得する。
3. 元のSQSキューに対して、同じ内容のメッセージを再送する。
 ・AWSコンソールから送信する
 ・または aws sqs send-message コマンドを実行する
4. 再送が成功し、不要になったDLQのメッセージを削除する(重複実行を避けるため)。
 ・AWSコンソールから送信する
 ・または aws sqs delete-message コマンドを実行する

 

例えば今回の例で試す場合、以下のようなコマンドを実行し、SQSを再実行します。
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の基本的な部分は一通りご紹介したので、これから学ぶ方はぜひ参考にしてみて下さい。

 

この記事を書いた人
Tomoyuki

SE→ブロガーを経て、現在はSoftware Engineer(Web/Gopher)をしています!

Tomoyukiをフォローする
3. 応用
スポンサーリンク
Tomoyukiをフォローする

コメント

タイトルとURLをコピーしました