Published on

Sliding Window System Design: Real-Time Data Processing in Go

Authors
  • Mehdi Akiki avatar
    Name
    Mehdi Akiki
    Twitter

The sliding window pattern shows up more than most people realise. Rate limiters use it ("no more than 100 requests in the last 60 seconds"). Monitoring systems use it ("alert if error rate exceeds 5% over the last 5 minutes"). Analytics pipelines use it constantly. Once you see it, you start recognising it everywhere.

This guide explains the pattern from scratch, then builds a real-time stock price monitor in Go — a program that tracks incoming prices and continuously computes the average over the last 5 seconds.

Contents

  1. What is a sliding window?
  2. The problem we're solving
  3. Design: the three pieces
  4. Step 1: The data structure
  5. Step 2: Generating prices
  6. Step 3: The window processor
  7. Step 4: Wiring it together
  8. Running it
  9. Taking it further

1. What is a sliding window?

Start with the naive approach: you want to know the average stock price over the last 5 seconds. You could keep a list of every price ever recorded, then filter down to the last 5 seconds and average whatever's left. That works, but the list grows forever and filtering it gets slower over time.

A sliding window solves this by maintaining only the data that's still relevant. As time moves forward, the window moves with it — old data falls out the back, new data comes in the front. At any given moment you have exactly the data you need and nothing you don't.

Here's what that looks like visually with a 3-second window:

Prices arriving every second:

Time:  t=1    t=2    t=3    t=4    t=5    t=6
Price: $72    $75    $68    $80    $71    $74

Window at t=3:  [$72, $75, $68]          avg = $71.67
Window at t=4:       [$75, $68, $80]     avg = $74.33
Window at t=5:            [$68, $80, $71] avg = $73.00
Window at t=6:                 [$80, $71, $74] avg = $75.00

As the clock ticks from t=3 to t=4, $72 drops off the back because it's now more than 3 seconds old. $80 joins at the front. The window is always exactly 3 seconds wide.

Two properties define any sliding window:

  • Window size — how much history to keep, measured in time or number of events
  • Slide interval — how often the window advances (here: every second)

2. The problem we're solving

We want to build a program that:

  1. Receives a continuous stream of stock prices — one price per second, timestamped when it arrives
  2. Maintains a 5-second sliding window over those prices
  3. Prints the average price in the current window every second

The stream is simulated with a goroutine that generates random prices. In a real system this would be a WebSocket feed, a Kafka topic, or a database polling loop — the window logic stays the same.


3. Design: the three pieces

Before writing any code, it helps to understand what we're building and why each piece exists.

The generator produces prices and puts them somewhere the processor can read. In Go, "somewhere" is a channel — a typed queue between goroutines. The generator runs in its own goroutine so it doesn't block anything else.

The channel connects the generator to the processor. It's buffered (size 10) so the generator can keep running even if the processor is briefly busy. No shared memory, no mutexes — just a channel.

The processor does two things on a loop:

  • When a new price arrives on the channel, add it to the window
  • Every second (on a ticker), evict prices older than 5 seconds and compute the average of what's left

The processor runs on the main goroutine and blocks forever — which is fine because the generator runs in the background.

Here's the data flow:

[generator goroutine]
    generates price every 1s
         |
         | channel (buffered, size 10)
         ↓
[processor goroutine]
    select {
        price arrives → add to window slice
        ticker fires  → evict old prices, compute average, print
    }

4. Step 1: The data structure

We need a type to hold a single price reading. The two fields that matter are the price itself and the exact time it was recorded — the timestamp is what the eviction logic uses to decide whether an entry is still inside the window.

type StockPrice struct {
    Value      float64
    RecordedAt time.Time
}

That's it. No ID, no symbol, no extra metadata — just what the algorithm needs.

The window itself will be a plain Go slice: []StockPrice. Slices in Go have a useful property for this pattern: window = window[i:] drops the first i elements in O(1) — it just moves the start pointer without copying data. Evicting old entries is cheap.


5. Step 2: Generating prices

The generator simulates a data stream. It runs in an infinite loop, creates a price with the current timestamp, and sends it down the channel. Then it sleeps for one second before repeating.

func generate(out chan<- StockPrice) {
    for {
        out <- StockPrice{
            Value:      50 + rand.Float64()*50, // random price between $50 and $100
            RecordedAt: time.Now(),
        }
        time.Sleep(time.Second)
    }
}

A few things to note:

  • chan<- StockPrice means the function can only send on this channel, not receive. This is Go's way of documenting intent in the type system — the generator has no business reading from its own output.
  • rand.Float64() returns a value in [0, 1), so 50 + rand.Float64()*50 gives us [50, 100).
  • The sleep is what makes this feel like a real stream. Remove it and prices flood in as fast as the CPU can generate them.

In production this function would instead be reading from a WebSocket, deserialising messages from Kafka, or calling an HTTP endpoint. The rest of the code doesn't care — it just reads from the channel.


6. Step 3: The window processor

This is the core of the pattern. The processor maintains a slice that represents the current window, and uses a select to handle two events: a new price arriving, and a ticker firing.

