Project Overview
Developed a robust, event-driven ETL pipeline that monitors filesystem events and automatically processes and uploads data to Azure Data Lake Storage Gen2. The system used YAML configuration files for pipeline definition, making it highly configurable and maintainable.
Business Context
The business needed a flexible solution to continuously monitor specific directories for new data files, process them according to predefined rules, and reliably upload the results to cloud storage. This enabled near real-time data processing without the complexity of a full streaming solution.
Technical Challenges
Challenge 1: Reliable Event Handling
File system events can be triggered multiple times in rapid succession when files are being written, which could lead to processing incomplete files or redundant processing.
Challenge 2: Debouncing File Events
Implementing effective debouncing of filesystem events to prevent multiple processing of the same file while maintaining low latency.
Challenge 3: Schema Validation
Ensuring that all YAML configuration files adhered to a strict schema for consistency and error prevention.
Architecture
flowchart TD
A[File System Events] --> B[Watchdog Observer]
B --> C[Event Handler]
C --> D{Event Type?}
D -->|Created| E[File Created Handler]
D -->|Modified| F[File Modified Handler]
D -->|Deleted| G[File Deleted Handler]
E --> H[Debouncing Logic]
F --> H
H --> I[YAML Config Processor]
I --> J[Schema Validator]
J --> K[Data Transformation]
K --> L[Data Quality Checks]
L --> M[Azure ADLS Upload]
N[YAML Config Files] --> I
O[Schema Definitions] --> J
subgraph "Error Handling"
P[Retry Mechanism]
Q[Dead Letter Queue]
R[Error Logging]
end
M -- Failure --> P
P -- Max Retries --> Q
P -- Retry --> M
Q --> RImplementation Details
Debouncing Implementation
| |
YAML Configuration Schema
| |
Results and Impact
Key Achievements
- Reliability: 99.9% successful file processing rate with built-in retry mechanisms
- Performance: Average processing time of less than 5 seconds per file
- Maintainability: Configuration changes could be made without code modifications
- Scalability: Successfully handled 10,000+ files per day across multiple directories
Business Impact
- Reduced data processing latency from hours to minutes
- Enabled near real-time analytics on incoming data
- Significantly reduced maintenance overhead with configuration-based pipelines
- Improved data quality through automated validation
Lessons Learned
Technical Insights
- Event-based systems need careful debouncing to prevent duplicate processing
- YAML schema validation is essential for configuration-driven applications
- Azure SDK authentication should use managed identities where possible
- Error handling strategies should be explicitly defined in configuration
Process Improvements
- Implementing semantic logging improved troubleshooting capabilities
- Unit tests for YAML schema validation caught many issues early
- Monitoring dashboards for pipeline health provided valuable insights
Future Enhancements
- Implement parallel processing for high-volume scenarios
- Add support for more complex transformation types
- Create a web UI for monitoring and configuration management
- Integrate with Azure Event Grid for more sophisticated event handling
Tools and Technologies Used
- Python: Core programming language
- Watchdog: File system event monitoring
- PyYAML: YAML parsing and generation
- Jsonschema: Schema validation
- Azure Data Lake Storage Gen2 SDK: Cloud storage integration
- Pandas: Data transformation and processing
- Pytest: Testing framework