こんにちは
今回は直前に書いた記事に追加で DynamoDB の構成を追加し、同一のIPがブロックされた場合でも、1通のみメールが送信されるようにしてみました
内容としては、DynamoDB内で当該ipを保持し、もし保持していた場合は送信をスキップするようにします
■構成
以前の構成にDynamoDBを追加したのみであり、基本的なレートベースルールや通知構成は全く同じなので、全体の構成の説明は割愛します
■DynamoDBの挙動
ipをパーティションキーとしてテーブルを作成し、通知前に attribute_not_exists 条件付きで書き込みます
そのうえで、すでに同じIPアドレスが存在していれば、ConditionalCheckFailedException を返し、メール通知をスキップしてみます
また、TTLを追加で設定し、TTLが失効した場合はDynamoDB側でキーを削除し、新規で登録できるようにします
■やってみた
1. DynamoDBの構築
まずは、DynamoDBを構築します
ipをパーティションキーとして作成し、expire_atの名前でTTLを有効化しました
2. IAMポリシーとロールの修正
次に Lambda からDynamoDBへ接続するために、LambdaのIAMロールへDynamoDB用ポリシーを追加します
下記を作成し、LambdaのIAMロールへ紐づけます
{
"Version": "2012-10-17",
"Statement": [
{
"Sid": "VisualEditor0",
"Effect": "Allow",
"Action": "dynamodb:PutItem",
"Resource": "*"
}
]
}
3. Lambda 関数の修正
既存のコードを下記へ書き換えます
import base64
import gzip
import json
import os
import time
from datetime import datetime, timedelta, timezone
import boto3
from botocore.exceptions import ClientError
sns = boto3.client("sns")
dynamodb = boto3.client("dynamodb")
SNS_TOPIC_ARN = os.environ["SNS_TOPIC_ARN"]
DEDUP_TABLE_NAME = os.environ["DEDUP_TABLE_NAME"]
COOLDOWN_SECONDS = int(os.environ.get("COOLDOWN_SECONDS", "300"))
# 表示タイムゾーン(既定は日本時間 JST = UTC+9)
_TZ_OFFSETS = {"Asia/Tokyo": 9, "UTC": 0}
_TZ_NAME = os.environ.get("TIMEZONE", "Asia/Tokyo")
_TZ = timezone(timedelta(hours=_TZ_OFFSETS.get(_TZ_NAME, 9)), _TZ_NAME)
# ISO 3166-1 alpha-2 → 日本語国名(主要国。未登録はコードをそのまま表示)
COUNTRY_NAMES = {
"JP": "日本", "US": "アメリカ合衆国", "CN": "中国", "KR": "韓国",
"TW": "台湾", "HK": "香港", "SG": "シンガポール", "IN": "インド",
"TH": "タイ", "VN": "ベトナム", "ID": "インドネシア", "PH": "フィリピン",
"MY": "マレーシア", "RU": "ロシア", "GB": "イギリス", "DE": "ドイツ",
"FR": "フランス", "NL": "オランダ", "BR": "ブラジル", "CA": "カナダ",
"AU": "オーストラリア", "UA": "ウクライナ", "TR": "トルコ", "IR": "イラン",
}
def _country_label(code):
"""国コードを『日本 (JP)』の形式に整形。未登録・不明はコードのみ。"""
if not code:
return "不明"
name = COUNTRY_NAMES.get(code.upper())
return f"{name} ({code})" if name else code
def _format_datetime(epoch_ms):
"""エポックミリ秒を表示タイムゾーンの文字列に整形。"""
if epoch_ms is None:
return "不明"
dt = datetime.fromtimestamp(epoch_ms / 1000, tz=_TZ)
return dt.strftime("%Y-%m-%d %H:%M:%S ") + _TZ_NAME
def _iter_waf_logs(event):
"""CloudWatch Logs サブスクリプションイベントを展開し、各 WAF ログ(dict)を返す。"""
data = event.get("awslogs", {}).get("data")
if not data:
return
decoded = gzip.decompress(base64.b64decode(data))
payload = json.loads(decoded)
for log_event in payload.get("logEvents", []):
try:
yield json.loads(log_event["message"])
except (KeyError, json.JSONDecodeError):
continue
def _try_reserve_ip(ip, now_epoch):
"""
DynamoDB に IP を条件付きで登録できたら True(=今回通知すべき)。
既に登録済み(クールダウン中)なら False。
TTL 属性 expire_at に「現在時刻 + COOLDOWN_SECONDS」を入れ、
期限切れは DynamoDB が自動削除する。念のため、期限切れレコードが
まだ残っている場合も上書きできるよう条件式に expire_at 比較を加える。
"""
expire_at = now_epoch + COOLDOWN_SECONDS
try:
dynamodb.put_item(
TableName=DEDUP_TABLE_NAME,
Item={
"ip": {"S": ip},
"expire_at": {"N": str(expire_at)},
"last_notified": {"N": str(now_epoch)},
},
# レコードが無い、または既存レコードの TTL が既に過ぎている場合のみ書き込む
ConditionExpression="attribute_not_exists(ip) OR expire_at < :now",
ExpressionAttributeValues={":now": {"N": str(now_epoch)}},
)
return True
except ClientError as e:
if e.response["Error"]["Code"] == "ConditionalCheckFailedException":
return False # クールダウン中 → スキップ
raise
def _build_message(waf_log):
"""WAF ログ 1 件から SNS 送信用の件名・本文を生成。"""
http_request = waf_log.get("httpRequest", {})
client_ip = http_request.get("clientIp", "不明")
country = _country_label(http_request.get("country"))
blocked_at = _format_datetime(waf_log.get("timestamp"))
rule_name = waf_log.get("terminatingRuleId", "不明")
subject = f"[WAF] レートベースルールでブロックしました (IP: {client_ip})"
body = (
"AWS WAF のレートベースルールにより、以下の通信をブロックしました。\n"
"\n"
"----------------------------------------\n"
f"ブロック日時 : {blocked_at}\n"
f"IP アドレス : {client_ip}\n"
f"国名 : {country}\n"
f"ルール名 : {rule_name}\n"
"----------------------------------------\n"
"\n"
"内容をご確認ください。"
)
return subject, body
def lambda_handler(event, context):
now_epoch = int(time.time())
published = 0
skipped = 0
seen_in_batch = set() # このバッチ内で処理済みの IP
for waf_log in _iter_waf_logs(event):
client_ip = waf_log.get("httpRequest", {}).get("clientIp", "不明")
# 1. バッチ内重複排除
if client_ip in seen_in_batch:
skipped += 1
continue
seen_in_batch.add(client_ip)
# 2. DynamoDB による確実な重複排除(クールダウン中はスキップ)
if not _try_reserve_ip(client_ip, now_epoch):
skipped += 1
print(f"Skip (cooldown): {client_ip}")
continue
subject, body = _build_message(waf_log)
sns.publish(
TopicArn=SNS_TOPIC_ARN,
# SNS の件名は最大 100 文字・ASCII 制約があるため安全な固定件名にする
Subject="[WAF] Rate-based rule blocked an IP",
Message=body,
)
published += 1
print(f"Published: {subject}")
return {"statusCode": 200, "published": published, "skipped": skipped}
→ DEDUP_TABLE_NAME のような環境変数を追加しており、lambda_handler ですでにIPが存在している場合は、メール送信をスキップするようにしています
次に、環境変数へ SNS_TOPIC_ARN、DEDUP_TABLE_NAME、COOLDOWN_SECONDS を設定します
4. 動作確認
それでは、最後に動作確認してみましょう
まずは、同一IPからF5連打し、レートベースのブロックを発生させたのち、DynamoDB にキーが登録されているかどうか確認します
→ IP自体はマスクしておりますが、キーとして登録されていました
最後にメールを確認したところ、下記のように1通のみ通知されていることを確認したので、重複して通知されていないのでOKですね
※CloudWatch Logs 上では下記のように大量の同一IPのブロックログが出力されていました
■最後に
いかがでしたでしょうか
DynamoDBを使うだけで、同一IPであれば、メールが重複して送信されないようにできました
そして、DnyamoDBの設定もそこまで複雑ではなく、Lamabda側の回収もそこまで多くはなかったので、そこも地味にうれしいポイントかなと思います
また、今回初めてがっつりDyanmoDBを利用したのですが、何らかの管理をする点においてかなり便利だと感じ、terraformのstateロックにも採用される理由が分かりました
次は DynamoDBに焦点を当てた何らかの記事を書いてみたいと思いまし





