← Back to list
AI & Data
#스트림처리#윈도잉#워터마크#이벤트시간#ExactlyOnce
Last updated · 2026-10-07

Windowing and Time Semantics in Stream Processing (Event Time, Watermarks, Exactly-Once)

1. Overview

A. Definition

Stream processing is a data-processing paradigm that continuously transforms, aggregates, and analyzes unbounded data—data with no fixed start or end—the instant it arrives, producing results with low latency.

Windowing is the boundary-setting technique that slices an unbounded stream into finite processing units by grouping events according to criteria such as time, count, or session, while time semantics is the model that guarantees aggregation accuracy by distinguishing "the time an event actually occurred (event time)" from "the time the engine processed it (processing time)."

The essential difficulty of stream processing lies in the question, "For data that never ends, when, what, and how accurately should a computation be finalized?" Batch processing has self-evident boundaries because it computes only after all input has gathered, but a stream flows ceaselessly and, due to network delays, arrives out-of-order, so the engine must itself define when to finalize an aggregation. Windowing decides "what to treat as a single group (where in event time)," and watermarks decide "when to close that group and emit a result (when in processing time)." From a professional engineer's perspective, this topic is not mere API usage but a real-time data-pipeline governance problem of designing the trade-off among latency, accuracy, and completeness.

Fundamentally, stream processing wrestles with the contradiction that "waiting until data is fully gathered is too late, yet not waiting leaves it incomplete." Where batch accepts latency for the sake of completeness and pure real-time abandons completeness to eliminate latency, modern stream processing parameterizes the middle ground through three knobs—windows, watermarks, and triggers—so that the same code can be tuned to business needs anywhere between a "fast but approximate result" and a "slow but accurate result." Therein lies its defining characteristic.

B. Background and Necessity

First, the time horizon of business decisions has shifted from "the next day" to "this very moment." Fraud detection systems (FDS), real-time recommendation, Internet of Things (IoT) equipment anomaly detection, and reflecting rapidly changing inventory and market prices all lose their value under nightly batch. Because decisions must be made within seconds to minutes of delay, a structure that processes data while it flows—rather than waiting for it to accumulate—became necessary.

Second, in a distributed ingestion environment data inherently arrives late and out of order. A mobile device that was offline and then recovers sends events from minutes ago belatedly, or differences in per-partition processing speed reverse the order. If one aggregates by processing time, "an order that occurred at 10:00 but arrived at 10:05" falls into the wrong window and distorts the result. Correcting this required event-time processing that uses the occurrence time carried by the event itself, together with the watermark concept that permits delay but draws a limit on it.

Third, being real-time is no excuse to abandon accuracy. If events are dropped or double-counted when a node dies or restarts due to a failure, the consequences are fatal in domains such as billing, settlement, and regulatory reporting. Accordingly, an exactly-once processing guarantee—whereby each event is reflected in the result exactly once even when a failure occurs—became a core requirement for stream engines.

Fourth, operational and cost demands have grown as well. A pipeline that runs nonstop around the clock must keep its state from growing without bound, and when logic changes it must be able to replay past data to regenerate results (reprocessing). Moreover, without backpressure control that propagates load upstream to protect itself during a sudden surge (traffic spike), the whole system collapses from memory exhaustion. Thus stream processing has evolved into a technology that demands not only accuracy but sustained operability.

C. The Spectrum of Processing Guarantees

The reliability of stream processing is defined by "how many times each event is reflected in the result upon failure." The minimal guarantee, at-most-once, performs no retransmission, so it permits loss but has no duplication. At-least-once prevents loss through reprocessing on failure, but the same event may be reflected twice, overstating sums and counts. Exactly-once guarantees that state and output are reflected exactly once, consistently across a failure. This does not mean the network transfer itself physically happens only once; rather, checkpoints and transactions make the "effect" happen once, so it is also called effectively-once. This distinction is realized through the checkpoint and transaction mechanisms covered in the deep-dive section below.

The choice of guarantee level is not free. The higher the guarantee, the greater the overhead of state snapshots, transaction coordination, and deferred commits, and the higher the latency. Therefore, rather than designing every pipeline uniformly as exactly-once, the standard practice is selective application: place exactly-once only on paths where accuracy is tied directly to financial or legal liability—billing, settlement, regulatory reporting—and place paths that tolerate a small amount of duplication or loss, such as dashboard metrics or recommendation signals, at at-least-once to optimize cost.

2. Overall Architecture and Time Semantics

A. The Stream-Processing Pipeline Architecture

A real-time pipeline is generally built around an immutable-log-based message broker, together with stateful processing operators, a checkpoint store, and idempotent or transactional sinks.

