0
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?

AWS lambda関数と SQS SNS との連携

0
Posted at

AWS lambda関数と SQS SNS との連携

lambda関数と SQS SNS との連携手順を解説します。
※ 操作はすべてAWS Python ライブラリ(boto3)を使用して実施します。

image

AWSのベストプラクティスに倣い、SNSイベントを直接Lambdaに送るのではなく、SQSを経由させてLambda関数につなぎます。この構成を採用することで、起動イベントが大量に重複・急増した場合でも、処理の取りこぼし(メッセージの紛失)を防ぎ、システムの堅牢性を高めることができます。

以降の内容は以下の記事を前提とします

ロールへの権限付与

SQS/SNS との連携に必要な権限をロールに付与します。

import boto3
from botocore.exceptions import ClientError

PROFILE_NAME = 'resources_dev_admin'
ROLE_NAME = 'ResourcesEnvMainRole'

POLICY_ARNS = [
    'arn:aws:iam::aws:policy/AmazonSQSFullAccess',
    'arn:aws:iam::aws:policy/AmazonSNSFullAccess'
]

session = boto3.Session(profile_name=PROFILE_NAME)
iam_client = session.client('iam')

print(f"Starting to attach policies to role '{ROLE_NAME}' using profile '{PROFILE_NAME}'...")

for policy_arn in POLICY_ARNS:
    try:
        iam_client.attach_role_policy(
            RoleName=ROLE_NAME,
            PolicyArn=policy_arn
        )
        print(f"Success: Attached {policy_arn.split('/')[-1]}.")
    except ClientError as e:
        print(f"Error: Failed to attach {policy_arn.split('/')[-1]}.")
        print(e.response['Error']['Message'])

print("Process completed.")

SQS を作成

import boto3

PROFILE_NAME = 'resources_dev'
ROLE_NAME = 'ResourcesEnvMainRole'
FUNCTION_NAME = 'func_test_01'

QUEUE_NAME = f"{FUNCTION_NAME}-sqs"

session = boto3.Session(
    profile_name=PROFILE_NAME
)

sqs = session.client('sqs')

create_response = sqs.create_queue(
    QueueName=QUEUE_NAME,
    Attributes={
        'VisibilityTimeout': '1080'
    }
)

url_response = sqs.get_queue_url(QueueName=QUEUE_NAME)
queue_url = url_response['QueueUrl']

attrs_response = sqs.get_queue_attributes(
    QueueUrl=queue_url,
    AttributeNames=['QueueArn']
)
queue_arn = attrs_response['Attributes']['QueueArn']

print(queue_url)
print(queue_arn)

SQSのサブクライブ

import boto3

PROFILE_NAME = 'resources_dev'
ROLE_NAME = 'ResourcesEnvMainRole'
FUNCTION_NAME = 'func_test_01'

QUEUE_NAME = f"{FUNCTION_NAME}-sqs"
BATCH_SIZE = 1

session = boto3.Session(
    profile_name=PROFILE_NAME
)

sqs = session.client('sqs')
url_response = sqs.get_queue_url(QueueName=QUEUE_NAME)
queue_url = url_response['QueueUrl']
attrs_response = sqs.get_queue_attributes(
    QueueUrl=queue_url,
    AttributeNames=['QueueArn']
)
queue_arn = attrs_response['Attributes']['QueueArn']

lambda_client = session.client('lambda')

response = lambda_client.create_event_source_mapping(
    FunctionName=FUNCTION_NAME,
    EventSourceArn=queue_arn,
    BatchSize=BATCH_SIZE
)

print(response['UUID'])

SQS メッセージの送信

import boto3

PROFILE_NAME = 'resources_dev'
ROLE_NAME = 'ResourcesEnvMainRole'
FUNCTION_NAME = 'func_test_01'
MESSAGE_BODY = 'Hello SQS!'

session = boto3.Session(
    profile_name=PROFILE_NAME
)

sqs = session.client('sqs')

queue_name = f"{FUNCTION_NAME}-sqs"
url_response = sqs.get_queue_url(QueueName=queue_name)
queue_url = url_response['QueueUrl']

response = sqs.send_message(
    QueueUrl=queue_url,
    MessageBody=MESSAGE_BODY
)

print(response['MessageId'])

SNSの作成

import boto3

PROFILE_NAME = 'resources_dev'
ROLE_NAME = 'ResourcesEnvMainRole'
FUNCTION_NAME = 'func_test_01'
TOPIC_NAME = f"{FUNCTION_NAME}-sns"

session = boto3.Session(profile_name=PROFILE_NAME)

sns = session.client('sns')

response = sns.create_topic(Name=TOPIC_NAME)

print(response['TopicArn'])

SQS に SNS をサブクライブする権限を付与

import json
import boto3

PROFILE_NAME = 'resources_dev'
ROLE_NAME = 'ResourcesEnvMainRole'
FUNCTION_NAME = 'func_test_01'
QUEUE_NAME = f"{FUNCTION_NAME}-sqs"
TOPIC_NAME = f"{FUNCTION_NAME}-sns"

session = boto3.Session(profile_name=PROFILE_NAME)

