Skip to content

Latest commit

 

History

8 Commits

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 

Repository files navigation

Architecture Diagram

Architecture Diagram

graph TD
    subgraph Data_Source ["Data Source"]
        API["Django REST API - /api/occupancy/"]
    end

    subgraph Ingestion_Layer ["Data Ingestion & Streaming"]
        Driver["Python Ingestion Driver - driver.py"]
        Kafka["Aiven Kafka Topic - occupancy_data"]
        Consumer["Kafka Stream Consumer - consumer.py"]
    end

    subgraph Storage_Layer ["Aiven PostgreSQL / TimescaleDB"]
        RawDB[("Table: occupancy_readings (Raw Data)")]
        AggDB[("Table: daily_occupancy (Hourly Aggregates)")]
    end

    subgraph Orchestration_Layer ["Orchestration & Transformation"]
        Airflow["Apache Airflow DAG - hourly_occupancy_aggregation"]
    end

    API -->|"1. Poll sensor data"| Driver
    Driver -->|"2. Produce stream events"| Kafka
    Kafka -->|"3. Consume messages"| Consumer
    Consumer -->|"4. Write raw readings"| RawDB
    Airflow -->|"5. Hourly batch aggregation"| RawDB
    Airflow -->|"6. Store aggregated stats"| AggDB
Loading

About

A poc to send fake timeseries data via kafka to be stored in a Db . And then Airflow dag runs on that data.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages