Process Synchronization with Locks
Simulating Ticket Booking System with Concurrency
Requirements:
- Check available tickets
- 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}')