Skip to content

Latest commit

Β 

History

2 Commits

Folders and files

NameName
Last commit message
Last commit date
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 

Repository files navigation

NotificationHub

A comprehensive notification processing system built with Python, FastAPI, and Kafka for scalable, real-time event-driven notifications across multiple channels.

πŸ—οΈ System Architecture

NotificationHub consists of two main services that work together to process events and deliver notifications:

β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”    β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”    β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
β”‚   Event Sources │───▢│  Core Service   │───▢│ Channels Serviceβ”‚
β”‚ (order, payment,β”‚    β”‚ (transformation β”‚    β”‚ (email, sms,    β”‚
β”‚ user events)    β”‚    β”‚ & routing)      β”‚    β”‚ webhook, slack) β”‚
β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜    β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜    β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
                              β”‚                        β”‚
                              β–Ό                        β–Ό
                       β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”    β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
                       β”‚   Kafka Topics  β”‚    β”‚  Notification   β”‚
                       β”‚ (input events)  β”‚    β”‚    Delivery     β”‚
                       β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜    β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜

πŸ”„ How It Works

  1. Event Ingestion: External systems send events (orders, payments, user actions) to Kafka input topics
  2. Event Processing: Core Service consumes events, transforms them into notifications, and applies routing rules
  3. Notification Delivery: Channels Service consumes notifications and delivers them through appropriate channels
  4. Multi-Channel Support: Each notification channel (email, SMS, webhook, Slack) runs as an independent consumer

πŸ“ Project Structure

NotificationHub/
β”œβ”€β”€ main.py                                    # Main application entry point
β”œβ”€β”€ requirements.txt                           # Python dependencies
β”œβ”€β”€ kafka-local/                              # Local Kafka setup
β”‚   β”œβ”€β”€ docker-compose.yml                    # Kafka cluster configuration
β”‚   └── scripts/                              # Kafka management scripts
β”œβ”€β”€ services/
β”‚   β”œβ”€β”€ notification-core-service/            # Event processing service
β”‚   β”‚   β”œβ”€β”€ app/
β”‚   β”‚   β”‚   β”œβ”€β”€ main.py                       # Core service entry point
β”‚   β”‚   β”‚   β”œβ”€β”€ core/                         # Business logic components
β”‚   β”‚   β”‚   β”‚   β”œβ”€β”€ consumer.py               # Kafka event consumer
β”‚   β”‚   β”‚   β”‚   β”œβ”€β”€ producer.py               # Kafka notification producer
β”‚   β”‚   β”‚   β”‚   └── transformer.py            # Event-to-notification transformer
β”‚   β”‚   β”‚   β”œβ”€β”€ config/                       # Configuration management
β”‚   β”‚   β”‚   └── utils/                        # Utility functions
β”‚   β”‚   β”œβ”€β”€ config.env                        # Service configuration
β”‚   β”‚   └── requirements.txt                  # Service dependencies
β”‚   └── notification-channels-service/        # Notification delivery service
β”‚       β”œβ”€β”€ consumers/                        # Channel-specific consumers
β”‚       β”‚   β”œβ”€β”€ email_consumer.py             # Email notification consumer
β”‚       β”‚   β”œβ”€β”€ sms_consumer.py               # SMS notification consumer
β”‚       β”‚   β”œβ”€β”€ webhook_consumer.py           # Webhook notification consumer
β”‚       β”‚   └── slack_consumer.py             # Slack notification consumer
β”‚       β”œβ”€β”€ common/                           # Shared utilities
β”‚       β”œβ”€β”€ config.env                        # Service configuration
β”‚       β”œβ”€β”€ docker-compose.yml                # Multi-container setup
β”‚       └── requirements.txt                  # Service dependencies
└── tests/                                    # Integration tests
    β”œβ”€β”€ e2e/                                  # End-to-end tests
    β”œβ”€β”€ run_e2e_test.py                       # Test runner
    └── requirements.txt                      # Test dependencies

