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

ई-कॉमर्स क्लिकस्ट्रीम एनालिटिक्स प्लेटफ़ॉर्म एक फुल-स्टैक डेटा इंजीनियरिंग प्रोजेक्ट है जो अक्टूबर–नवंबर 2019 के Kaggle डेटासेट से लगभग 285 मिलियन ई-कॉमर्स क्लिकस्ट्रीम इवेंट्स को प्रोसेस करता है। ## अवलोकन प्लेटफ़ॉर्म कच्चे CSV इवेंट डेटा (~14 GB) को S3-कम्पैटिबल ऑब्जेक्ट स्टोर (MinIO) में इनजेस्ट करता है, Apache Spark द्वारा बैच एनालिटिक्स चलाता है, Kafka और Spark Structured Streaming के माध्यम से लाइव स्ट्रीमिंग करता है, Hadoop MapReduce जॉब्स निष्पादित करता है, और FastAPI बैकएंड से React फ्रंटएंड डैशबोर्ड तक परिणाम सेवा करता है। ## आर्किटेक्चर घटक **इनजेस्टन एवं भंडारण:** - MinIO (S3-कम्पैटिबल ऑब्जेक्ट स्टोरेज) कच्चे CSV और Parquet आउटपुट रखता है - PostgreSQL 16.3 एग्रेगेटेड KPIs और क्वेरी-तैयार परिणामों को संग्रहीत करता है - Redis 7.4.1 कैशिंग लेयर प्रदान करता है - Apache Kafka 3.9 (KRaft मोड) स्ट्रीमिंग डेटा के लिए मैसेज बस के रूप में कार्य करता है **प्रसंस्करण:** - Apache Spark 3.5 पाँच PySpark जॉब्स के साथ बैच ETL संभालता है: दैनिक KPI कंप्यूटेशन, फ़नल विश्लेषण, कार्ट परित्याग डिटेक्शन, प्रोडक्ट अफ़िनिटी स्कोरिंग, और MapReduce के लिए डेटा तैयारी - YARN Streaming चलाने वाला Apache Hadoop 3.4.1 दो MapReduce जॉब्स निष्पादित करता है (श्रेणी इवेंट गिनती और ब्रांड रेवेन्यू एग्रीगेशन) - Spark Structured Streaming Kafka से उपभोग करता है, सक्रिय सत्रों को हर 5 सेकंड में एग्रेगेट करता है और PostgreSQL में upsert करता है **ऑर्केस्ट्रेशन और API:** - Apache Airflow 2.9.3 पूरे बैच पाइपलाइन को DAG के रूप में ऑर्केस्ट्रेट करता है - FastAPI (asyncpg के साथ) डैशबोर्ड के लिए REST एंडपॉइंट्स expos करता है - Terraform टेम्पलेट्स AWS और GCP डिप्लॉयमेंत को कवर करते हैं **फ्रंटएंड:** - React 18.3 TypeScript, Tailwind CSS, और Recharts के साथ चार पेज रेंडर करता है: - Overview: KPI कार्ड्स और दैनिक प्रवृत्ति चार्ट्स - Funnel Analysis: व्यू-टू-कार्ट-टू-खरीद फ़्लो और श्रेणी द्वारा परित्याग - Product Analytics: शीर्ष उत्पाद, ब्रांड, और श्रेणियाँ - Live Monitor: Server-Sent Events के माध्यम से रियल-टाइम सत्र ट्रैकिंग ## डेटा पाइपलाइन 1. **बैच एनालिटिक्स**: MinIO से कच्चा CSV पढ़ा जाता है; दैनिक मेट्रिक्स (इवेंट्स, उपयोगकर्ता, रेवेन्यू, कन्वर्ज़न दर), शीर्ष उत्पाद/ब्रांड/श्रेणियों की गणना करके PostgreSQL में लिखा जाता है। एक अलग जॉब प्रति-श्रेणी फ़नल दरों की गणना करता है, कार्ट परित्याग को पहचानता है, और प्रोडक्ट अफ़िनिटी lift scores की गणना करता है (कार्यक्षमता के लिए 10% सैंपल्ड)। 2. **MapReduce**: तैयार Parquet फ़ाइलों को CSV में परिवर्तित कर HDFS पर अपलोड किया जाता है, जहाँ Hadoop Streaming जॉब्स प्रत्येक श्रेणी/इवेंट_टाइप में इवेंट्स की गिनती और प्रति ब्रांड रेवेन्यू का योग करते हैं। 3. **लाइव स्ट्रीमिंग**: एक Kafka प्रोड्यूसर CSV को कॉन्फ़िगर करने योग्य दर पर (डिफ़ॉल्ट 200 इवेंट्स/सेकंड) पुनः चलाता है। एक Spark स्ट्रीमिंग कंज़्यूमर 5-सेकंड माइक्रो-बैचेस में सक्रिय सत्रों को एग्रेगेट करता है और उन्हें `live_sessions` टेबल में संग्रहीत करता है। ## डिप्लॉयमेंट सभी 13 सेवाएँ Docker Compose के माध्यम से चलती हैं। पूर्व-अवश्यकताओं में कम से कम 16 GB RAM और 30 GB डस्क स्पेस के साथ Docker Desktop >= 4.30 शामिल है। Windows उपयोगकर्ताओं को कम से कम 12 GB मेमोरी के साथ WSL2 की आवश्यकता होती है। Makefile सेवाएँ शुरू करने, डेटा अपलोड करने, पाइपलाइन्स चलाने और परीक्षण निष्पादित करने के लिए टार्गेट प्रदान करता है। शुरुआत के बाद एक्सेस URLs: डैशबोर्ड (पोर्ट 3000), FastAPI docs (पोर्ट 8000), Airflow (पोर्ट 8080), MinIO Console (पोर्ट 9001), Spark Master UI (पोर्ट 8081), HDFS NameNode (पोर्ट 9870), और YARN (पोर्ट 8088)। ## रिपॉज़िटरी संरचना मुख्य निर्देशिकाओं में `spark/jobs/` पाँच PySpark बैच स्क्रिप्ट्स के लिए, `spark/streaming/` Kafka कंज़्यूमर के लिए, `kafka/producer/` रीप्ले प्रोड्यूसर के लिए, `api/` routed एंडपॉइंट्स के साथ FastAPI सेवा के लिए, `frontend/` React अनुप्रयोग के लिए, `hadoop/mapreduce/` Streaming जॉब स्क्रिप्ट्स के लिए, `airflow/dags/` ऑर्केस्ट्रेशन के लिए, `terraform/` IaC टेम्पलेट्स के लिए, और `docker/` Hadoop और Kafka सेवा परिभाषाओं के लिए शामिल हैं। लाइसेंस: MIT