ব্যাচ বনাম স্ট্রিম প্রসেসিং
এই পাঠে যা শিখবেন
- ব্যাচ প্রসেসিং ও স্ট্রিম প্রসেসিং-এর মৌলিক পার্থক্য এবং প্রতিটির trade-off
- কখন ব্যাচ, কখন স্ট্রিম বেছে নেবেন — বাস্তব উদাহরণসহ
- Windowing-এর প্রাথমিক ধারণা — Tumbling বনাম Sliding window
- Python দিয়ে একই ডেটার উপর ব্যাচ ও স্ট্রিম অ্যাগ্রিগেট হিসাব করে ফলাফল মিলিয়ে দেখা
১ · দুটি ভিন্ন প্রসেসিং মডেল
L22-L23-এ আমরা দেখেছি ইভেন্ট/বার্তা কীভাবে তৈরি ও ডেলিভার হয়। কিন্তু একবার ইভেন্ট এসে পৌঁছালে, সেটি কখন প্রসেস করা হবে তার দুটি মৌলিকভাবে ভিন্ন পদ্ধতি আছে।
একটি বড়, নির্দিষ্ট ডেটাসেট (যেমন গতকালের সব লেনদেন) জমা করে একটি শিডিউলে (রাতে, প্রতি ঘণ্টায়) একসাথে প্রসেস করা — Hadoop/Spark batch jobs।
ডেটা আসার সাথে সাথে, একটির পর একটি, ক্রমাগত প্রসেস করা — Kafka Streams, Apache Flink। নিয়ার-রিয়েল-টাইম ফলাফল।
ব্যাচ প্রসেসিং বেশি থ্রুপুট-দক্ষ — একসাথে অনেক ডেটা প্রসেস করলে প্রতি-রেকর্ড ওভারহেড কম হয়, এবং সম্পূর্ণ ডেটাসেট হাতে থাকায় জটিল অ্যাগ্রিগেশন/জয়েন সহজ। কিন্তু ফলাফল পেতে ঘণ্টা বা দিন লেগে যেতে পারে। স্ট্রিম প্রসেসিং ফলাফল সেকেন্ডের মধ্যে দেয়, কিন্তু প্রতিটি ইভেন্ট আলাদাভাবে (বা ছোট মাইক্রো-ব্যাচে) সামলাতে হয়, আউট-অফ-অর্ডার ইভেন্ট সামলাতে হয়, এবং ইনফ্রাস্ট্রাকচার বেশি জটিল ও ব্যয়বহুল হয়।
২ · কখন কোনটি বেছে নেবেন
রাতের বিক্রয় রিপোর্ট, মাসিক বিলিং সামারি, মেশিন লার্নিং মডেল রি-ট্রেনিং — এসব ক্ষেত্রে কয়েক ঘণ্টার দেরি সম্পূর্ণ সহনীয়, তাই ব্যাচ প্রসেসিং সঠিক ও সস্তা পছন্দ। কিন্তু ক্রেডিট কার্ড ফ্রড ডিটেকশন, লাইভ স্টক টিকার, রাইড-শেয়ারিং অ্যাপে রিয়েল-টাইম চাহিদা-দাম (surge pricing), বা লাইভ ড্যাশবোর্ডে — এসব ক্ষেত্রে সেকেন্ডের দেরিও অগ্রহণযোগ্য, তাই স্ট্রিম প্রসেসিং দরকার।
৩ · Windowing — সংক্ষেপে
স্ট্রিম প্রসেসিংয়ে প্রায়ই "গত ৫ মিনিটে গড় কত" জাতীয় প্রশ্নের উত্তর দিতে হয় — এর জন্য ডেটাকে উইন্ডোতে ভাগ করা হয়। Tumbling window — নির্দিষ্ট, অ-ওভারল্যাপিং সময়-বাক্স (যেমন প্রতি ৫ মিনিট আলাদা আলাদা উইন্ডো)। Sliding window — ওভারল্যাপিং, ক্রমাগত এগিয়ে চলা উইন্ডো (যেমন "শেষ ৫ মিনিট", প্রতি সেকেন্ডে রিক্যালকুলেটেড)। Tumbling কম গণনা লাগে, Sliding বেশি স্মুথ/আপ-টু-ডেট ফলাফল দেয় কিন্তু বেশি রিসোর্স নেয়।
৪ · Python-এ ব্যাচ ও স্ট্রিম অ্যাগ্রিগেট মিলিয়ে দেখা
নিচে একই ১০টি টাইমস্ট্যাম্প-করা লেনদেন ইভেন্টের উপর দুইভাবে যোগফল বের করছি — একবার ব্যাচ পদ্ধতিতে (সব ইভেন্ট জমা হওয়ার পর একবারে), একবার স্ট্রিম পদ্ধতিতে (প্রতিটি ইভেন্ট আসার সাথে সাথে running total আপডেট)।
events = [
(1, 120), (2, 85), (3, 200), (4, 50), (5, 175),
(6, 90), (7, 60), (8, 220), (9, 40), (10, 130),
] # (timestamp, amount) - টাকা-লেনদেনের একটি স্ট্রিম
# ব্যাচ প্রসেসিং — সব ইভেন্ট জমা হওয়ার পর একবারে হিসাব করা
batch_total = sum(amount for _, amount in events)
# স্ট্রিম প্রসেসিং — প্রতিটি ইভেন্ট আসার সাথে সাথে running total আপডেট
running_total = 0
print("স্ট্রিম প্রসেসিং — প্রতিটি ইভেন্টের পর running total:")
for ts, amount in events:
running_total += amount
print(f" t={ts}: +{amount:>3} -> running total = {running_total}")
print(f"\nব্যাচ টোটাল (সব শেষে একবারে): {batch_total}")
print(f"স্ট্রিম চূড়ান্ত টোটাল (ধাপে ধাপে): {running_total}")
print(f"দুটো ফলাফল মেলে কি? {batch_total == running_total}")
ব্যাচ ও স্ট্রিম একে অপরের বিকল্প নয় — অনেক বাস্তব সিস্টেম দুটোই ব্যবহার করে (Lambda আর্কিটেকচার): স্ট্রিম দ্রুত-কিন্তু-আনুমানিক ফলাফল দেয়, আর ব্যাচ পরে নিখুঁত/সংশোধিত ফলাফল দিয়ে ওভাররাইট করে। প্রশ্নটা "কোনটি সবসময় ভালো" নয়, বরং "এই নির্দিষ্ট ব্যবহারের জন্য লেটেন্সি বনাম জটিলতার ট্রেড-অফে কোনটা মানানসই।"
ভাবনার প্রশ্ন
প্রতিটি প্রশ্ন নিজে কিছুক্ষণ ভাবুন — তারপর "→ উত্তর" চাপুন।
প্র ০১ একটি ব্যাংক কেন লেনদেন-ভিত্তিক ফ্রড ডিটেকশনের জন্য ব্যাচ প্রসেসিং ব্যবহার করবে না, যদিও ব্যাচ থ্রুপুট-দক্ষ?
ফ্রড ডিটেকশনের মূল্য তার গতিতে — একটি জালিয়াতিমূলক লেনদেন যত দ্রুত ধরা পড়ে, তত কম ক্ষতি হয়। ব্যাচ প্রসেসিংয়ে যদি রাতের জব চলে, তাহলে একটি জালিয়াতিমূলক লেনদেন ঘটার ২৩ ঘণ্টা পর্যন্তও ধরা নাও পড়তে পারে — ততক্ষণে টাকা তুলে নেওয়া হয়ে যেতে পারে। তাই এখানে থ্রুপুট-দক্ষতার চেয়ে লেটেন্সি (কত দ্রুত ধরা যায়) অনেক বেশি গুরুত্বপূর্ণ, এবং স্ট্রিম প্রসেসিং প্রতিটি লেনদেন ঘটার মুহূর্তেই মূল্যায়ন করে তাৎক্ষণিক ব্লক করার সুযোগ দেয়।
প্র ০২ স্ট্রিম প্রসেসিংয়ে "আউট-অফ-অর্ডার ইভেন্ট" সমস্যাটি আসলে কী, এবং ব্যাচ প্রসেসিংয়ে এটি কেন সাধারণত সমস্যা নয়?
নেটওয়ার্ক লেটেন্সি বা একাধিক সোর্স থাকার কারণে ইভেন্টগুলো কখনো কখনো তাদের প্রকৃত ঘটার ক্রম অনুযায়ী না এসে এলোমেলো ক্রমে পৌঁছাতে পারে (যেমন ৩ নম্বর ইভেন্টের আগেই ৫ নম্বর ইভেন্ট চলে আসা)। স্ট্রিম প্রসেসিংয়ে যেহেতু প্রতিটি ইভেন্ট আসার সাথে সাথেই প্রসেস করতে হয়, তাই একটি "দেরিতে আসা" পুরনো ইভেন্ট আগের হিসাবকে ভুল করে দিতে পারে। ব্যাচ প্রসেসিংয়ে যেহেতু সব ইভেন্ট প্রসেসিং শুরুর আগেই হাতে জমা থাকে, প্রসেসিং শুরুর আগে টাইমস্ট্যাম্প অনুযায়ী সর্ট করে নেওয়া যায় — তাই ক্রম নিয়ে সমস্যাই হয় না।
প্র ০৩ একটি নিউজফিড অ্যাপ (L46-এ বিস্তারিত) কি একটি ইউজারের আনরিড নোটিফিকেশন কাউন্ট আপডেট করার জন্য ব্যাচ নাকি স্ট্রিম ব্যবহার করা উচিত?
স্ট্রিম প্রসেসিং — কারণ একজন ইউজার একটি নতুন লাইক/কমেন্ট পাওয়ার সাথে সাথে তার আনরিড কাউন্ট আপডেট হওয়া উচিত, কয়েক ঘণ্টা পর নয়। এখানে ইউজার-অভিজ্ঞতার সাথে সরাসরি সম্পর্কিত রিয়েল-টাইম আপডেট দরকার। বিপরীতে, "এই সপ্তাহে কোন পোস্টগুলো সবচেয়ে বেশি এনগেজমেন্ট পেয়েছে" জাতীয় সাপ্তাহিক ট্রেন্ড রিপোর্ট ব্যাচ প্রসেসিং দিয়েই যথেষ্ট, কারণ এখানে সেকেন্ডের নিখুঁততা কোনো বাস্তব মূল্য যোগ করে না।
অনুশীলন
-
চিন্তা করুন: একটি ভিডিও স্ট্রিমিং প্ল্যাটফর্মের (L48) জন্য ৩টি ব্যবহারের ক্ষেত্র চিহ্নিত করুন যেখানে ব্যাচ প্রসেসিং উপযুক্ত, এবং ৩টি যেখানে স্ট্রিম প্রসেসিং উপযুক্ত।
ব্যাচ উপযুক্ত: মাসিক ভিউ-কাউন্ট রিপোর্ট জেনারেট করা, প্রতিদিনের রেকমেন্ডেশন মডেল রি-ট্রেনিং, পুরনো লগ থেকে সাপ্তাহিক অ্যানালিটিক্স ড্যাশবোর্ড তৈরি করা। স্ট্রিম উপযুক্ত: লাইভ স্ট্রিমের দর্শক সংখ্যা রিয়েল-টাইমে দেখানো, অ্যাডাপটিভ বিটরেট স্ট্রিমিং-এ ব্যান্ডউইথ পরিবর্তন সনাক্ত করে কোয়ালিটি সমন্বয় করা (L48), এবং লাইভ চ্যাট/কমেন্ট সেকশনে বার্তা রিয়েল-টাইমে দেখানো।
-
কোড পরিবর্তন করুন: উপরের কোড সেলে একটি "tumbling window" যোগ করুন — প্রতি ৩টি ইভেন্টকে একটি উইন্ডো ধরে প্রতিটি উইন্ডোর যোগফল আলাদাভাবে প্রিন্ট করুন (সবশেষে অসম্পূর্ণ উইন্ডো থাকলেও প্রিন্ট করুন), তারপর সব উইন্ডোর যোগফল মিলিয়ে দেখুন তা ব্যাচ টোটালের সমান কিনা।
সমাধানের কাঠামো:
window_size = 3। ১০টি ইভেন্টে ৩-সাইজের উইন্ডো হলে ৪টি উইন্ডো তৈরি হবে (৩+৩+৩+১টি ইভেন্ট) — যোগফল যথাক্রমে ৪০৫, ৩১৫, ৩২০, ১৩০ (মোট ১,১৭০), যা ব্যাচ টোটালের (১,১৭০) সমান — উইন্ডোয়িং শুধু একই ডেটাকে ছোট ছোট অংশে ভাগ করে, মোট পরিমাণ বদলায় না।
for i in range(0, len(events), window_size):
window = events[i:i+window_size]
window_total = sum(a for _, a in window)
print(f"window {i//window_size+1}: {window_total}")
আরও পড়ুন · ABCL TECH-এ আপনার পরবর্তী পদক্ষেপ
- কোর্সের সম্পূর্ণ সিলেবাস দেখুন ৫১টি পাঠ পরবর্তী পাঠ — মনোলিথ বনাম মাইক্রোসার্ভিস, যেখানে M7 আর্কিটেকচার প্যাটার্নস মডিউল শুরু হচ্ছে।
- ইভেন্ট-ড্রিভেন আর্কিটেকচার L23 স্ট্রিম প্রসেসিং কীভাবে ইভেন্ট-ড্রিভেন সিস্টেমের উপর ভিত্তি করে কাজ করে তা আবার দেখে নিন।
- সব Courses দেখুন ABCL TECH C, C++, Python, Java, JavaScript, DSA, DBMS, Discrete Mathematics ও System Design — সব এক জায়গায়।