عن المشروع

RealTime-ECommerce-Analytics (RTEA) هو خط أنابيب تدفق عالي الأداء مكتوب بلغة Python يلتقط بيانات تدفق النقرات للتجارة الإلكترونية من واجهة برمجة تطبيقات REST ويوجهها عبر Confluent Cloud Kafka إلى ClickHouse للاستعلامات التحليلية في الوقت الفعلي. ## البنية يتكون خط الأنابيب من أربع طبقات: **مصدر البيانات** — واجهة برمجة تطبيقات مولد تدفق نقرات محاكاة مستضافة على Railway.app تنتج حوالي 10,000 رسالة في الثانية بتنسيق JSON. **وسيط الرسائل** — Confluent Cloud Kafka مع موضوع `clickstream-events` عبر 6 أقسام، مصادقة SASL_SSL، ضغط Snappy، واحتفاظ لمدة 7 أيام. **طبقة المعالجة** — خدمة منتج باستخدام `kafka-python` بإعدادات محسنة (دفعات 16 كيلوبايت، ضغط gzip، مخزن مؤقت 32 ميجابايت، محدود بمعدل 10 آلاف رسالة/ثانية) وخدمة مستهلك مع معالجة دفعية (1000 حدث لكل دفعة، تأكيد إزاحة تلقائي، تحويل الحقول). **تخزين التحليلات** — ClickHouse Cloud مع جدول محرك MergeTree مرتب حسب `(timestamp, user_id)`، ضغط LZ4، و TTL لمدة 30 يومًا. ## تدفق البيانات 1. يستعلم المنتج واجهة برمجة التطبيقات الخارجية كل ثانية، يجمع الرسائل في طابور، وينشرها دفعيًا إلى موضوع Kafka. 2. تشترك عدة مثيلات مستهلك في الأقسام المخصصة، وتفك ترميز JSON، وتحول الطوابع الزمنية والحقول، وتكتب دفعات إلى ClickHouse. 3. يستعلم المستخدمون جدول `ecommerce.clickstream_events` الناتج عبر SQL. ## التكوين الرئيسي ### المنتج - حجم الدفعة: 16 كيلوبايت - الضغط: gzip - ذاكرة المخزن المؤقت: 32 ميجابايت - محاولات إعادة: 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`) ويغطي أنماط الفشل الشائعة بما في ذلك مهلة اتصال واجهة برمجة التطبيقات، تأخر مستهلك Kafka، أخطاء فك ترميز الرسائل، وعدم تطابق مخطط/إدراج ClickHouse.