πŸš€ Quick Start

Prerequisites

  • Docker & Docker Compose (for Kafka)
  • Python 3.10+
  • Git

1. Clone and Setup

git clone <repository-url>
cd NotificationHub

# Create virtual environment
python3 -m venv venv
source venv/bin/activate  # On Windows: venv\Scripts\activate

# Install main dependencies
pip install -r requirements.txt

2. Start Kafka Cluster

cd kafka-local
docker-compose up -d

This starts:

  • Zookeeper (port 2181)
  • Kafka broker (port 9094)
  • Kafka UI (port 8080)

3. Start Core Service

cd services/notification-core-service
pip install -r requirements.txt
python app/main.py

4. Start Channels Service

cd services/notification-channels-service
pip install -r requirements.txt
docker-compose up -d

This starts all notification consumers:

  • Email consumer
  • SMS consumer
  • Webhook consumer
  • Slack consumer

5. Verify System

# Run integration tests
cd tests
pip install -r requirements.txt
python run_e2e_test.py

πŸ”§ Configuration

Kafka Topics

The system uses these Kafka topics:

Input Topics (Core Service consumes):

  • order.events - Order-related events
  • payment.events - Payment-related events
  • user.events - User-related events

Output Topics (Channels Service consumes):

  • notifications.email - Email notifications
  • notifications.sms - SMS notifications
  • notifications.webhook - Webhook notifications
  • notifications.slack - Slack notifications

Event Routing Rules

Event Type Channels
order.created email, webhook, slack
order.cancelled email, sms
payment.success email
payment.failed email, sms
user.registered email

πŸ“Š Monitoring

Service Health

  • Core Service: Check logs for event processing
  • Channels Service: Check Docker container status
  • Kafka: Access Kafka UI at http://localhost:8080

Logs

# Core Service logs
tail -f services/notification-core-service/logs/app.log

# Channels Service logs
docker-compose -f services/notification-channels-service/docker-compose.yml logs -f

πŸ§ͺ Testing

Integration Tests

cd tests
python run_e2e_test.py

Tests verify:

  • Event ingestion from input topics
  • Event transformation and routing
  • Notification delivery through all channels
  • End-to-end message flow

Manual Testing

# Send test events to Kafka
cd kafka-local/scripts
./producer.sh order.events '{"event_type":"order.created","user_id":"test123","order_id":"ORD001","amount":99.99}'

πŸ› οΈ Development

Adding New Channels

  1. Create consumer in services/notification-channels-service/consumers/
  2. Add topic configuration in common/config.py
  3. Update routing rules in Core Service
  4. Add to Docker Compose configuration

Adding New Event Types

  1. Update routing rules in services/notification-core-service/app/config/settings.py
  2. Test with integration tests
  3. Update documentation

πŸ” Troubleshooting

Common Issues

Kafka Connection Failed

# Check Kafka status
docker-compose -f kafka-local/docker-compose.yml ps

# Check logs
docker-compose -f kafka-local/docker-compose.yml logs kafka

Messages Not Consumed

# Check consumer group offsets
cd kafka-local/scripts
./consumer.sh notifications.email

Services Not Starting

# Check dependencies
pip install -r requirements.txt

# Check configuration
cat services/*/config.env

πŸ“ˆ Performance

  • Throughput: Handles thousands of events per second
  • Latency: Sub-second notification delivery
  • Scalability: Horizontal scaling via Kafka partitioning
  • Reliability: At-least-once delivery guarantees

πŸ”’ Security

  • Environment-based configuration
  • Kafka SASL authentication support
  • Input validation on all events
  • Secure Docker container deployment

🀝 Contributing

  1. Fork the repository
  2. Create a feature branch
  3. Add tests for new functionality
  4. Ensure all tests pass
  5. Submit a pull request

πŸ“„ License

MIT License - see LICENSE file for details.


Built with ❀️ using Python, FastAPI, and Kafka

About

No description, website, or topics provided.

Resources

Stars

1 star

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages