Coordination

Coordination: queues & backpressure

Producer-consumer for LLD: why busy-wait fails, blocking queues, backpressure, graceful shutdown, actors, async requests, and burst buffering.

Handoff without burning CPU or memory

Producer consumer
API producers → bounded queue → workers.
Coordination is threads communicating and sequencing work. Signup needs a welcome email; upload needs resize; report needs minutes. Handlers enqueue; workers process asynchronously.
Empty queue + busy-wait = 100% CPU for nothing. Sleep-polling trades CPU for latency. Fast producers + unbounded queue = OOM and a dead API. You need: efficient waiting, backpressure, and a thread-safe handoff.

Analogy: take-a-number counter

Queue ticket counter
Producers take numbers; workers call the next.

Coordination is a bakery take-a-number line: customers (producers) don’t shout at the clerk; they queue. If the ticket dispenser is full (bounded queue), new customers wait outside — that’s backpressure — instead of stuffing infinite people into the shop (OOM).

Condition variables (know, rarely code)

Wait releases the lock and parks; notify wakes waiters who reacquire and recheck the condition (spurious wakeups + lost races). Prefer notify_all / separate conditions for “not empty” vs “not full” unless you’re sure one waiter can proceed. In interviews, jump to blocking queues.

Blocking queues (default answer)

Thread-safe queue: put blocks when full, get/take blocks when empty. Synchronization, waiting, and backpressure built in.
from queue import Queue
from threading import Thread


class TaskScheduler:
    def __init__(self, capacity: int = 1000):
        self._q: Queue = Queue(maxsize=capacity)

    def submit(self, task) -> None:
        self._q.put(task)  # blocks if full

    def worker_loop(self) -> None:
        while True:
            task = self._q.get()
            try:
                task()
            finally:
                self._q.task_done()
import java.util.concurrent.*;

class TaskScheduler {
    private final BlockingQueue<Runnable> q;

    TaskScheduler(int capacity) {
        this.q = new ArrayBlockingQueue<>(capacity);
    }

    void submit(Runnable task) throws InterruptedException {
        q.put(task);  // blocks if full
    }

    void workerLoop() throws InterruptedException {
        while (true) {
            Runnable task = q.take();
            task.run();
        }
    }
}
  • Always bound capacity — size ≈ worker_rate × burst_seconds (e.g. 100/s × 10s = 1000).
  • Full queue options: block (put) for internal pipelines; timeout/reject on request paths; drop+log for lossy analytics.
  • Shutdown: interrupt / poison pill / get with timeout + flag.

Message passing / actors (alternative)

Each actor owns private state and a mailbox; processes one message at a time — no locks inside business logic. Great for many independent stateful entities (chat sessions, game rooms, order books). For simple “background tasks,” a queue + workers is simpler.

Pattern: process requests asynchronously

Handler does the minimum, enqueues slow work, returns fast. Email, image pipeline, payment side-effects, long reports — same shape.
def signup(self, email: str, name: str) -> None:
    self.users.save(email, name)
    self.email_queue.put(EmailTask(email, "welcome", name))
    # return immediately — worker sends mail later
void signup(String email, String name) throws InterruptedException {
    users.save(email, name);
    emailQueue.put(new EmailTask(email, "welcome", name));
    // return immediately — worker sends mail later
}

Pattern: absorb bursty traffic

Size workers for normal load; let a bounded queue absorb spikes. News spikes, ticket onsale, webhook floods, batch completion fan-out. Reject with timeout when the buffer is exhausted rather than OOM.
# Pseudocode: offer with timeout on the request path
ok = purchase_queue.offer(request, timeout_ms=100)
if not ok:
    raise ServiceUnavailable("try again")
// Pseudocode: offer with timeout on the request path
boolean ok = purchaseQueue.offer(request, 100, TimeUnit.MILLISECONDS);
if (!ok) {
    throw new ServiceUnavailableException("try again");
}

Decision tree

Worked mini-example: bounded queue backpressure

q = BoundedQueue(maxsize=100)

def producer(item):
    q.put(item)  # blocks when full — backpressure

def consumer():
    while True:
        item = q.take()  # blocks when empty
        process(item)
BlockingQueue<Item> q = new ArrayBlockingQueue<>(100);

void producer(Item item) throws InterruptedException {
    q.put(item);  // blocks when full — backpressure
}

