
API 的部分使用 API Gateway ,當然也可以使用 ALB ,需要設定 path 和要 trigger 的 AWS 服務。當 API Gateway 收到 AWS 服務的 response 時,如果 API 是同步的,可以直接回傳收到的結果,但如果 API 是非同步的,舉例來說,只是把 request 送到 SQS,最後在 API Gateway 設定要回什麼給 client 。
AdClickApi:
Type: AWS::ApiGateway::RestApi
Properties:
Name: !Sub 'ad-click-api-${Environment}'
Description: API for receiving ad click events
EndpointConfiguration:
Types:
- REGIONAL
ClicksResource:
Type: AWS::ApiGateway::Resource
Properties:
RestApiId: !Ref AdClickApi
ParentId: !GetAtt AdClickApi.RootResourceId
PathPart: clicks
ClicksPostMethod:
Type: AWS::ApiGateway::Method
Properties:
RestApiId: !Ref AdClickApi
ResourceId: !Ref ClicksResource
HttpMethod: POST
AuthorizationType: NONE
Integration:
Type: AWS
IntegrationHttpMethod: POST
Uri: !Sub 'arn:aws:apigateway:${AWS::Region}:sqs:path/${AWS::AccountId}/${AdClickQueue.QueueName}'
Credentials: !GetAtt ApiGatewaySQSRole.Arn
RequestTemplates:
application/json: 'Action=SendMessage&MessageBody=$input.body'
IntegrationResponses:
- StatusCode: 200
ResponseTemplates:
application/json: '{"message": "Click event queued successfully"}'
這邊有設定 DLQ ,方便查看失敗的 message 有什麼問題。
AdClickQueue:
Type: AWS::SQS::Queue
Properties:
QueueName: !Sub 'ad-click-queue-${Environment}'
VisibilityTimeout: 300
MessageRetentionPeriod: 1209600 # 14 days
ReceiveMessageWaitTimeSeconds: 20
RedrivePolicy:
deadLetterTargetArn: !GetAtt AdClickDLQ.Arn
maxReceiveCount: 3
AdClickDLQ:
Type: AWS::SQS::Queue
Properties:
QueueName: !Sub 'ad-click-dlq-${Environment}'
MessageRetentionPeriod: 1209600 # 14 days
Lambda 的部分,需要撰寫程式碼把資料寫到 dynamodb ,如果 API 是非同步的,寫完之後就不需要回傳結果,最後的 return 不一定是必要的。
AdClickProcessorFunction:
Type: AWS::Lambda::Function
Properties:
FunctionName: !Sub 'ad-click-processor-${Environment}'
Runtime: python3.11
Handler: index.lambda_handler
Role: !GetAtt LambdaExecutionRole.Arn
Timeout: 60
MemorySize: 256
Environment:
Variables:
DYNAMODB_TABLE: !Ref AdClickTable
ENVIRONMENT: !Ref Environment
Code:
ZipFile: |
import json
import boto3
import os
import uuid
from datetime import datetime
from decimal import Decimal
dynamodb = boto3.resource('dynamodb')
table = dynamodb.Table(os.environ['DYNAMODB_TABLE'])
def lambda_handler(event, context):
print(f"Processing {len(event['Records'])} records")
processed = 0
failed = 0
for record in event['Records']:
try:
# Parse the message body
body = json.loads(record['body'])
print(f"Received body: {json.dumps(body)}")
# Extract ad click data with new field names
ad_id = body.get('ad_id')
impression_id = body.get('impression_id')
user_id = body.get('user_id', 'anonymous')
timestamp_str = body.get('timestamp')
source_url = body.get('source_url')
# Validate required fields
if not ad_id:
raise ValueError("ad_id is required")
if not impression_id:
raise ValueError("impression_id is required")
# Generate unique click_id using UUID
click_id = str(uuid.uuid4())
# Convert ISO timestamp to Unix timestamp
if timestamp_str:
dt = datetime.fromisoformat(timestamp_str.replace('Z', '+00:00'))
timestamp = int(dt.timestamp())
else:
timestamp = int(datetime.now().timestamp())
# Prepare item for DynamoDB
item = {
'ad_id': ad_id, # Partition Key (HASH)
'click_id': click_id, # Sort Key (RANGE)
'impression_id': impression_id,
'user_id': user_id,
'timestamp': timestamp,
'timestamp_iso': timestamp_str,
'source_url': source_url,
'processedAt': datetime.now().isoformat()
}
# Store in DynamoDB
table.put_item(Item=item)
processed += 1
print(f"Successfully processed click for ad: {ad_id}, impression: {impression_id}")
except Exception as e:
failed += 1
print(f"Error processing record: {str(e)}")
print(f"Record: {record}")
return {
'statusCode': 200,
'body': json.dumps({
'processed': processed,
'failed': failed,
'total': len(event['Records'])
})
}
Tags:
- Key: Environment
Value: !Ref Environment
- Key: Application
Value: AdClickAggregator
今天先簡單介紹 API Gateway 、 SQS 和 Lambda 這一段,明天再接著介紹 DynamoDB Table 和 IAM Role 。