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

रियल-टाइम-ई-कॉमर्स-एनालिटिक्स (RTEA) पायथन में लिखी गई एक हाई-परफॉर्मेंस स्ट्रीमिंग पाइपलाइन है जो REST API से ई-कॉमर्स क्लिकस्ट्रीम डेटा capture करती है और उसे क्लिकहाउस के माध्यम से रियल-टाइम एनालिटिकल क्वेरी के लिए Confluent Cloud Kafka में route करती है। ## आर्किटेक्चर पाइपलाइन चार लेयर से मिलकर बनी है: **डेटा स्रोत** — Railway.app पर होस्ट किया गया एक सिमुलेटेड क्लिकस्ट्रीम जनरेटर API जो JSON फॉर्मेट में प्रति सेकंड लगभग 10,000 मेसेजेज उत्पन्न करता है। **मेसेज ब्रोकर** — 6 पार्τιशन वाले `clickstream-events` टॉपिक के साथ Confluent Cloud Kafka, SASL_SSL ऑथेंटिकेशन, Snappy कंप्रेशन और 7-दिन की रिटेंशन। **प्रोसेसिंग लेयर** — अनुकूलित सेटिंग्स के साथ kafka-python का उपयोग करने वाला प्रोड्यूसर सर्विस (16 KB बैच, gzip कंप्रेशन, 32 MB बफर, 10K msg/sec तक रेट-लिमिटेड) और बैच प्रोसेसिंग के साथ कंज़्यूमर सर्विस (प्रात्येक बैच में 1000 इवेंट्स, ऑटो-कमिट ऑफसेट, फ़ील्ड ट्रांसफॉर्मेशन)। **एनालिटिक्स स्टोरेज** — MergeTree इंजन टेबल के साथ ClickHouse Cloud, `(timestamp, user_id)` द्वारा ऑर्डर, LZ4 कंप्रेशन और 30-दिन की TTL। ## डेटा फ्लो 1. प्रोड्यूसर प्रति सेकंड बाहरी API से poll करता है, संदेशों को कतार में एकत्र करता है और उन्हें बैच में Kafka टॉपिक पर प्रकाशित करता है। 2. कई कंज़्यूमर इंस्टेंस प्रत्येक असाइन किए गए पार्टिशन को सब्सक्राइब करते हैं, JSON deserialize करते हैं, समय और फ़ील्डों में परिवर्तन करते हैं, और ClickHouse में बैच लिखते हैं। 3. उपयोगकर्ता SQL के माध्यम से `ecommerce.clickstream_events` टेबल पर क्वेरी करते हैं। ## प्रमुख कॉन्फ़िगरेशन ### प्रोड्यूसर - बैच साइज़: 16 KB - कंप्रेशन: gzip - बफ़र मेमोरी: 32 MB - पुनः प्रयास: 3 ### कंज़्यूमर - बैच साइज़: 1000 इवेंट्स - बैच टाइमआउट: 5 सेकंड - ऑफसेट ऑटो-कमिट सक्षम ### ClickHouse टेबल ```sql CREATE TABLE clickstream_events ( user_id String, session_id String, timestamp UInt64, event_type String, product_id String, product_category String, price Float64, quantity Int32, source String, received_time DateTime DEFAULT now() ) ENGINE = MergeTree() ORDER BY (timestamp, user_id) ``` ## रिपोर्ट की गई प्रदर्शन - निरंतर थ्रूपुट: 650+ इवेंट्स/सेकंड - कुल प्रसंस्कृत इवेंट्स: 510,414+ - सब-सेकंड एंड-टू-एंड लेटेन्सी ## त्वरित शुरुआत 1. `uv sync` के साथ निर्भरताएं इन्स्टॉल करें 2. `.env.example` को `.env` में कॉपी करें और Confluent Cloud व ClickHouse क्रेडेंशियल्स भरें 3. प्रोड्यूसर चलाएं: `uv run python services/producer/src/simple_producer.py` 4. कंज़्यूमर चलाएं: `uv run python services/consumer/src/simple_consumer.py` 5. डेटा की पुष्टि करें: `uv run python check_clickhouse.py` ## समस्या समाधान रिपॉजिटरी में डायग्नोज्टिक स्क्रिप्ट (`debug_kafka.py`, `test_clickhouse_insert.py`) शामिल हैं और API कनेक्शन टाइमआउट, Kafka कंज़्यूमर लैग, मेसेज डिज़ीरियलाइज़ेशन एरर और ClickHouse इंsert/स्कीमा असंगतता सहित सामान्य विफलता स्थितियों को कवर करते हैं।