पाठ 20 / 25

Stream Processing: Kafka Streams, Flink और ksqlDB

Events को filter, join और aggregate करने के लिए stream-processing तरीक़ा चुनें।

Stream पर गणना

Stream processors topics पढ़ते हैं, events बदलते हैं और नतीजे दूसरे topics में लिखते हैं। Kafka Streams Java library है जो आपके application के भीतर चलती है, changelog topics द्वारा समर्थित स्थानीय state के साथ stateful ऑपरेशन (aggregations, joins, windows) देती है। ksqlDB SQL-जैसी परत देता है। Apache Flink बड़े, जटिल streaming jobs के लिए अलग, शक्तिशाली engine है। मूल विचार: event time बनाम processing time, windows (tumbling, hopping, session), state और देर से आया डेटा। सबसे सरल उपयुक्त साधन से शुरू करें: सरल रूपांतरण के लिए सादा consumer, joins, windows और state चाहिए तो Streams या Flink।

Windowed गिनती, चलाकर

मैंने एक-मिनट की tumbling window का सादा-Python मॉडल चलाया जो प्रति window शुरुआत events गिनता है। असली engines state stores, देर से आए डेटा का सँभालना और दोष-सहिष्णुता जोड़ते हैं।

from collections import Counter
events = [5, 20, 59, 61, 75, 130]            # event times in seconds
window = 60
counts = Counter((t // window) * window for t in events)
print(dict(sorted(counts.items())))

Output:

{0: 3, 60: 2, 120: 1}

देर से और क्रम से बाहर आए events की योजना बनाएँ

Networks और retries के कारण events देर से आते हैं। तय करें कि window बंद करने से पहले आप कितना इंतज़ार करेंगे (grace period), और बाद में आए events का क्या करेंगे।

त्वरित जाँच: Tumbling window क्या करती है?

  • सारे events गिराती है
  • Events को निश्चित, गैर-अतिव्यापी समय intervals में समूहित करती है
  • हर window को 99% अतिव्यापी करती है
  • Topics वर्णानुक्रम से छाँटती है
Answer

Events को निश्चित, गैर-अतिव्यापी समय intervals में समूहित करती है — Tumbling windows समय को लगातार निश्चित खानों में बाँटती हैं।