पाठ 20 / 25

Event-Time Windows और Watermarks

Event time से aggregate करें और देर से आए डेटा के लिए watermark से state सीमित करें।

आगमन का नहीं, घटना का समय

असली events देर से और क्रम से बाहर आते हैं, इसलिए आप आम तौर पर processing time की जगह event time (घटना कब हुई, डेटा का एक column) से aggregate करते हैं। groupBy(F.window("event_time", "5 minutes")) tumbling windows बनाता है। चूँकि देर से आए events पुरानी windows को हमेशा अपडेट कर सकते हैं, Spark को सारा state रखना पड़ता। Watermark (withWatermark("event_time", "10 minutes")) कहता है "मैं अधिकतम 10 मिनट देर से आया डेटा स्वीकार करता हूँ": Spark window state को watermark के उससे आगे निकलने तक रखता है फिर हटा देता है, जिससे memory सीमित रहती है, और उससे देर से आए events अनदेखे हो सकते हैं। देरी वह चुनें जितनी देर आपका डेटा असल में आता है।

Watermark के साथ windowed aggregation (उदाहरण)

events streaming DataFrame है जिसमें event_time और user_id columns हैं। यहाँ चलाया नहीं गया।

counts = (events
    .withWatermark("event_time", "10 minutes")
    .groupBy(F.window("event_time", "5 minutes"), "user_id")
    .count())

query = (counts.writeStream
    .outputMode("append")                 # emit a window once the watermark passes it
    .format("delta").option("path", "s3a://lake/gold/user_counts")
    .option("checkpointLocation", "s3a://lake/_chk/user_counts")
    .start())

त्वरित जाँच: Watermark Spark को क्या करने देता है?

  • Events को तेज़ पहुँचाना
  • सीमित देरी सहते हुए पुराना window state हटाना और memory सीमित रखना
  • Stream encrypt करना
  • Windows बंद करना
Answer

सीमित देरी सहते हुए पुराना window state हटाना और memory सीमित रखना — यह तय करता है कि डेटा कितना देर से आ सकता है, ताकि पुराना state साफ़ हो सके।