EventBridge Integration Sending Duplicate Events — Deduplication Strategy
What You Will Build
- A Python-based event consumer that detects and filters duplicate events from AWS EventBridge before processing them.
- This solution uses the AWS SDK for Python (Boto3) and the EventBridge API surface for event ingestion and deduplication logic.
- The programming language covered is Python 3.9+.
Prerequisites
- AWS Credentials: An IAM user or role with
eventbridge:PutEventspermissions ands3:GetObjectif using S3 as a deduplication store, ordynamodb:PutItem/Dynamodb:GetItemfor DynamoDB-based deduplication. - SDK Version:
boto3>=1.26.0andbotocore>=1.29.0. - Runtime Requirements: Python 3.9 or higher.
- External Dependencies:
pydanticfor data validation,hashlib(standard library),uuid(standard library).
Authentication Setup
AWS EventBridge does not use OAuth 2.0 in the traditional web sense. It relies on IAM authentication via AWS Signature Version 4. Boto3 handles this automatically when credentials are configured via environment variables, shared credentials file (~/.aws/credentials), or IAM instance profiles.
Ensure your environment has the following variables set:
export AWS_ACCESS_KEY_ID="AKIAIOSFODNN7EXAMPLE"
export AWS_SECRET_ACCESS_KEY="wJalrXUtnFEMI/K7MDENG/bPxRfiCYEXAMPLEKEY"
export AWS_DEFAULT_REGION="us-east-1"
For local development, install AWS SAM CLI or use aws configure to set up credentials. In production, assign an IAM Role to the EC2 instance, Lambda function, or ECS task running the consumer.
Implementation
Step 1: Define the Event Schema and Deduplication Key
EventBridge events are JSON documents. To deduplicate, you must identify a unique identifier. If the source system provides a unique ID (e.g., messageId, correlationId), use it. If not, generate a hash of the immutable content of the event.
We will define a Pydantic model to validate the incoming event and extract the deduplication key.
import hashlib
import json
import logging
from datetime import datetime, timezone
from typing import Optional, Dict, Any
from pydantic import BaseModel, Field, validator
logger = logging.getLogger(__name__)
class EventBridgeEvent(BaseModel):
"""
Represents a simplified EventBridge event structure.
Real events may have additional fields like 'source', 'account', etc.
"""
id: str = Field(..., description="Unique ID provided by the event source")
source: str = Field(..., description="Source of the event")
detail_type: str = Field(..., description="Type of the event")
detail: Dict[str, Any] = Field(..., description="Payload of the event")
event_time: Optional[datetime] = Field(None, description="Timestamp of the event")
class Config:
validate_assignment = True
@validator('event_time', pre=True)
def parse_event_time(cls, v):
if isinstance(v, str):
try:
# Handle ISO 8601 format
return datetime.fromisoformat(v.replace('Z', '+00:00'))
except ValueError:
return None
return v
def get_dedup_key(self) -> str:
"""
Generates a unique hash for deduplication.
Strategy: Hash the combination of 'id', 'source', and a hash of 'detail'.
"""
# If the source provides a reliable unique ID, use it directly
if self.id and self.source:
# Fallback to content hash if ID is not globally unique
content_hash = hashlib.sha256(
json.dumps(self.detail, sort_keys=True).encode('utf-8')
).hexdigest()
return f"{self.source}:{self.id}:{content_hash}"
# If no ID, hash the entire detail payload
return hashlib.sha256(
json.dumps(self.detail, sort_keys=True).encode('utf-8')
).hexdigest()
Step 2: Implement the Deduplication Store
Deduplication requires a stateful store. For high-throughput systems, DynamoDB is preferred due to its low-latency read/write capabilities and support for Time-To-Live (TTL) to automatically clean up old records.
We will create a DeduplicationStore class using Boto3’s DynamoDB client.
import boto3
from botocore.exceptions import ClientError
from typing import Optional
class DeduplicationStore:
def __init__(self, table_name: str, region: str = "us-east-1", ttl_seconds: int = 3600):
self.dynamodb = boto3.resource('dynamodb', region_name=region)
self.table = self.dynamodb.Table(table_name)
self.ttl_seconds = ttl_seconds
self.ttl_attribute = "ttl"
def is_duplicate(self, dedup_key: str) -> bool:
"""
Checks if the event has already been processed within the TTL window.
Returns True if duplicate, False if new.
"""
try:
response = self.table.get_item(Key={"dedup_key": dedup_key})
item = response.get('Item')
if item:
# Check if the record has expired based on TTL
# Note: DynamoDB TTL is asynchronous, so we check manually for immediate feedback
if self.ttl_attribute in item:
expiry_time = item[self.ttl_attribute]
current_time = int(datetime.now(timezone.utc).timestamp())
if current_time > expiry_time:
# Record expired, treat as new
return False
else:
# Record still valid, treat as duplicate
return True
else:
return True
return False
except ClientError as e:
logger.error(f"DynamoDB GetItem error: {e.response['Error']['Message']}")
# Fail open: allow processing if check fails
return False
def mark_as_processed(self, dedup_key: str) -> bool:
"""
Records the event as processed with a TTL.
Returns True if successful, False otherwise.
"""
try:
current_time = int(datetime.now(timezone.utc).timestamp())
expiry_time = current_time + self.ttl_seconds
self.table.put_item(
Item={
"dedup_key": dedup_key,
"processed_at": current_time,
self.ttl_attribute: expiry_time
}
)
return True
except ClientError as e:
logger.error(f"DynamoDB PutItem error: {e.response['Error']['Message']}")
return False
Step 3: Process Events with Deduplication
Now we combine the event model and the store into a processor. This processor simulates receiving events from EventBridge (e.g., via a Lambda trigger or a self-hosted consumer using the EventBridge PutEvents API for testing).
class EventProcessor:
def __init__(self, dedup_store: DeduplicationStore):
self.dedup_store = dedup_store
def process_event(self, raw_event: Dict[str, Any]) -> Dict[str, Any]:
"""
Validates, deduplicates, and processes an event.
"""
try:
event = EventBridgeEvent(**raw_event)
dedup_key = event.get_dedup_key()
if self.dedup_store.is_duplicate(dedup_key):
logger.info(f"Duplicate event detected: {dedup_key}")
return {"status": "skipped", "reason": "duplicate"}
# Mark as processed BEFORE business logic to prevent race conditions
if not self.dedup_store.mark_as_processed(dedup_key):
logger.warning(f"Failed to mark event as processed: {dedup_key}")
# Depending on strategy, you may want to raise an error here
return {"status": "failed", "reason": "dedup_mark_failure"}
# Business Logic Here
result = self._execute_business_logic(event)
return {"status": "processed", "result": result}
except Exception as e:
logger.error(f"Error processing event: {str(e)}")
return {"status": "error", "reason": str(e)}
def _execute_business_logic(self, event: EventBridgeEvent) -> Dict[str, Any]:
"""
Simulates business logic. Replace with actual integration code.
"""
# Example: Send to another API, update database, etc.
return {
"event_id": event.id,
"source": event.source,
"detail_summary": f"Processed {event.detail_type}"
}
Step 4: Simulate EventBridge Integration
To test this locally, we can simulate receiving events. In a real scenario, this code would run in an AWS Lambda function triggered by EventBridge, or in a container receiving events via a message queue.
def simulate_eventbridge_trigger():
"""
Simulates receiving multiple events, including duplicates.
"""
# Initialize the deduplication store
table_name = "EventDedupTable"
store = DeduplicationStore(table_name=table_name, ttl_seconds=60)
processor = EventProcessor(dedup_store=store)
# Sample events
events = [
{
"id": "evt-123",
"source": "my.app",
"detail_type": "OrderCreated",
"detail": {"orderId": "ORD-001", "amount": 100.00},
"event_time": "2023-10-27T10:00:00Z"
},
# Duplicate event
{
"id": "evt-123",
"source": "my.app",
"detail_type": "OrderCreated",
"detail": {"orderId": "ORD-001", "amount": 100.00},
"event_time": "2023-10-27T10:00:00Z"
},
# Unique event
{
"id": "evt-124",
"source": "my.app",
"detail_type": "OrderCreated",
"detail": {"orderId": "ORD-002", "amount": 200.00},
"event_time": "2023-10-27T10:01:00Z"
}
]
results = []
for event in events:
result = processor.process_event(event)
results.append(result)
print(f"Result: {result}")
return results
if __name__ == "__main__":
# Ensure DynamoDB table exists
# You must create the table 'EventDedupTable' with partition key 'dedup_key' (String)
# and enable TTL on the 'ttl' attribute.
simulate_eventbridge_trigger()
Complete Working Example
Below is the full, copy-pasteable script. This example uses in-memory deduplication for simplicity if DynamoDB is not available, but the production code above uses DynamoDB. For this complete example, we will use a simple in-memory cache to demonstrate the logic without external dependencies.
import hashlib
import json
import logging
from datetime import datetime, timezone
from typing import Dict, Any, Optional
from pydantic import BaseModel, Field, validator
# Configure logging
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)
# In-memory store for demonstration. Replace with DynamoDB/Redis in production.
class InMemoryDedupStore:
def __init__(self, ttl_seconds: int = 60):
self.store: Dict[str, float] = {}
self.ttl_seconds = ttl_seconds
def is_duplicate(self, dedup_key: str) -> bool:
current_time = datetime.now(timezone.utc).timestamp()
if dedup_key in self.store:
expiry_time = self.store[dedup_key]
if current_time < expiry_time:
return True
else:
del self.store[dedup_key]
return False
def mark_as_processed(self, dedup_key: str) -> bool:
current_time = datetime.now(timezone.utc).timestamp()
self.store[dedup_key] = current_time + self.ttl_seconds
return True
class EventBridgeEvent(BaseModel):
id: str = Field(..., description="Unique ID provided by the event source")
source: str = Field(..., description="Source of the event")
detail_type: str = Field(..., description="Type of the event")
detail: Dict[str, Any] = Field(..., description="Payload of the event")
event_time: Optional[datetime] = Field(None, description="Timestamp of the event")
class Config:
validate_assignment = True
@validator('event_time', pre=True)
def parse_event_time(cls, v):
if isinstance(v, str):
try:
return datetime.fromisoformat(v.replace('Z', '+00:00'))
except ValueError:
return None
return v
def get_dedup_key(self) -> str:
if self.id and self.source:
content_hash = hashlib.sha256(
json.dumps(self.detail, sort_keys=True).encode('utf-8')
).hexdigest()
return f"{self.source}:{self.id}:{content_hash}"
return hashlib.sha256(
json.dumps(self.detail, sort_keys=True).encode('utf-8')
).hexdigest()
class EventProcessor:
def __init__(self, dedup_store: InMemoryDedupStore):
self.dedup_store = dedup_store
def process_event(self, raw_event: Dict[str, Any]) -> Dict[str, Any]:
try:
event = EventBridgeEvent(**raw_event)
dedup_key = event.get_dedup_key()
if self.dedup_store.is_duplicate(dedup_key):
logger.info(f"Duplicate event detected: {dedup_key}")
return {"status": "skipped", "reason": "duplicate"}
if not self.dedup_store.mark_as_processed(dedup_key):
logger.warning(f"Failed to mark event as processed: {dedup_key}")
return {"status": "failed", "reason": "dedup_mark_failure"}
result = self._execute_business_logic(event)
return {"status": "processed", "result": result}
except Exception as e:
logger.error(f"Error processing event: {str(e)}")
return {"status": "error", "reason": str(e)}
def _execute_business_logic(self, event: EventBridgeEvent) -> Dict[str, Any]:
return {
"event_id": event.id,
"source": event.source,
"detail_summary": f"Processed {event.detail_type}"
}
def main():
store = InMemoryDedupStore(ttl_seconds=60)
processor = EventProcessor(dedup_store=store)
events = [
{
"id": "evt-123",
"source": "my.app",
"detail_type": "OrderCreated",
"detail": {"orderId": "ORD-001", "amount": 100.00},
"event_time": "2023-10-27T10:00:00Z"
},
{
"id": "evt-123",
"source": "my.app",
"detail_type": "OrderCreated",
"detail": {"orderId": "ORD-001", "amount": 100.00},
"event_time": "2023-10-27T10:00:00Z"
},
{
"id": "evt-124",
"source": "my.app",
"detail_type": "OrderCreated",
"detail": {"orderId": "ORD-002", "amount": 200.00},
"event_time": "2023-10-27T10:01:00Z"
}
]
for event in events:
result = processor.process_event(event)
print(f"Result: {result}")
if __name__ == "__main__":
main()
Common Errors & Debugging
Error: pydantic.ValidationError
- What causes it: The incoming event does not match the expected schema (e.g., missing
idorsource). - How to fix it: Ensure the EventBridge rule or source sends all required fields. Add optional fields to the Pydantic model if they are not guaranteed.
- Code showing the fix:
# Make fields optional if they might be missing
id: Optional[str] = None
Error: DynamoDB ConditionalCheckFailedException
- What causes it: When using conditional writes for deduplication, another process has already written the record.
- How to fix it: This is expected behavior in a distributed system. Handle it by treating the event as a duplicate and skipping processing.
- Code showing the fix:
try:
self.table.put_item(
Item={"dedup_key": dedup_key, "processed_at": current_time},
ConditionExpression="attribute_not_exists(dedup_key)"
)
except ClientError as e:
if e.response['Error']['Code'] == 'ConditionalCheckFailedException':
return True # Duplicate detected
Error: High Latency in Deduplication Check
- What causes it: DynamoDB throttling or network latency.
- How to fix it: Use provisioned capacity with autoscaling or on-demand capacity. Implement retry logic with exponential backoff for
ThrottlingException. - Code showing the fix:
import time
from botocore.exceptions import ClientError
def mark_as_processed_with_retry(self, dedup_key: str, retries: int = 3) -> bool:
for attempt in range(retries):
try:
return self.mark_as_processed(dedup_key)
except ClientError as e:
if e.response['Error']['Code'] == 'ThrottlingException':
time.sleep(2 ** attempt) # Exponential backoff
else:
raise
return False