Process Synchronization with Locks, Queues, Producer-Consumer Pattern, and Threading

Process Synchronization with Locks

Simulating Ticket Booking System with Concurrency

Requirements:

  1. Check available tickets
  2. Purchase tickets

Concurrent Ticket Purchase Leading to Data Corruptoin

import multiprocessing
import time
import json

def check_tickets(user_id):
    with open('ticket_data.json', 'r') as file:
        data = json.load(file)
        available = data['available']
        print(f'User {user_id} sees {available} tickets available')
        return available

def purchase_ticket(user_id):
    available = check_tickets(user_id)
    if available <= 0:
        print('No tickets available!')
        return
    time.sleep(0.5)
    available -= 1
    with open('ticket_data.json', 'w') as file:
        json.dump({'available': available}, file)
        print(f'User {user_id} purchased successfully. Remaining: {available}')

for i in range(10):
    proc = multiprocessing.Process(target=purchase_ticket, args=(i,))
    proc.start()

Process Lock for Synchronization

import multiprocessing
import time
import json

def check_availability(user_id):
    time.sleep(0.5)
    with open('ticket_data.json', 'r') as file:
        data = json.load(file)
        available = data['available']
        print(f'User {user_id} sees {available} tickets')

def update_tickets(user_id):
    time.sleep(0.5)
    with open('ticket_data.json', 'r') as file:
        data = json.load(file)
        available = data['available']
    if available <= 0:
        print('No tickets left!')
        return
    time.sleep(0.5)
    available -= 1
    with open('ticket_data.json', 'w') as file:
        json.dump({'available': available}, file)
        print(f'User {user_id} purchased. Remaining: {available}')

def book_ticket(user_id, lock):
    check_availability(user_id)
    lock.acquire()
    update_tickets(user_id)
    lock.release()

lock = multiprocessing.Lock()
for i in range(10):
    p = multiprocessing.Process(target=book_ticket, args=(i, lock))
    p.start()

Thread Lock with Context Manager

import threading
import time

counter = 20

def modify_counter(thread_id):
    global counter
    with lock:
        current = counter
        time.sleep(0.1)
        counter = current - 1
        print(counter)

lock = threading.Lock()
threads = []
for i in range(10):
    t = threading.Thread(target=modify_counter, args=(i,))
    t.start()
    threads.append(t)

for t in threads:
    t.join()

print('Final value:', counter)

Process Queues

Queue Characteristics

  • FIFO data structure in memory
  • put() blocks when full
  • put_nowait() raises exception when full
  • get() blocks when empty
  • get_nowait() raises exception when empty
import multiprocessing

queue = multiprocessing.Queue(2)
queue.put(1)
queue.put(2)
print(queue.empty())
print(queue.full())

Inter-Process Communication via Queues

import multiprocessing

class Producer(multiprocessing.Process):
    def __init__(self, queue):
        super().__init__()
        self.queue = queue
        
    def run(self):
        data = f'Data from process {self.pid}'
        self.queue.put(data)
        print(f'Process {self.pid} sent data')

class Consumer(multiprocessing.Process):
    def __init__(self, queue):
        super().__init__()
        self.queue = queue
        
    def run(self):
        data = self.queue.get()
        print(f'Process {self.pid} received: {data}')

q = multiprocessing.Queue()
p1 = Producer(q)
p2 = Consumer(q)
p1.start()
p2.start()

Producer-Consumer Pattern

Basic Implementation

import multiprocessing
import random
import time

def producer(queue):
    for _ in range(10):
        time.sleep(random.random())
        value = random.randint(1, 100)
        queue.put(value)
        print(f'Produced: {value}')
    queue.put(None)

def consumer(queue):
    while True:
        time.sleep(random.uniform(0, 2))
        value = queue.get()
        if value is None:
            print('Consumption complete')
            return
        print(f'Consumed: {value}')

q = multiprocessing.Queue()
p = multiprocessing.Process(target=producer, args=(q,))
c = multiprocessing.Process(target=consumer, args=(q,))
p.start()
c.start()

JoinableQueue for Multiple Producers/Consumers

import multiprocessing
import random
import time

def producer(queue, id):
    for _ in range(5):
        time.sleep(random.random())
        value = random.randint(100, 200)
        queue.put(value)
        print(f'Producer {id} produced: {value}')

def consumer(queue, id):
    while True:
        time.sleep(random.uniform(0, 2))
        value = queue.get()
        queue.task_done()
        print(f'Consumer {id} consumed: {value}')

q = multiprocessing.JoinableQueue()
producers = []
for i in range(3):
    p = multiprocessing.Process(target=producer, args=(q, i))
    p.start()
    producers.append(p)

for i in range(2):
    c = multiprocessing.Process(target=consumer, args=(q, i))
    c.daemon = True
    c.start()

for p in producers:
    p.join()

q.join()

Threading Fundamentals

Thread Creation Methods

Using Thread Class

import threading
import time

def thread_task(thread_id):
    print(f'Thread {thread_id} started')
    time.sleep(1)
    print(f'Thread {thread_id} completed')

t = threading.Thread(target=thread_task, args=(1,))
t.start()

Custom Thread Class

import threading
import time

class CustomThread(threading.Thread):
    def __init__(self, value):
        super().__init__()
        self.value = value
        
    def run(self):
        print(f'Thread with value {self.value} started')
        time.sleep(1)
        print('Thread completed')

t = CustomThread(42)
t.start()

Performance Comparison: Proecsses vs Threads

import threading
import multiprocessing
import time
import os

def cpu_intensive():
    total = 0
    for _ in range(10000000):
        total += 1

def io_intensive():
    time.sleep(1)

# Test with 50 processes/threads
start = time.time()
processes = []
for i in range(50):
    p = multiprocessing.Process(target=cpu_intensive)
    # p = threading.Thread(target=io_intensive)
    p.start()
    processes.append(p)

for p in processes:
    p.join()

print(f'Execution time: {time.time() - start}')

Thread Synchronization

Shared Data Access

import threading

shared_data = 100

def modify_data():
    global shared_data
    shared_data = 200

t = threading.Thread(target=modify_data)
t.start()
t.join()
print(shared_data)  # Output: 200

Thread Lock for Data Consistency

import threading

counter = 0
lock = threading.Lock()

def increment():
    global counter
    with lock:
        for _ in range(200000):
            counter += 1

threads = []
for i in range(2):
    t = threading.Thread(target=increment)
    t.start()
    threads.append(t)

for t in threads:
    t.join()

print(f'Final counter value: {counter}')

Tags: multiprocessing threading Synchronization producer-consumer Queue

Posted on Sat, 12 Sep 2026 16:56:40 +0000 by napier_matt