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'])