This project implements a scalable and highly available data processing pipeline using Kubernetes, Kafka, and Neo4j. The pipeline ingests a stream of data (NYC Yellow Cab dataset) and performs near-real-time graph analytics, including PageRank and Breadth-First Search (BFS).
.
├── Dockerfile
├── init.sh
├── kafka-neo4j-connector.yaml
├── kafka-setup.yaml
├── neo4j-service.yaml
├── neo4j-values.yaml
├── zookeeper-setup.yaml
├── sink.neo4j.json
├── interface.py
├── data_producer.py
├── grader.md
└── README.md
- Minikube as Kubernetes orchestrator
- Kafka (using Zookeeper)
- Neo4j with Graph Data Science plugin
- Kafka Connect for Neo4j integration
- NYC Yellow Cab March 2022 dataset
- Docker
- Kubernetes
- Apache Kafka
- Neo4j (Graph Data Science)
kubectl apply -f zookeeper-setup.yaml
kubectl apply -f kafka-setup.yamlhelm install my-neo4j-release neo4j/neo4j -f neo4j-values.yaml
kubectl apply -f neo4j-service.yamlkubectl apply -f kafka-neo4j-connector.yamlkubectl port-forward svc/neo4j-service 7474:7474 7687:7687
kubectl port-forward svc/kafka-service 9092:9092
python3 data_producer.pyImplemented Graph Algorithms:
- PageRank: Determines node importance based on inbound relationships.
- Breadth-First Search (BFS): Finds shortest paths between nodes.
Run analytics through:
python3 interface.py- Kafka Ports: 9092 (external), 29092 (internal)
- Neo4j Ports: 7474 (HTTP), 7687 (Bolt)
- Kafka Connect Port: 8083
- Ensure Minikube has sufficient resources.
- Check logs with
kubectl logs <pod-name>for troubleshooting.
- Kavish Patel

