nyc-taxi
Production-grade streaming ETL pipeline for NYC Yellow Taxi data using Kafka, DDD-based validation and enrichment, ClickHouse for sub-second analytics, and real-time Grafana dashboards.
File Explorer
- ci.yml
- integration.yml
- 001_initial_schema.sql
- 002_materialized_views.sql
- materialized_views.sql
- queries.sql
- schema.sql
- .gitkeep
- taxi_zone_lookup.csv
- .gitkeep
- Dockerfile
- config.xml
- users.xml
- Dockerfile
- etl_overview.json
- dashboard.yml
- datasource.yml
- Dockerfile
- kafka-topics.sh
- prometheus.yml
- backend.hcl
- terraform.tfvars
- backend.hcl
- terraform.tfvars
- backend.hcl
- terraform.tfvars
- aws-load-balancer-controller-iam-policy.json
- main.tf
- outputs.tf
- variables.tf
- main.tf
- outputs.tf
- variables.tf
- main.tf
- outputs.tf
- variables.tf
- .terraform.lock.hcl
- backend.tf
- main.tf
- outputs.tf
- providers.tf
- variables.tf
- versions.tf
- architecture.md
- deployment.md
- event_flow.md
- observability.md
- replay_strategy.md
- scaling_notes.md
- __init__.py
- process_batch.py
- process_trip.py
- replay_dlq.py
- __init__.py
- enrichment_service.py
- ingestion_service.py
- validation_service.py
- __init__.py
- __init__.py
- settings.py
- __init__.py
- models.py
- services.py
- __init__.py
- deduplicator.py
- enrichers.py
- events.py
- exceptions.py
- models.py
- normalizers.py
- repositories.py
- services.py
- validators.py
- __init__.py
- __init__.py
- consumer.py
- health_server.py
- producer.py
- schema_apply.py
- __init__.py
- client.py
- inserter.py
- repository.py
- schema_manager.py
- __init__.py
- consumer.py
- dead_letter_publisher.py
- producer.py
- serializer.py
- topic_manager.py
- __init__.py
- health.py
- kafka_lag.py
- logging.py
- metrics.py
- __init__.py
- parquet_reader.py
- zone_lookup.py
- __init__.py
- __init__.py
- correlation.py
- structured_logging.py
- tracing.py
- __init__.py
- batching.py
- lag_monitor_loop.py
- lifecycle.py
- retry.py
- shutdown.py
- __init__.py
- compression.py
- datetime.py
- hashing.py
- json.py
- __init__.py
- apply_schema.sh
- bootstrap.sh
- check_kafka_lag.sh
- create_topics.sh
- download_data.sh
- replay_dlq.sh
- smoke_test.sh
- invalid_trips.parquet
- sample_trips.parquet
- taxi_zone_lookup.csv
- test_topic_creation.py
- conftest.py
- test_duplicate_pipeline.py
- test_invalid_trip_pipeline.py
- test_replay_dlq_pipeline.py
- test_shutdown_pipeline.py
- test_trip_pipeline.py
- conftest.py
- test_dlq_flow.py
- test_kafka_to_clickhouse.py
- test_materialized_views.py
- test_replay_flow.py
- test_process_batch.py
- test_process_trip.py
- test_replay_dlq.py
- test_enrichment_service.py
- test_ingestion_service.py
- test_validation_service.py
- test_deduplicator.py
- test_normalizers.py
- test_services.py
- test_validators.py
- test_inserter.py
- test_repository.py
- test_schema_manager.py
- test_consumer.py
- test_dead_letter_publisher.py
- test_serializer.py
- test_topic_manager.py
- test_health.py
- test_kafka_lag.py
- test_logging.py
- test_structured_logging.py
- test_batching.py
- test_lag_monitor_loop.py
- test_lifecycle.py
- test_retry.py
- test_shutdown.py
- test_compression.py
- test_datetime.py
- test_hashing.py
- test_json.py
- __init__.py
- .dockerignore
- .env.example
- .gitignore
- docker-compose.yml
- LICENSE
- Makefile
- poetry.lock
- pyproject.toml
- README.md
- requirements.txt
# Use via CDN
jsDelivrjsDelivr serves any public GitHub repository as a CDN with zero setup. Pick a version and a file to get a ready-to-paste link and snippet.
Link
Example
// repository documentation
Was this content helpful?
(0 ratings)
