EventBridge Integration Sending Duplicate Events — Deduplication Strategy

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:PutEvents permissions and s3:GetObject if using S3 as a deduplication store, or dynamodb:PutItem/Dynamodb:GetItem for DynamoDB-based deduplication.
  • SDK Version: boto3>=1.26.0 and botocore>=1.29.0.
  • Runtime Requirements: Python 3.9 or higher.
  • External Dependencies: pydantic for 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 id or source).
  • 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

Official References