इस प्रोजेक्ट के बारे में

यह प्रोजेक्ट एक पूरी तरह से 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