How a pipe, a lock, and a slow consumer deadlocked a crawler — and how to fix it
A production crawler froze mid-run: the process was alive, using ~0% CPU, but producing no logs, no errors, and no exit.
Not a crash — the process was blocked, waiting for something that would never finish.
Progress logs stopped because the print came after a blocking write().
One stalled worker turned into a total freeze when every other thread waited behind it.
This happened in the Indeed crawler: crawler | feeder, where the feeder did slow S3/SQS/MySQL work.
Workers write NDJSON to stdout (a pipe). The feeder reads it and does slow I/O in between.
A classic concurrency pattern: producers generate data, consumers process it, and a bounded buffer connects them.
When the consumer is slower than the producers, the buffer fills — this is called backpressure.
A lock guarantees mutual exclusion: only one thread enters a critical section at a time.
Rule of thumb: never hold a lock across blocking I/O, because the lock then blocks everyone else too.
write() to a full pipe does not return until the reader drains space. The OS gives you no timeout and no error — it just blocks.
This bug is a form of deadlock: one thread holds a lock while blocked on I/O, and the rest wait on the lock — nobody can proceed.
write() blocks while holding the lockOne blocking write() under a shared lock = a frozen thread pool.
data = build_record()
queue.put(data) # fast, non-blocking
Producers enqueue and return. One dedicated thread does the I/O.
with lock:
write(data) # blocks if full!
Holding the lock across write() couples every worker to the slowest I/O.
Every job record was written under one shared lock:
record = self.extract_job(job, query, location, crawl_id)
line = json.dumps(record, ensure_ascii=False) + "\n"
with self._output_lock:
for out_f in outputs:
out_f.write(line) # blocks if pipe is full
out_f.flush()
total_written += 1
The bug is not specific to crawlers. Here is the smallest faithful reproduction.
import threading
import time
lock = threading.Lock()
pipe = [] # stand-in for the 64 KB OS pipe
def producer():
for i in range(50):
with lock: # BUG: lock held across a blocking write
while len(pipe) >= 20:
time.sleep(0.01) # pretend write() blocks on a full pipe
pipe.append("job-" + str(i))
time.sleep(0.05)
def consumer():
while True:
if pipe:
pipe.pop(0)
time.sleep(0.15) # slow consumer -> backpressure
threading.Thread(target=consumer, daemon=True).start()
producer()
Run this and the producer eventually stalls inside the lock while the slow consumer drains — the classic freeze.
Move the blocking I/O off the worker threads onto one dedicated writer thread fed by a bounded queue.
class NdjsonWriter:
def __init__(self, streams, maxsize=20000):
self.queue = queue.Queue(maxsize=maxsize)
self._thread = threading.Thread(
target=self._drain, daemon=True)
self._thread.start()
def write(self, line):
self.queue.put(line) # non-blocking for workers
def _drain(self):
while True:
line = self.queue.get()
if line is None:
return
for stream in self.streams:
stream.write(line)
stream.flush()
Worker threads enqueue and return immediately — no lock is held across I/O.
A single drain thread performs all writes, so a stall affects only that thread.
The writer warns on stderr when the queue backs up, so a stall is never silent.
| Aspect | ❌ Before | ✅ After |
|---|---|---|
| Lock held across write | Yes | No |
| Buffer headroom | ~20 records (64 KB pipe) | 20,000 records (queue) |
| Stall visibility | Silent freeze | Explicit warning |
| Consumer stall impact | Whole pool freezes | One drain thread blocks |