func process(in <-chan StockPrice, windowSize time.Duration) {
    var window []StockPrice
    tick := time.NewTicker(time.Second)
    defer tick.Stop()

    for {
        select {
        case p := <-in:
            // A new price arrived. Add it to the window.
            window = append(window, p)

        case <-tick.C:
            // One second has passed. Time to evict and compute.

            // Find the first entry that's still inside the window.
            cutoff := time.Now().Add(-windowSize)
            i := 0
            for i < len(window) && window[i].RecordedAt.Before(cutoff) {
                i++
            }
            // Drop everything before that point.
            window = window[i:]

            if len(window) == 0 {
                fmt.Println("No data in window yet.")
                continue
            }

            // Compute the average of what remains.
            sum := 0.0
            for _, p := range window {
                sum += p.Value
            }
            avg := sum / float64(len(window))

            fmt.Printf("[%s] %d prices in window  avg = $%.2f\n",
                time.Now().Format("15:04:05"), len(window), avg)
        }
    }
}

Let's walk through the eviction logic since that's where it's easy to get confused.

cutoff := time.Now().Add(-windowSize) computes a point in time 5 seconds ago. Any price recorded before that point is outside the window and should be dropped.

We then scan the window from the front — since prices are appended in order, the oldest ones are always at the start. Once we find the first entry whose timestamp is not before the cutoff, everything from that point onward is still valid.

window = window[i:] does the eviction. It doesn't copy anything — it just returns a slice starting at index i. The old entries at the front become unreachable and will be garbage collected.

The select statement is what allows the processor to respond to both the channel and the ticker without blocking on either. If a price arrives while we're waiting for the ticker, we handle it immediately. If the ticker fires while we're processing a price, we handle that next. Go's select picks a random branch when multiple are ready at the same time, which is fine here.


7. Step 4: Wiring it together

The main function creates the channel, starts the generator in a goroutine, and hands control to the processor.

func main() {
    ch := make(chan StockPrice, 10)

    go generate(ch) // runs in background, sends prices every second

    process(ch, 5*time.Second) // runs forever on main goroutine
}

The channel is buffered with size 10 as a small cushion. If the processor is briefly slow — say it's computing a more expensive metric — the generator can keep running without blocking, and the backed-up prices will be processed when the processor catches up. In practice with a 1-second sleep, this buffer will almost never be used.

Here's the complete program in one file:

package main

import (
    "fmt"
    "math/rand"
    "time"
)

type StockPrice struct {
    Value      float64
    RecordedAt time.Time
}

func generate(out chan<- StockPrice) {
    for {
        out <- StockPrice{
            Value:      50 + rand.Float64()*50,
            RecordedAt: time.Now(),
        }
        time.Sleep(time.Second)
    }
}

func process(in <-chan StockPrice, windowSize time.Duration) {
    var window []StockPrice
    tick := time.NewTicker(time.Second)
    defer tick.Stop()

    for {
        select {
        case p := <-in:
            window = append(window, p)

        case <-tick.C:
            cutoff := time.Now().Add(-windowSize)
            i := 0
            for i < len(window) && window[i].RecordedAt.Before(cutoff) {
                i++
            }
            window = window[i:]

            if len(window) == 0 {
                fmt.Println("No data in window yet.")
                continue
            }

            sum := 0.0
            for _, p := range window {
                sum += p.Value
            }
            avg := sum / float64(len(window))

            fmt.Printf("[%s] %d prices in window  avg = $%.2f\n",
                time.Now().Format("15:04:05"), len(window), avg)
        }
    }
}

func main() {
    ch := make(chan StockPrice, 10)
    go generate(ch)
    process(ch, 5*time.Second)
}

8. Running it

Save the file as main.go and run it:

go run main.go

For the first few seconds, the window fills up as prices arrive. After 5 seconds, the window reaches its full size and you'll see old prices start dropping off:

[14:32:01] 1 prices in window  avg = $73.41
[14:32:02] 2 prices in window  avg = $76.88
[14:32:03] 3 prices in window  avg = $71.25
[14:32:04] 4 prices in window  avg = $74.60
[14:32:05] 5 prices in window  avg = $72.94
[14:32:06] 5 prices in window  avg = $75.12   ← $73.41 dropped, new price added
[14:32:07] 5 prices in window  avg = $73.88
[14:32:08] 5 prices in window  avg = $74.31

The window size stays at 5 once it's full. Each second, one old price falls off and one new price comes in — the window slides.

To change the window size, change 5*time.Second in main. To change how often prices arrive, change the time.Sleep in generate. These two values are independent: a 10-second window with prices arriving every 100ms would hold up to 100 prices at a time.


9. Taking it further

The core pattern here is simple enough that you can extend it without much difficulty.

Track more metrics. The loop that computes the average is the right place to also compute min, max, standard deviation, or percentiles. Nothing about the window structure changes.

Count-based windows. If you want "the last N prices" instead of "the last N seconds", replace the timestamp-based eviction with a length check: if len(window) > maxSize { window = window[1:] }. The rest stays the same.

Multiple windows. Nothing stops you from running two processors on the same channel — one with a 5-second window for short-term averages, one with a 60-second window for trend detection. Since Go channels can be read by only one goroutine, you'd fan out through two separate channels.

Production scale. For high-volume streams, the in-process slice works well up to tens of thousands of events per second. Beyond that, the usual path is to push the stream into Kafka and use a streaming framework like Apache Flink or Redpanda's WASM-based transforms to manage the windows in a distributed way. The sliding window concept stays the same — only the infrastructure changes.

I build and scale reliable production systems. Open to full-time and freelance work with U.S.-based teams that value ownership and execution.

Got something in mind?

Book a Discovery Call