A real-time clinical data pipeline that simulates streaming EHR events from MIMIC-IV. Patient events are validated, transformed, scored for severity using a trained Random Forest classifier, and written to Snowflake for future downstream analytics.
Streaming alternative to traditional overnight EHR batch processing using SQL, Python, and Snowflake.
| Feature Group | Features |
|---|---|
| Vitals | Heart rate, temperature, respiratory rate, vitals sum, vitals range |
| Labs | Lab result, lab-present flag |
| Clinical Notes | Note length, word count, average word length |
| Categorical | Department, event type |
The model then uses these engineered features to estimate the probability of a severe event.
stream_producer.py simulates a real-time stream of patient events from the
MIMIC-IV dataset and publishes them to Kafka.
realtime_etl.py consumes each event and validates and enriches (later assigned a severity score).
| subject_id | admit_type | vitals_score | severity | p_severe |
|---|---|---|---|---|
| 1001 | EMERGENCY | 2.84 | severe | 0.87 |
| 1002 | ELECTIVE | 0.91 | non-severe | 0.08 |
| 1003 | EMERGENCY | 2.17 | severe | 0.71 |
| 1004 | URGENT | 1.42 | non-severe | 0.29 |
| 1005 | EMERGENCY | 0.63 | non-severe | 0.14 |
p_severe represents the model's predicted probability of the severe class,
while severity represents the resulting binary classification.
- Python — ETL, feature engineering, model inference
- Apache Kafka — event streaming
- scikit-learn — Random Forest classification
- Snowflake — analytical data warehouse
- Pandas — data preparation and transformation
pip install -r requirements.txtAdd the required Snowflake connection settings to config.py:
SNOWFLAKE_ACCOUNT
SNOWFLAKE_USER
SNOWFLAKE_PASSWORD
SNOWFLAKE_DATABASE
SNOWFLAKE_SCHEMA
SNOWFLAKE_WAREHOUSE
Run the training pipeline to generate the serialized model artifact:
python train_model.pypython stream_producer.pyThen start the Kafka consumer and scoring pipeline:
python realtime_etl.pyPotential extensions include:
- Schema enforcement with a formal event contract
- Dead-letter Kafka topic for failed events
- Model version tracking and registry integration
- Automated data-quality monitoring
- Kafka consumer lag monitoring
- Snowflake incremental transformation layer
- Model performance and drift monitoring
- Containerized deployment with Docker
- CI/CD and automated pipeline tests
MIT