Skip to content

Latest commit

 

History

24 Commits

Folders and files

Repository files navigation

Clinical Data Streaming Pipeline (MIMIC-IV)

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.

Overview

Streaming alternative to traditional overnight EHR batch processing using SQL, Python, and Snowflake.

Features

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).

Example Output

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.

Tech Stack

  • Python — ETL, feature engineering, model inference
  • Apache Kafka — event streaming
  • scikit-learn — Random Forest classification
  • Snowflake — analytical data warehouse
  • Pandas — data preparation and transformation

Setup

1. Install Dependencies

pip install -r requirements.txt

2. Configure Snowflake

Add the required Snowflake connection settings to config.py:

SNOWFLAKE_ACCOUNT
SNOWFLAKE_USER
SNOWFLAKE_PASSWORD
SNOWFLAKE_DATABASE
SNOWFLAKE_SCHEMA
SNOWFLAKE_WAREHOUSE

3. Model Training

Run the training pipeline to generate the serialized model artifact:

python train_model.py

4. Starting the Streaming Pipeline

python stream_producer.py

Then start the Kafka consumer and scoring pipeline:

python realtime_etl.py

Future Improvements

Potential 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

License

MIT