aws_region = session.region_name
sts_client = session.client('sts')
aws_account_id = sts_client.get_caller_identity()['Account']

sqs = session.client('sqs')

url_response = sqs.get_queue_url(QueueName=QUEUE_NAME)
queue_url = url_response['QueueUrl']

queue_arn = f"arn:aws:sqs:{aws_region}:{aws_account_id}:{QUEUE_NAME}"
sns_arn = f"arn:aws:sns:{aws_region}:{aws_account_id}:{TOPIC_NAME}"

policy_dict = {
    "Version": "2012-10-17",
    "Statement": [
        {
            "Sid": "Allow-SNS-SendMessage",
            "Effect": "Allow",
            "Principal": "*",
            "Action": "sqs:SendMessage",
            "Resource": queue_arn,
            "Condition": {
                "ArnEquals": {
                    "aws:SourceArn": sns_arn
                }
            }
        }
    ]
}

sqs.set_queue_attributes(
    QueueUrl=queue_url,
    Attributes={
        'Policy': json.dumps(policy_dict)
    }
)

print(f"Policy applied to: {queue_url}")

SQS から SNS のサブクライブ

import boto3

PROFILE_NAME = 'resources_dev'
ROLE_NAME = 'ResourcesEnvMainRole'
FUNCTION_NAME = 'func_test_01'
QUEUE_NAME = f"{FUNCTION_NAME}-sqs"
TOPIC_NAME = f"{FUNCTION_NAME}-sns"

session = boto3.Session(profile_name=PROFILE_NAME)

aws_region = session.region_name
sts_client = session.client('sts')
aws_account_id = sts_client.get_caller_identity()['Account']

sns = session.client('sns')

queue_arn = f"arn:aws:sqs:{aws_region}:{aws_account_id}:{QUEUE_NAME}"
sns_arn = f"arn:aws:sns:{aws_region}:{aws_account_id}:{TOPIC_NAME}"

response = sns.subscribe(
    TopicArn=sns_arn,
    Protocol='sqs',
    Endpoint=queue_arn,
    Attributes={
        'RawMessageDelivery': 'true'
    }
)

print(response['SubscriptionArn'])

SNSの発行

import boto3

PROFILE_NAME = 'resources_dev'
ROLE_NAME = 'ResourcesEnvMainRole'
FUNCTION_NAME = 'func_test_01'
MESSAGE = 'SNS message test!!!'

session = boto3.Session(profile_name=PROFILE_NAME)

aws_region = session.region_name
sts_client = session.client('sts')
aws_account_id = sts_client.get_caller_identity()['Account']

sns = session.client('sns')

topic_name = f"{FUNCTION_NAME}-sns"
sns_arn = f"arn:aws:sns:{aws_region}:{aws_account_id}:{topic_name}"

response = sns.publish(
    TopicArn=sns_arn,
    Message=MESSAGE
)

print(response['MessageId'])

ログの確認

from datetime import datetime, timedelta, timezone
import boto3

PROFILE_NAME = 'resources_dev'
ROLE_NAME = 'ResourcesEnvMainRole'
FUNCTION_NAME = 'func_test_01'

LOG_GROUP_NAME = f"/aws/lambda/{FUNCTION_NAME}"

session = boto3.Session(profile_name=PROFILE_NAME)
client = session.client("logs")

SINCE_MINUTES = 10
JST = timezone(timedelta(hours=9))  # JST (UTC+9)

start_time = int(
    (datetime.now(JST) - timedelta(minutes=SINCE_MINUTES)).timestamp() * 1000
)

paginator = client.get_paginator("filter_log_events")

for page in paginator.paginate(
    logGroupName=LOG_GROUP_NAME, startTime=start_time
):
    for event in page.get("events", []):
        dt = datetime.fromtimestamp(event["timestamp"] / 1000, JST)
        print(f"[{dt.isoformat()}] {event['message'].rstrip()}")

lambda - SNS SQS 連携の一括作成(SQS作成/SNS作成/権限付与/サブクライブ)

lambda - SNS SQS 連携を一括作成します

import json
import boto3
import time
from datetime import datetime, timedelta, timezone
from botocore.exceptions import ClientError

PROFILE_NAME_ADMIN = 'resources_dev_admin'
ROLE_NAME = 'ResourcesEnvMainRole'
FUNCTION_NAME = 'func_test_01'

QUEUE_NAME = f"{FUNCTION_NAME}-sqs"
TOPIC_NAME = f"{FUNCTION_NAME}-sns"

POLICY_ARNS = [
    'arn:aws:iam::aws:policy/AmazonSQSFullAccess',
    'arn:aws:iam::aws:policy/AmazonSNSFullAccess'
]

session = boto3.Session(profile_name=PROFILE_NAME_ADMIN)
iam_client = session.client('iam')

print(f"Starting to attach policies to role '{ROLE_NAME}' using profile '{PROFILE_NAME_ADMIN}'...")

for policy_arn in POLICY_ARNS:
    try:
        iam_client.attach_role_policy(
            RoleName=ROLE_NAME,
            PolicyArn=policy_arn
        )
        print(f"Success: Attached {policy_arn.split('/')[-1]}.")
    except ClientError as e:
        print(f"Error: Failed to attach {policy_arn.split('/')[-1]}.")
        print(e.response['Error']['Message'])

