Safe Inter-Thread Data Exchange in Python Using Queue

In multi-threaded programming environments, directly sharing mutable state among processes creates significant risks of race conditions. The most reliable solution for passing data between concurrent execution paths is the queue.Queue class found in the standard library. By instantiating a shared Queue object, individual threads can safely interact with data through put() and get() operations. The internal implementation handles all necessary locking mechanisms, ensuring that access remains atomic and secure without manual semaphore management.

Architectural Overview

The core architecture involves separating concerns between data generators and data processors. This separation minimizes coupling:

  • Data Generator: Creates information and pushes it into the queue buffer.
  • Data Processor: Retrieves available items from the buffer for execution.

This model supports scaling; multiple generators can feed into a single queue, and multiple processors can consume from the same stream.

Implementation Example: Basic Producer-Consumer

Below is a demonstration establishing the fundamental communication pattern. Note the use of a termination flag to ensure clean shutdown rather than infinite looping.

import threading
import queue
import time

# Initialize a FIFO queue with a defined capacity limit
shared_queue = queue.Queue(maxsize=10)

class ItemProducer(threading.Thread):
    def __init__(self, name, stop_event):
        super().__init__()
        self.name = name
        self.stop_event = stop_event

    def run(self):
        counter = 1
        while not self.stop_event.is_set():
            print(f"[{self.name}] Generating item")
            item = f"Item-{counter}"
            
            # Blocks if queue is full (maxsize reached)
            shared_queue.put(item) 
            print(f"[{self.name}] Produced: {item}")
            
            counter += 1
            time.sleep(2)

class ItemConsumer(threading.Thread):
    def __init__(self, name, stop_event):
        super().__init__()
        self.name = name
        self.stop_event = stop_event

    def run(self):
        while not self.stop_event.is_set():
            try:
                # Blocks if queue is empty until data arrives
                data = shared_queue.get(timeout=1)
                print(f"[{self.name}] Consumed: {data}")
                
                # Signal that processing of this specific task is finished
                shared_queue.task_done()
                time.sleep(1)
            except queue.Empty:
                continue

if __name__ == "__main__":
    # Create an event for signaling termination
    terminate_signal = threading.Event()

    # Instantiate producer and three consumers
    producer = ItemProducer("Factory-A", terminate_signal)
    
    consumers = []
    for i in range(3):
        c = ItemConsumer(f"Worker-{i+1}", terminate_signal)
        consumers.append(c)

    # Launch all threads
    producer.start()
    for c in consumers:
        c.start()

    # Let the system run briefly then signal shutdown
    time.sleep(5)
    terminate_signal.set()
    
    # Wait for threads to finish gracefully
    producer.join()
    for c in consumers:
        c.join()
        
    print("System termination complete.")

API Configuration and Methods

Understanding the configuration parameters of the Queue module is critical for preventing deadlocks and managing backpressure effectively.

Method Behavior Description
Queue(size) Initializes the buffer. If zero, size is unbounded.
put(item, block=True, timeout=None) Adds a element. If block=True, waits for space. If block=False, raises Full immediately if buffer is saturated.
get(block=True, timeout=None) Retrieves an element. If block=True, waits for content. If block=False, raises Empty immediately if buffer is vacant.
qsize() Returns approximate count of items currently in queue.
empty() Boolean check: Returns true if buffer contains no items.
full() Boolean check: Returns true if buffer has reached max capacity.
task_done() Indicates that a previously retrieved item has been fully processed. Required for accurate tracking of pending work.
join() Blocks the calling thread until every item pushed to the queue has been popped and marked as task_done.

Complex Scenario: Batch Processing with Metrics

This example addresses a high-volume production requirement where a master thread orchestrates data creation across batches, while several worker threads handle extraction. It also captures execution latency to analyze throughput performance.

from queue import Queue
import threading
import time

# Setup global queue
batch_queue = Queue()
total_items_expected = 0
items_processed = 0
lock = threading.Lock()

def execute_production_batch():
    """
    Generates data in batches.
    Each cycle produces 15 chunks, paused for 0.5 seconds between batches.
    """
    global total_items_expected
    batch_count = 4
    items_per_batch = 15
    
    total_items_expected = batch_count * items_per_batch

    for i in range(batch_count):
        for j in range(items_per_batch):
            msg = f"Batch-{i}_Seq-{j}"
            batch_queue.put(msg)
        
        time.sleep(0.5)

def execute_worker_process(worker_id):
    """
    Continuously attempts to pull items for processing.
    Uses a timeout to periodically check if production is halted.
    """
    processed_local = 0
    
    while True:
        try:
            data = batch_queue.get(timeout=2)
            # Simulate heavy computation
            time.sleep(0.2) 
            
            print(f"[Worker-{worker_id}] Received: {data}")
            batch_queue.task_done()
            processed_local += 1
            
        except queue.Empty:
            # Check if producer finished and queue is empty
            with lock:
                if batch_queue.empty():
                    pass
            break
        finally:
            # Update global stats atomically
            with lock:
                global items_processed
                items_processed += 1

if __name__ == "__main__":
    start_time = time.time()

    # Define worker pool
    active_workers = []
    worker_threads = 5
    
    for idx in range(worker_threads):
        t = threading.Thread(target=execute_worker_process, args=(idx,))
        t.daemon = False
        t.start()
        active_workers.append(t)

    # Start the production pipeline
    producer_thread = threading.Thread(target=execute_production_batch)
    producer_thread.start()
    producer_thread.join()

    # Wait for all workers to drain the queue
    for w in active_workers:
        w.join()

    # Ensure queue tasks are fully acknowledged
    batch_queue.join()

    end_time = time.time()
    duration = end_time - start_time

    print(f"Pipeline Finished. Total time: {duration:.2f} seconds")
    print(f"Total items handled: {items_processed}")

Tags: python multithreading queue-class concurrency-control producer-consumer-pattern

Posted on Tue, 06 Oct 2026 16:00:45 +0000 by eludlow