メインコンテンツまでスキップ

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日

ITとソフトウェアの人気オンラインコースHP Directplus -HP公式オンラインストア-