AWS lambda関数と SQS SNS との連携
lambda関数と SQS SNS との連携手順を解説します。
※ 操作はすべてAWS Python ライブラリ(boto3)を使用して実施します。
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 organization を利用して新しいAWS環境をゼロから作成する
- AWS lambda関数の作成から実行まで
- AWS lambda関数と SQS SNS との連携
- AWS APIGateway からの SNS 連携
- AWS APIGateway からの lambda 連携
- AWS Lambda のソースを更新する
- AWS Lambda でLine bot を作成する
- AWS Lambda で Telegram bot を作成する
AWS lambda関数と SQS SNS との連携
更新日:2026年08月15日