প্রকল্প সম্পর্কে

রিয়ালটাইম-ই-কমার্স-অ্যানালিটিক্স (RTEA) একটি হাই-পারফরম্যান্স স্ট্রিমিং পাইপলাইন যা পাইথনে লেখা, যেটি ই-কমার্স ক্লিকস্ট্রিম ডেটা একটি REST API থেকে বন্ধু করে Confluent Cloud Kafka-তে পাঠায় এবং ClickHouse-এ রিয়েল-টাইম অ্যানালিটিক্যাল কোয়েরি করার ব্যবস্থা করে। ## আর্কিটেকচার পাইপলাইনটি চারটি স্তর নিয়ে গঠিত: **ডেটা সোর্স** — Railway.app-এ হোস্ট করা একটি সিমেটেটেড ক্লিকস্ট্রিম জেনারেটর API যা প্রতি সেকেন্ডে প্রায় ১০,০০০ বার্তা JSON ফর্ম্যাটে উৎপাদন করে। **মেসেজ ব্রোকার** — `clickstream-events` টপিকযুক্ত Confluent Cloud Kafka, ৬টি পার্টিশন সহ, SASL_SSL অথেন্টিকেশন, Snappy কম্প্রেশন এবং ৭ দিনের রেটেনশন। **প্রসেসিং লেয়ার** — অপ্টিমাইজড সেটিংস সহ `kafka-python` ব্যবহারকারী প্রোডিউসার সার্ভিস (১৬ KB ব্যাচ, gzip কম্প্রেশন, ৩২ MB বাফার, ১০K msg/sec রেট-লিমিটেড) এবং ব্যাচ প্রসেসিং সহ কনজিউমার সার্ভিস (প্রতি ব্যাচে ১০০০ ইভেন্ট, অটো-কমিট অফসেট, ফিল্ড ট্রান্সফরমেশন)। **অ্যানালিটিক্স স্টোরেজ** — MergeTree ইঞ্জিন টেবিল সহ ClickHouse Cloud, `(timestamp, user_id)` দ্বারা অর্ডার করা, LZ4 কম্প্রেশন এবং ৩০ দিনের TTL। ## ডেটা ফ্লো ১. প্রোডিউসার প্রতি সেকেন্ডে এক্সটার্নাল API পোল করে, মেসেজগুলো একটি কোয়েতে সংগ্রহ করে এবং ব্যাচে Kafka টপিকে প্রকাশ করে। ২. একাধিক কনজিউমার ইনস্ট্যান্স অ্যাসাইন্ড পার্টিশনে সাবস্ক্রাইব করে, JSON ডিজেসেরিয়ালাইজ করে, টাইমস্ট্যাম্প এবং ফিল্ড ট্রান্সফর্ম করে এবং ব্যাচ ClickHouse-এ লিখে। ৩. ব্যবহারকারীরা SQL-এর মাধ্যমে `ecommerce.clickstream_events` টেবিল কোয়ারি করে। ## কী কনফিগারেশন ### প্রোডিউসার - ব্যাচ সাইজ: ১৬ KB - কম্প্রেশন: gzip - বাফার মেমোরি: ৩২ MB - রিট্রাইজ: ৩ ### কনজিউমার - ব্যাচ সাইজ: ১০০০ ইভেন্ট - ব্যাচ টাইমআউট: ৫ সেকেন্ড - অফসেট অটো-কমিট সক্রিয় ### 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) ``` ## প্রতিবেদিত পারফরম্যান্স - স্থায়ী থ্রুপুট: সেকেন্ডে ৬৫০+ ইভেন্ট - মোট ইভেন্ট প্রক্রিয়া: ৫১০,৪১৪+ - সাব-সেকেন্ড এন্ড-টু-এন্ড লেটেন্সি ## কোয়াইক স্টার্ট ১. `uv sync` দিয়ে ডিপেন্ডেন্সি ইনস্টল করুন ২. `.env.example` থেকে `.env`-তে কপি করুন এবং Confluent Cloud ও ClickHouse ক্রেডেনশিয়াল পূরণ করুন ৩. প্রোডিউসার রান করুন: `uv run python services/producer/src/simple_producer.py` ৪. কনজিউমার রান করুন: `uv run python services/consumer/src/simple_consumer.py` ৫. ডেটা ভেরিফাই করুন: `uv run python check_clickhouse.py` ## ট্রাবুলশুটিং রেপোজিটরি `debug_kafka.py`, `test_clickhouse_insert.py` সহ ডায়াগনস্টিক স্ক্রিপ্ট ধারণ করে এবং API কানেকশন টাইমআউট, Kafka কনজিউমার ল্যাগ, মেসেজ ডিজেরিয়ালাইজেশন এরর, এবং ClickHouse ইনসার্ট/স্কিমা মিসম্যাচ সহ সাধারণ ব্যর্থতা মোড কভার করে।