void consumer() throws InterruptedException {
    while (true) {
        Item item = q.take();  // blocks when empty
        process(item);
    }
}

Coordination is about handoff timing — not about whether balance += 1 is atomic.

Coordination anti-patterns

  • Unbounded queues as “async” — memory death under burst.
  • wait/notify without while-predicate loop (spurious wakeups).
  • Fire-and-forget threads with no lifecycle.
  • Actors namedrop with no mailbox semantics.

Coordination checklist

  1. Who waits for whom?
  2. What happens when the buffer is full?
  3. Is work CPU-bound or I/O-bound for async choice?
  4. How do you drain on shutdown?

Common interview pitfalls

These mistakes show up constantly on this prompt. Name the trap, then show the fix in your design — don’t wait for the interviewer to catch you.

  • Unbounded queue as default async.
  • Busy-wait instead of condition wait.
  • No backpressure story.
  • Actor buzzwords without mailbox.
  • Ignoring shutdown/drain.

Interview script (say this)

Read this once out loud before a mock. It’s the spine of a strong answer — not a script to recite robotically.

  1. Coordination is producer/consumer handoff.
  2. Bounded queue gives backpressure when consumers are slow.
  3. wait/notify or blocking queue primitives.
  4. Async pipelines for email/resize/reports — not for seat invariants.
  5. Burst buffers absorb spikes then drain.
  6. Actors: single-threaded mailbox per entity as an option.
  7. I’ll pick queue vs callback based on overload behavior needed.
  8. Trace a full buffer: producer blocks or drops per policy.

Extra verification traces

Walk these three traces on the board. If you can narrate them cleanly, your implementation section usually follows.

Staff-level follow-ups

At staff+, they twist the prompt. Answer in one sentence that names the seam — don’t redesign the whole board.

  • Drop vs block? — Policy by product — telemetry may drop; payments shouldn’t.
  • Kafka? — Same coordination idea at infra zoom — say so, stay LLD if asked.
  • Priority queues? — Still a handoff structure; starvation needs aging.

Complete solution: async logger (bounded queue)

API threads must not block on disk. Complete design: enqueue formatted lines; one worker drains to the sink; full queue = drop or block (pick a policy).
from queue import Queue, Full
from threading import Thread

class AsyncLogger:
    def __init__(self, sink, maxsize: int = 10000):
        self._q: Queue = Queue(maxsize=maxsize)
        self._sink = sink
        self._alive = True
        self._worker = Thread(target=self._run, daemon=True)
        self._worker.start()

    def info(self, msg: str) -> None:
        try:
            self._q.put_nowait(msg)          # backpressure: drop newest
        except Full:
            # alternative: put() to block caller — say which you chose
            pass

    def _run(self) -> None:
        while self._alive or not self._q.empty():
            line = self._q.get()
            try:
                self._sink.write(line + "\n")
            finally:
                self._q.task_done()

    def shutdown(self) -> None:
        self._alive = False
        self._q.put("__STOP__")              # wake worker
        self._worker.join(timeout=5)

# Concurrency: many producers, one consumer — Queue is the lock.
# Never share the file handle across threads without its own lock.
import java.util.concurrent.*;
import java.io.Writer;

class AsyncLogger {
    private final BlockingQueue<String> q;
    private final Writer sink;
    private volatile boolean alive = true;
    private final Thread worker;

    AsyncLogger(Writer sink, int maxsize) {
        this.q = new ArrayBlockingQueue<>(maxsize);
        this.sink = sink;
        this.worker = new Thread(this::run, "async-logger");
        this.worker.setDaemon(true);
        this.worker.start();
    }

    void info(String msg) {
        if (!q.offer(msg)) {          // backpressure: drop newest
            // alternative: q.put(msg) to block caller — say which you chose
        }
    }

    private void run() {
        try {
            while (alive || !q.isEmpty()) {
                String line = q.take();
                if ("__STOP__".equals(line)) break;
                sink.write(line + "\n");
            }
        } catch (Exception e) {
            // log / rethrow as needed
        }
    }

    void shutdown() throws InterruptedException {
        alive = false;
        q.offer("__STOP__");              // wake worker
        worker.join(5000);
    }
}

// Concurrency: many producers, one consumer — BlockingQueue is the lock.
// Never share the file handle across threads without its own lock.

← Lattice