flowchart LR
  SRC["Event source(apps, IoT, logs)"] --> BROKER["Message broker(immutable log, partitions)"]
  BROKER --> OP["Stream operator(windows, aggregation)"]
  OP --> ST["State store(Keyed State)"]
  OP --> CP["Checkpoint(snapshot)"]
  CP --> DFS["Durable store(distributed file system)"]
  OP --> SINK["Sink(idempotent, transactional)"]
  SINK --> SERVE["Serving(DB, dashboard, alerts)"]

A message broker (e.g., Apache Kafka) provides an immutable log whose order is guaranteed per partition, becoming the foundation for reprocessing by rewinding the offset to replay on failure. Setting the retention period to, say, 7 days lets you rewind the consumption position to any point within that window and reapply the logic, while the partition count determines the unit of parallelism and thus the basis for throughput scaling. The stream operator maintains state per key and performs windowed aggregation, periodically taking state snapshots as checkpoints and storing them in a distributed file system. On failure it restores state from the most recent checkpoint and resumes consumption from that point's offset, so at-least-once reprocessing is the default. Combining this with an idempotent or transactional sink removes duplicate reflection as well, completing exactly-once. In this structure the checkpoint interval governs how finely spaced the recovery points are, and the performance of the state backend governs aggregation throughput, so tuning the two together is the starting point of operational design.

B. Event Time, Processing Time, and Watermarks

The starting point of time semantics is distinguishing three times. Event time is the time the event actually occurred (e.g., the payment-approval time), embedded in the data as a timestamp. Ingestion time is the time it entered the broker or engine, and processing time is the time the operator actually processed it. Reproducibility and accuracy of results come from event-time-based processing, because when the same input is replayed (reprocessing) the processing time differs each time, whereas event-time-based aggregation always yields the same result.

The problem is that closing a window by event time requires the judgment that "no more events belonging to that window will arrive." The device that estimates this is the watermark. A watermark W(t) is the engine's progress indicator asserting that "data before event time t has (for the most part) all arrived," and the moment this watermark passes the window's end time, the window's aggregation is finalized and fired (triggered). In practice, a bounded out-of-orderness watermark—the maximum observed event time minus the allowed delay—is commonly used. For example, with an allowed delay of 5 seconds, W = max_event_time - 5s.

sequenceDiagram
  participant E as Event stream
  participant W as Watermark generator
  participant WIN as Window operator
  participant OUT as Result sink
  E->>W: Event arrives(includes event-time timestamp)
  W->>WIN: Watermark propagated("arrival before t complete")
  WIN->>WIN: Compare window boundary with watermark
  alt watermark > window end
    WIN->>OUT: Finalize and fire window aggregation
  else late data arrives
    WIN->>WIN: Update if within allowed lateness
    WIN->>OUT: Separate out via side output
  end

For reference, ingestion time is a compromise between event time and processing time; when the source has no trustworthy timestamp, it is a practical alternative that uses the broker-entry time as the basis to gain some order stability. It does not, however, guarantee full reproducibility, since results may differ on reprocessing.

A watermark is essentially the knob that reconciles the trade-off between completeness and latency. Setting a large allowed delay makes the result accurate by including even late-arriving data, but window finalization is delayed, increasing latency. Setting it small makes things fast but misses late data, leaving the result incomplete. Many engines therefore continue to update the result with late data during an allowed lateness period even after the window is finalized, and separate data later than that via a side output into a distinct correction path.

C. Types of Windows

Windows are divided according to the criterion by which the unbounded stream is grouped. The choice of type depends on the nature of the question (whether a periodic aggregation or an analysis of activity intervals); the table below is only a supplementary summary, and the reason for choosing each type is explained in prose.

Window type Definition Characteristics Representative use
Tumbling Fixed size, non-overlapping Boundaries do not overlap, so each event falls in exactly 1 window 1-minute sales aggregation, hourly reports
Sliding Fixed size, shifting at a fixed interval Windows overlap, so an event belongs to several windows 5-minute moving average (every minute)
Session Based on the gap between activities Variable size, automatically separates activity intervals User sessions, equipment operating intervals
Global No boundary, custom trigger The user defines the trigger Count-based aggregation

A tumbling window cuts into fixed lengths such as 1 minute or 1 hour without overlap, and each event belongs to exactly one window. It suits clear periodic reports such as "transaction count every minute." Its boundaries are simple, so state-management cost is lowest, and when a window closes its state can be discarded immediately, making memory predictable. On the other hand, because the boundaries are fixed, one must account in design for the fact that if a threshold phenomenon is split across two windows, each falls below the threshold and detection can be missed.

A sliding window separates size (e.g., 5 minutes) from the shift interval (e.g., 1 minute) so that windows overlap, and is used for continuous trend monitoring such as "a 5-minute moving average refreshed every minute." Because one event is included redundantly in (size ÷ interval) windows, the state size and computation grow correspondingly. In the example above one event belongs to 5 windows, so there is roughly 5× the state and computation burden of simple tumbling. The trade-off between trend sensitivity (a shorter interval) and resource cost is the key design point.

