Data Flow¶
How data moves through the Fleet Decision Platform from ingestion to output.
End-to-End Data Pipeline¶
flowchart LR
subgraph Sources["Data Sources"]
NYC[NYC Taxi API/Files]
NASA[NASA Dataset]
SIM[Simulation]
end
subgraph Ingestion["Ingestion Layer"]
DOWNLOAD[Download]
PARSE[Parse/Clean]
VALIDATE[Validate]
end
subgraph Storage["Storage"]
RAW[(Raw/Parquet)]
PROC[(Processed)]
MODEL[(Models)]
end
subgraph Transform["Transformation"]
AGG[Aggregate]
FEAT[Feature Engineering]
SPLIT[Train/Test Split]
end
subgraph ML["ML Pipeline"]
TRAIN[Training]
PREDICT[Prediction]
EVAL[Evaluation]
end
subgraph Output["Output"]
ALLOC[Allocation Plan]
KPI[KPIs]
EXPLAIN[Explanations]
end
Sources --> Ingestion
Ingestion --> RAW
RAW --> Transform
Transform --> PROC
PROC --> ML
ML --> MODEL
MODEL --> Output
Data Formats¶
Raw Data¶
| Source | Format | Location |
|---|---|---|
| NYC Taxi | CSV/Parquet | data/raw/nyc_taxi/ |
| NASA Turbofan | CSV | data/raw/nasa_turbofan/ |
| Contracts | data/raw/contracts/ |
Processed Data¶
| Data Type | Format | Location |
|---|---|---|
| Aggregated Demand | Parquet | data/processed/demand/ |
| Fleet State | Parquet | data/processed/fleet_state/ |
| Network Costs | NumPy (.npy) | data/processed/network/ |
| Features | Parquet | data/processed/features/ |
Model Artifacts¶
| Artifact | Format | Location |
|---|---|---|
| XGBoost Model | Pickle/JSON | data/models/demand_forecast/ |
| Risk Model | Pickle | data/models/risk_scoring/ |
| Metadata | JSON | data/models/*/metadata.json |
Stage-by-Stage Flow¶
Stage 1: Data Ingestion¶
graph LR
A[External Source] -->|Download| B[Raw Files]
B -->|Parse| C[DataFrames]
C -->|Validate| D[Clean Data]
D -->|Save| E[(Parquet)]
Input: External APIs/files
Output: Clean Parquet files in data/raw/
# Example: NYC Taxi ingestion
from src.data.ingestion import DataIngestion
ingestion = DataIngestion(config)
raw_data = ingestion.load_nyc_taxi()
# Returns: DataFrame with columns [pickup_datetime, dropoff_datetime,
# pickup_location_id, dropoff_location_id, trip_distance, ...]
Stage 2: Feature Engineering¶
graph LR
A[(Raw Data)] -->|Load| B[DataFrame]
B -->|Time Features| C[hour, day, month]
B -->|Lag Features| D[lag_1h, lag_24h]
B -->|Aggregations| E[zone_demand]
C & D & E -->|Combine| F[(Feature Matrix)]
Input: Raw data Parquet files Output: Feature matrix for ML
# Example: Feature engineering
from src.data.feature_engineering import FeatureEngineer
engineer = FeatureEngineer(config)
features = engineer.create_features(raw_data)
# Returns: DataFrame with columns [location_id, timestamp, hour,
# day_of_week, lag_1h, lag_24h, demand, ...]
Stage 3: Model Training¶
graph LR
A[(Features)] -->|Split| B[Train Set]
A -->|Split| C[Test Set]
B -->|Fit| D[Model]
D -->|Evaluate| E[Metrics]
D -->|Save| F[(Model Artifact)]
Input: Feature matrix Output: Trained model + metrics
# Example: Model training
from src.forecasting import ModelTrainer
trainer = ModelTrainer(config)
model, metrics = trainer.train(features)
trainer.save(model, "data/models/demand_forecast/")
# Metrics: {"rmse": 5.2, "mae": 3.8, "mape": 0.12}
Stage 4: Prediction¶
graph LR
A[(New Features)] -->|Load| B[Feature Vector]
C[(Trained Model)] -->|Load| D[Model]
B & D -->|Predict| E[Forecasts]
Input: Features for forecast horizon Output: Demand forecasts per location
# Example: Prediction
from src.forecasting import DemandPredictor
predictor = DemandPredictor(config)
predictor.load_model("data/models/demand_forecast/")
forecasts = predictor.predict(features, horizon_days=7)
# Returns: Dict[str, np.ndarray] - location_id -> hourly forecasts
Stage 5: Optimization¶
graph LR
A[Forecasts] --> D[Optimizer]
B[Fleet State] --> D
C[Constraints] --> D
D -->|Solve| E[Allocation Plan]
E -->|Calculate| F[KPIs]
Input: Forecasts, fleet state, constraints Output: Allocation plan + KPIs
# Example: Optimization
from src.optimization import CascadingOptimizer
optimizer = CascadingOptimizer(config)
result = optimizer.optimize(
demand_forecast=forecasts,
fleet_state=fleet_state,
network_costs=network_costs,
constraints=constraints
)
# Returns: OptimizationResult with allocation_plan, total_cost, kpis
Data Validation¶
Each stage includes validation:
# Schema validation example
from pydantic import BaseModel, validator
class DemandRecord(BaseModel):
location_id: int
timestamp: datetime
demand: float
@validator('demand')
def demand_must_be_positive(cls, v):
if v < 0:
raise ValueError('demand must be non-negative')
return v
Caching Strategy¶
graph TB
A[Request] -->|Check| B{Cache Hit?}
B -->|Yes| C[Return Cached]
B -->|No| D[Compute]
D -->|Store| E[(Cache)]
D --> F[Return Result]
Caching is implemented at:
- Data layer: Processed data cached in Parquet
- Model layer: Model artifacts cached to disk
- API layer: Response caching with Redis (Phase 2+)
Next Steps¶
- Module Design - Detailed module architecture
- API Reference - API data formats
- Data Formats Reference - Complete schema docs