print("Process completed.")


print("Waiting for IAM propagation...")
time.sleep(10)

sqs = session.client('sqs')

create_response = sqs.create_queue(
    QueueName=QUEUE_NAME,
    Attributes={
        'VisibilityTimeout': '1080'
    }
)

url_response = sqs.get_queue_url(QueueName=QUEUE_NAME)
queue_url = url_response['QueueUrl']

attrs_response = sqs.get_queue_attributes(
    QueueUrl=queue_url,
    AttributeNames=['QueueArn']
)
queue_arn = attrs_response['Attributes']['QueueArn']

print(queue_url)
print(queue_arn)

BATCH_SIZE = 1

lambda_client = session.client('lambda')

response = lambda_client.create_event_source_mapping(
    FunctionName=FUNCTION_NAME,
    EventSourceArn=queue_arn,
    BatchSize=BATCH_SIZE
)

print(response['UUID'])

sns = session.client('sns')

response = sns.create_topic(Name=TOPIC_NAME)

print(response['TopicArn'])

aws_region = session.region_name
sts_client = session.client('sts')
aws_account_id = sts_client.get_caller_identity()['Account']

sns_arn = f"arn:aws:sns:{aws_region}:{aws_account_id}:{TOPIC_NAME}"

policy_dict = {
    "Version": "2012-10-17",
    "Statement": [
        {
            "Sid": "Allow-SNS-SendMessage",
            "Effect": "Allow",
            "Principal": "*",
            "Action": "sqs:SendMessage",
            "Resource": queue_arn,
            "Condition": {
                "ArnEquals": {
                    "aws:SourceArn": sns_arn
                }
            }
        }
    ]
}

sqs.set_queue_attributes(
    QueueUrl=queue_url,
    Attributes={
        'Policy': json.dumps(policy_dict)
    }
)

print(f"Policy applied to: {queue_url}")

response = sns.subscribe(
    TopicArn=sns_arn,
    Protocol='sqs',
    Endpoint=queue_arn,
    Attributes={
        'RawMessageDelivery': 'true'
    }
)

print(response['SubscriptionArn'])

MESSAGE = 'SNS message test!!!'
response = sns.publish(
    TopicArn=sns_arn,
    Message=MESSAGE
)

print(response['MessageId'])

LOG_GROUP_NAME = f"/aws/lambda/{FUNCTION_NAME}"

client = session.client("logs")

SINCE_MINUTES = 10
JST = timezone(timedelta(hours=9))

start_time = int(
    (datetime.now(JST) - timedelta(minutes=SINCE_MINUTES)).timestamp() * 1000
)

paginator = client.get_paginator("filter_log_events")

for page in paginator.paginate(
    logGroupName=LOG_GROUP_NAME, startTime=start_time
):
    for event in page.get("events", []):
        dt = datetime.fromtimestamp(event["timestamp"] / 1000, JST)
        print(f"[{dt.isoformat()}] {event['message'].rstrip()}")

SNS SQS の削除(構成が不要になった場合のみ)

import boto3

PROFILE_NAME = 'resources_dev'
ROLE_NAME = 'ResourcesEnvMainRole'
FUNCTION_NAME = 'func_test_01'
QUEUE_NAME = f"{FUNCTION_NAME}-sqs"
TOPIC_NAME = f"{FUNCTION_NAME}-sns"

session = boto3.Session(profile_name=PROFILE_NAME)

aws_region = session.region_name
sts_client = session.client('sts')
aws_account_id = sts_client.get_caller_identity()['Account']

sqs = session.client('sqs')
sns = session.client('sns')

# 1. SQSキューの削除 (get_queue_url から URL を取得)
try:
    url_response = sqs.get_queue_url(QueueName=QUEUE_NAME)
    queue_url = url_response['QueueUrl']
    sqs.delete_queue(QueueUrl=queue_url)
    print(f"Successfully deleted SQS queue: {QUEUE_NAME}")
except sqs.exceptions.QueueDoesNotExist:
    print(f"SQS queue '{QUEUE_NAME}' does not exist. Skipped.")
except Exception as e:
    print(f"Error deleting SQS queue: {e}")

# 2. SNSトピックの削除 (安全に組み立てた ARN を使用)
sns_arn = f"arn:aws:sns:{aws_region}:{aws_account_id}:{TOPIC_NAME}"

try:
    sns.delete_topic(TopicArn=sns_arn)
    print(f"Successfully deleted SNS topic: {TOPIC_NAME}")
except sns.exceptions.NotFoundException:
    print(f"SNS topic '{TOPIC_NAME}' does not exist. Skipped.")
except Exception as e:
    print(f"Error deleting SNS topic: {e}")

関連記事

AWS lambda関数と SQS SNS との連携

更新日:2026年08月15日

0
0
0

Register as a new user and use Qiita more conveniently

  1. You get articles that match your needs
  2. You can efficiently read back useful information
  3. You can use dark theme
What you can do with signing up
0
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?