A session window divides intervals not by fixed boundaries but by the gap between activities. For example, with a gap threshold of 30 minutes, if a user's continuous clicks break for 30 minutes or more, the session ends and a new one begins. It is natural for user-behavior analysis or extracting equipment operating sessions. Because the window length is determined dynamically by the data, when a late event fills the gap between two sessions a complexity arises of having to merge two already-fired sessions into one. For this reason the session window is the type whose state and trigger design is most demanding.

A global window fires via a user-defined trigger such as "every 100 records," without a time boundary. It responds flexibly to count- or threshold-condition-based processing that is hard to express as time-based aggregation, but because the user is wholly responsible for the trigger and state cleanup, overuse causes state leaks and memory growth.

D. Triggers and Accumulation Modes

Even for the same window, "when and how many times to emit a result" is designed separately as the trigger, and "how to combine with the previous result when re-firing" as the accumulation mode. In environments where late data matters, it is useful to emit a fast first (speculative) result when the watermark arrives, then re-fire the window to correct the result when late data later arrives. Here the accumulating mode emits the full recomputed value each time, while the accumulating-and-retracting mode emits a correction record that cancels the previous value together with the new value, so that downstream does not double-aggregate. This "early result + after-the-fact correction" pattern is the core contribution of the Dataflow model and is the key to obtaining low latency without sacrificing completeness.

3. Comparison and Applied Cases

A. Comparison of Representative Engines and the Reasons for Differences

Dimension Apache Flink Kafka Streams Spark Structured Streaming
Processing model True per-record streaming Per-record (library) Micro-batch (default), continuous (experimental)
State and checkpoint Asynchronous Barrier Snapshotting (ABS) Changelog topic + RocksDB Checkpoint + WAL
exactly-once Guaranteed via state + transactional sink exactly_once_v2 (within Kafka) Based on idempotent or transactional sinks
Latency profile Millisecond-level low latency Millisecond-level (embedded in app) As much as the batch interval (hundreds of ms to seconds)

The differences among the three engines stem from the fundamental view of "what streaming is." Flink treats a stream as a first-class citizen and processes it continuously per record, so it is strong at low latency and sophisticated event-time and watermark control. Spark Structured Streaming traditionally treats a stream as a succession of small batches—a micro-batch model—so it is advantageous for throughput and integration with the batch ecosystem, but at the cost of latency on the order of the batch interval (continuous-processing mode lowers latency but has constraints on the guarantee level). Kafka Streams is a library embedded in the application without a separate cluster, so operations are simple; it replicates state to a changelog topic for restoration and provides exactly-once for I/O internal to Kafka. In short, Flink is the reasonable choice when low latency and sophisticated time control are central, Kafka Streams for lightweight embedded processing centered on Kafka, and Spark when integration with the batch and ML ecosystem matters.

B. Concrete Applied Cases

Consider real-time fraud detection (FDS). Card-approval events are grouped by the card-number key, approval counts and amounts are aggregated in a 1-minute tumbling window, and a blocking alert is fired if a threshold is exceeded. Because, by the nature of mobile payment, some events arrive several seconds late, the watermark delay is set to 5 seconds to include late arrivals, while events later than that are sent to a side output and reflected in after-the-fact settlement. The judgment of "the same 1-minute interval" is consistent only when based on the approval time (event time), not the processing time.

As an industrial IoT case, separating operating sessions with a session window (gap threshold of 10 minutes) from the sensor streams of thousands of machines allows early detection of anomalies from vibration and temperature trends within a session. For equipment whose operation and stoppage are irregular, a session window reflects the real operating intervals more naturally than a fixed window. A global ride-hailing service is also known to operate Flink-based stream processing at large scale for real-time supply-demand matching and ETA prediction, in the same vein.

In e-commerce recommendation, a sliding window is used. Aggregating each user's clicks over the last 10 minutes with a 1-minute sliding interval to reflect the trend of products of interest in real time makes the recommendation respond sensitively to the latest behavior as the session lengthens. Here, because an event is redundantly included in several windows and state grows, one designs the state TTL together with the checkpoint interval (e.g., 10 seconds) to manage memory and recovery time. For example, with 1 million active users × 1 KB average state, roughly 1 GB of per-key state is maintained at all times, so the decision to move the state backend from memory to RocksDB and to expire idle users' state by TTL directly affects throughput and cost.

4. Deep Dive: Implementation Mechanisms of Exactly-Once and Recent Trends

