इस प्रोजेक्ट के बारे में
यह प्रोजेक्ट एक पूरी तरह से Docker-आधारित, स्थानीय रियल-टाइम स्ट्रीमिंग एनालिटिक्स पाइपलाइन प्रदान करता है जो कोई भी क्लाउड खाता किए बिना Spark Structured Streaming की अवधारणाओं को प्रदर्शित करती है।
## यह क्या करता है
एक Python प्रोड्यूसर Redpanda टॉपिक (Kafka-संगत) में सिम्युलेटेड क्लिकस्ट्रीम इवेंट भेजता है, और एक PySpark Structured Streaming जॉब उन इवेंट्स को कॉन्जूम करता है, लैट इवेंट्स के लिए वॉटरमार्किंग के साथ 1-मिनट के टंबलिंग विंडो एग्रीग्रेशन लागू करता है, और परिणामों को कंसोल या Parquet फ़ाइलों में आउटपुट करता है।
## आर्किटेक्चर
- **Python प्रोड्यूसर**: URL और टाइमस्टैम्प डेटा वाले सिम्युलेटेड क्लिकस्ट्रीम JSON इवेंट जनरेट करता है
- **Redpanda**: Docker में चलने वाला Kafka-संगत ब्रोकर, कोई Zookeeper आवश्यक नहीं
- **PySpark Structured Streaming**: Kafka से एक अनबाउंडेड DataFrame के रूप में पढ़ता है, स्पष्ट स्कीमा के साथ JSON विघटित करता है, event-time विंडोइंग वॉटरमार्किंग के साथ लागू करता है, और प्रति विंडो प्रति URL पेज-व्यू गणना का एग्रीग्रेशन करता है
## प्रमुख अवधारणाएँ
- Spark Structured Streaming के साथ Kafka टॉपिक को एक अनबाउंडेड DataFrame के रूप में पढ़ना
- schema-on-read के साथ JSON इवेंट पेलोड का विघटन
- लैट-आगमन इवेंट्स को हैंडल करने के लिए वॉटरमार्किंग के साथ event-time विंडोइंग (1-मिनट के टंबलिंग विंडो)
- प्रति URL पेज व्यू का निरंतर माइक्रो-बैच एग्रीग्रेशन
- शिक्षण के लिए कंसोल में आउटपुट, स्थायीकरण के लिए Parquet सिंक विकल्प के साथ
## प्रोजेक्ट संरचना
```
streaming-pipeline-kafka-spark/
├── docker-compose.yml # Redpanda + Spark कंटेनर्स
├── producer/
│ └── click_event_producer.py # सिम्युलेटेड क्लिकस्ट्रीम प्रोड्यूसर
├── spark_app/
│ └── streaming_aggregation.py # Structured Streaming जॉब
└── requirements.txt
```
## शुरुआत कैसे करें
```bash
git clone https://github.com/Kornelius99/streaming-pipeline-kafka-spark.git
cd streaming-pipeline-kafka-spark
docker-compose up -d
# टर्मिनल 1: इवेंट प्रोड्यूस करना शुरू करें
python producer/click_event_producer.py
# टर्मिनल 2: Spark Streaming जॉब चलाएं
docker-compose exec spark spark-submit \
--packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.1 \
/opt/spark_app/streaming_aggregation.py
```
परिणाम 1-मिनट के विंडो-आधारित पेज-व्यू गणना के रूप में कंसोल पर दिखाई देते हैं, जब नए इवेंट आते हैं तो अपडेट होते रहते हैं।
## Redpanda क्यों
Redpanda Kafka प्रोटोकॉल का उपयोग करता है इसलिए मौजूदा Kafka क्लाइंट कोड बिना किसी बदलाव के काम करता है, लेकिन यह Zookeeper निर्भरता के बिना सिंगल कंटेनर के रूप में आता है, जिससे सीखने के उद्देश्य के लिए docker-compose सरल रहता है।
## लाइसेंस
MIT
Comments
0 Rating appears after 10 ratings
Sign in to join the discussion.