Exactly-once is the problem of guaranteeing "the effect only once" rather than "transmission only once." Flink implements this via Asynchronous Barrier Snapshotting. When the source periodically inserts a special barrier into the data flow, each operator takes a snapshot of its own state as the barrier passes and stores it. When the snapshots of all operators are gathered, a globally consistent checkpoint is completed, and on failure state and source offset are rewound together to this point. The key is that performance degradation is small because snapshots are taken asynchronously without halting processing. This is an adaptation of distributed-snapshot theory (the Chandy-Lamport algorithm) to stream data flows.

Since state restoration alone cannot prevent duplication of output that has already left for the sink, the output side combines idempotent writes or transactional writes (two-phase commit). For example, Flink's TwoPhaseCommitSinkFunction commits the external transaction in step with checkpoint completion, atomically binding the checkpoint and output visibility. Kafka provides exactly-once for I/O internal to Kafka via the read-process-write pattern, which binds a transaction (a producer transaction) and the consumer offset commit into a single transaction, and Kafka Streams simplified this with the processing.guarantee=exactly_once_v2 setting (improvements in the KIP-447 line are reported to have improved efficiency when processing many partitions).

Meanwhile, Spark Structured Streaming guarantees failure recovery based on a WAL (Write-Ahead Log) that records offsets and state in checkpoints, securing end-to-end accuracy when combined with an idempotent or transactional sink. Kafka Streams replicates state to a changelog topic and restores it on another instance, so a stateful application can be scaled out and relocated like a stateless service. Although the mechanism's implementation differs per engine, the key is to understand that the common principle—"a consistent state snapshot + idempotent or transactional guarantee of the output"—is the same.

Several currents stand out among recent trends. First, the four axes established in Google's Dataflow model—"What, Where (windows), When (watermarks and triggers), and How (correcting late data)"—have spread through Apache Beam as an engine-neutral abstraction, strengthening the direction of reusing a single pipeline definition across multiple execution engines. Second, streaming SQL (Flink SQL and the like), which treats a stream like a table via SQL, has matured, raising accessibility for non-developer roles. Third, there is active effort toward a state-compute separation architecture that separates the state store into cloud object storage to scale compute and state independently, and toward unifying stream and batch over a single table format (e.g., a lakehouse such as Apache Iceberg). That said, these detailed figures and versions change quickly, so it is advisable to verify them against the latest official documentation at design time.

5. Considerations and Implications

First, design by translating the latency-completeness trade-off into business requirements. The watermark's allowed delay and the window size are, before being technical parameters, business judgments of "up to how late is data worth reflecting in the result." Where completeness matters, as in settlement, set a generous allowed delay and provide an after-the-fact correction path; where immediacy matters, as in alerting, set it short and correct late arrivals separately. From a professional engineer's perspective, the advisable approach is to first define the SLA and accuracy requirements and then back-calculate the parameters.

Second, make the scope and cost of exactly-once explicit. Exactly-once is usually limited to "engine state and a specific sink" and does not automatically extend to the entire external system. External integrations for which idempotent design is impossible (e.g., outbound email or payment) are hard to bind in a transaction, so they must be complemented with at-least-once plus an idempotency-key design. Also, shortening the checkpoint interval reduces recovery time but increases overhead, so one must find the balance point between the recovery time objective (RTO) and throughput.

Third, design state management and the reprocessing strategy together. To keep per-key state from growing without bound, put in place TTLs, a state backend (e.g., RocksDB), and state-size monitoring, and prepare in advance the broker retention period, an offset-rewind strategy, and an idempotent regeneration structure for result views so that past data can be safely reprocessed when logic changes. Designing in connection with the reprocessing philosophy of the Kappa architecture ([[lambda-kappa-architecture]]) can lower operational complexity.

Fourth, build in operational observability and backpressure handling. A sudden surge in watermark delay, the checkpoint failure rate, consumer lag, and intervals where backpressure occurs are health indicators of a real-time pipeline. By combining the backpressure control of reactive streams ([[reactive-streams-backpressure]]) with the partition and offset observation of message queues ([[message-queue-kafka]]), one must automatically throttle throughput during a surge and keep state from straying beyond the recoverable range.

Fifth, outlook and related technologies. Streaming SQL, Beam-based engine-neutral abstractions, and lakehouse-based stream-batch unification are converging in the direction of "blurring the boundary between real-time and batch." Viewed together with data-consistency models such as eventual consistency ([[eventual-consistency]]) and CQRS with event sourcing ([[cqrs-event-sourcing]]), stream processing should be understood not as a single feature but as a design axis threading through the entire real-time data platform.

References


In one line: Stream processing groups unbounded data into finite units via windowing, determines when to finalize aggregation using event time and watermarks, and guarantees exactly-once processing via checkpoints and transactional sinks—a real-time data-processing technology for designing the trade-off among latency, accuracy, and completeness.