1. Tổng Quan Về Luồng Trong Python
Mô hình thực thi của Python được điều khiển bởi Máy Ảo Python (PyVM). Ngay từ thiết kế ban đầu, ngôn ngữ này đã quy định chỉ cho phép một luồng chạy trên máy chủ tại bất kỳ thời điểm nào. Mặc dù có thể khởi tạo nhiều luồng, nhưng cơ chế Global Interpreter Lock (GIL) đảm bảo rằng chỉ có duy nhất một luồng thao tác với mã bytecode của Python trong mỗi chu kỳ.
GIL (Global Interpreter Lock)
Cơ chế khóa toàn cục này kiểm soát quyền truy cập vào máy ảo. Quy trình xử lý đa luồng diễn ra như sau:
- Khoa GIL được thiết lập.
- Luồng hiện tại được chọn để chạy.
- Luồng thực thi một số lượng bytecode nhất định hoặc tự nhả控制权 (ví dụ: gọi
time.sleep()). - Luồng chuyển sang trạng thái chờ.
- Mở khóa GIL.
- Lặp lại các bước trên.
Khi gọi các hàm mở rộng viết bằng C/C++, GIL thường bị giữ nguyên cho đến khi hàm kết thúc. Các nhà phát triển thư viện cấp thấp có thể tùy chỉnh hành vi này để giải phóng khóa cho các tác vụ nặng.
Lựa Chọn Mô-đun
Python cung cấp các công cụ cho lập trình song song bao gồm threading, thread (đã lỗi thời) và queue. Chúng tôi khuyến nghị sử dụng threading vì tính năng quản lý phong phú hơn so với thread. Đặc biệt, thread không hỗ trợ luồng bảo vệ (daemon thread) và sẽ chấm dứt mọi luồng con khi luồng chính tắt mà không thông báo, trong khi threading xử lý vấn đề này an toàn hơn.
2. Làm Việc Với Module Threading
Khai Báo Luồng
Có hai cách phổ biến để khởi tạo đối tượng luồng:
Cách 1: Sử dụng hàm mục tiêu
from threading import Thread
import logging
def worker_process(id):
logging.info(f"Chức năng đang chạy ID: {id}")
# Giả lập xử lý
import time
time.sleep(1.5)
if __name__ == "__main__":
# Khởi tạo luồng
task_thread = Thread(target=worker_process, args=(101,))
task_thread.start()
print("Luồng chính đang thực thi tiếp tục...")
Cách 2: Kế thừa lớp Thread
from threading import Thread
import time
class CustomWorker(Thread):
def __init__(self, worker_id):
super().__init__()
self.id = worker_id
def run(self):
time.sleep(1.5)
print(f"Dữ liệu luồng {self.id} đã hoàn tất.")
if __name__ == "__main__":
my_worker = CustomWorker(205)
my_worker.start()
print("Luồng chính tiếp tục xử lý.")
So Sánh Luồng Và Tiến Trình
Sự khác biệt lớn nhất nằm ở việc chia sẻ bộ nhớ và ID tiến trình (PID).
- Luồng: Chia sẻ cùng không gian địa chỉ bộ nhớ với luồng chính, cùng chung PID.
- Tiến trình: Có không gian bộ nhớ độc lập, PID riêng biệt do hệ điều hành fork.
import os
from threading import Thread
from multiprocessing import Process
def check_id():
print(f"ID bên trong: {os.getpid()}")
if __name__ == '__main__':
parent_pid = os.getpid()
# Kiểm tra luồng
t = Thread(target=check_id)
t.start()
print(f"PID Luồng chính: {parent_pid}")
# Kết quả: Cùng PID
# Kiểm tra tiến trình
p = Process(target=check_id)
p.start()
print(f"PID Luồng chính: {parent_pid}")
# Kết quả: Luồng phụ có PID khác
Về hiệu suất, việc tạo mới luồng nhẹ hơn so với tạo tiến trình. Tuy nhiên, việc chia sẻ dữ liệu giữa các luồng cần sự đồng bộ cẩn thận để tránh xung đột.
shared_resource = 500
def modify_resource():
global shared_resource
temp_val = shared_resource
shared_resource = temp_val - 100
if __name__ == '__main__':
thread_a = Thread(target=modify_resource)
thread_b = Thread(target=modify_resource)
thread_a.start()
thread_b.start()
thread_a.join()
thread_b.join()
print(f"Giá trị còn lại: {shared_resource}")
Socket Đa Luồng Mẫu
Dưới đây là ví dụ về server chấp nhận nhiều kết nối đồng thời bằng cách sử dụng luồng riêng cho mỗi client.
Server
import socket
import threading
def handle_connection(sock_conn):
while True:
try:
data = sock_conn.recv(1024)
if not data: break
print(f"Nhận dữ liệu: {data}")
sock_conn.sendall(data.upper())
except:
break
sock_conn.close()
server_sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
server_sock.bind(('localhost', 9900))
server_sock.listen(5)
while True:
conn, addr = server_sock.accept()
client_thread = threading.Thread(target=handle_connection, args=(conn,))
client_thread.start()
Client
import socket
client = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
client.connect(('localhost', 9900))
msg = input("Nhập lệnh gửi: ")
client.send(msg.encode('utf-8'))
response = client.recv(1024)
print(f"Phản hồi server: {response}")
client.close()
Các Phương Thức Quản Lý Luồng
Lớp Thread và module threading cung cấp nhiều tiện ích:
is_alive(): Kiểm tra trạng thái hoạt động.join(): Chặn luồng chính cho đến khi luồng con kết thúc.setDaemon(True): Đặt luồng thành luồng bảo vệ (kết thúc cùng luồng chính).enumerate(): Danh sách các luồng đang chạy.
import threading
import time
def monitor_task():
time.sleep(5)
print(threading.current_thread().getName())
if __name__ == "__main__":
worker = threading.Thread(target=monitor_task)
worker.setName("Monitoring_Service")
worker.start()
print(f"Tên luồng hiện tại: {threading.current_thread().getName()}")
print(f"Số lượng luồng hoạt động: {threading.active_count()}")
print(f"Danh sách luồng: {threading.enumerate()}")
# Đợi luồng con hoàn thành
worker.join()
print(f"Trạng thái luồng con: {worker.is_alive()}")
Luồng Bảo Vệ (Daemon Thread)
Luồng bảo vệ không ngăn chặn chương trình kết thúc. Khi luồng chính dừng, các luồng daemon sẽ bị hủy bỏ ngay lập tức. Điều này hữu ích cho các nhiệm vụ nền như theo dõi tài nguyên.
def background_service():
while True:
print("Đang giám sát...")
time.sleep(2)
d_thread = threading.Thread(target=background_service, daemon=True)
d_thread.start()
print("Luồng chính kết thúc")
# Background service sẽ dừng ngay khi dòng này chạy xong
3. Khóa và Đồng Bộ Hóa (Locking)
Vấn Đề Xung Đột Dữ Liệu
Khi nhiều luồng thao tác cùng một biến toàn cục mà không có cơ chế bảo vệ, giá trị cuối cùng có thể không chính xác do tình trạng cạnh tranh (race condition).
# Không dùng khóa - Kết quả không ổn định
balance = 10000
threads_list = []
def unsafe_withdraw(amount):
global balance
local_check = balance
time.sleep(0.05) # Giả lập trễ IO
balance = local_check - amount
for _ in range(10):
t = threading.Thread(target=unsafe_withdraw, args=(1000,))
threads_list.append(t)
t.start()
for t in threads_list:
t.join()
print(balance)
Giải pháp là sử dụng threading.Lock để biến khối mã thành đoạn tín hiệu (critical section).
from threading import Thread, Lock
import time
lock_mutex = Lock()
safe_balance = 10000
def safe_withdraw(value):
global safe_balance
lock_mutex.acquire()
try:
current = safe_balance
time.sleep(0.05)
safe_balance = current - value
finally:
lock_mutex.release() # Đảm bảo giải phóng khóa
# Sử dụng context manager tốt hơn
def better_safe_withdraw(value):
global safe_balance
with lock_mutex:
curr = safe_balance
time.sleep(0.05)
safe_balance = curr - value
Khóa Đệ Quy (RLock)
Lock thông thường sẽ gây deadlock nếu cùng một luồng cố gắng khóa lại nhiều lần. RLock (Reentrant Lock) cho phép điều này miễn là số lần giải phóng bằng số lần khóa.
import threading
recursive_lck = threading.RLock()
def inner_call():
recursive_lck.acquire()
print("Hàm con lấy khóa")
recursive_lck.release()
def outer_call():
recursive_lck.acquire()
print("Hàm chính lấy khóa")
inner_call()
recursive_lck.release()
outer_call()
Ví Dụ Deadlock Kinh Điển
Scenario: Hai người muốn ăn mì cần cả đũa và bát. Nếu người A giữ đũa chờ bát, người B giữ bát chờ đũa -> Tắc nghẽn.
# Giải quyết Deadlock bằng RLock hoặc thứ tự khóa thống nhất
r_lock_noodle = threading.RLock()
r_lock_fork = threading.RLock()
def eat_person_1(name):
with r_lock_noodle:
print(f"{name}: Nắm bát")
with r_lock_fork:
print(f"{name}: Nắm đũa -> Ăn")
def eat_person_2(name):
with r_lock_fork:
print(f"{name}: Nắm đũa")
with r_lock_noodle:
print(f"{name}: Nắm bát -> Ăn")
4. Các Công Cụ Đồng Bộ Nâng Cao
Signal Nhẹ (Event)
Event cho phép các luồng chờ một cờ báo hiệu trở nên đúng (true) trước khi tiếp tục.
import threading
import random
signal_event = threading.Event()
def database_checker():
print("Kiểm tra DB...")
time.sleep(random.randint(2, 5))
signal_event.set() # Đánh thức các luồng chờ
def app_worker():
while not signal_event.is_set():
print("Đợi DB sẵn sàng...")
signal_event.wait(timeout=1)
print("Kết nối DB thành công")
t_checker = threading.Thread(target=database_checker)
t_worker = threading.Thread(target=app_worker)
t_worker.start()
t_checker.start()
Semaphore
Điều giới hạn số lượng luồng truy cập tài nguyên cùng lúc (ví dụ: giới hạn kết nối max 5).
semaphore_limit = threading.Semaphore(5)
def limit_task():
with semaphore_limit:
print(f"{threading.current_thread().name} đang xử lý")
time.sleep(1.5)
for i in range(8):
t = threading.Thread(target=limit_task)
t.start()
Condition
Cho phép luồng chờ cho đến khi một điều kiện cụ thể thay đổi, thường dùng cho Producer-Consumer.
cond_var = threading.Condition()
def producer():
with cond_var:
print("Sản phẩm đã sẵn sàng")
cond_var.notify_all()
def consumer():
with cond_var:
cond_var.wait()
print("Tiêu thụ sản phẩm")
# Khởi chạy trong vòng lặp tương tự
Hàng Đợi An Toàn (Queue)
Module queue cung cấp các cấu trúc dữ liệu an toàn cho luồng để trao đổi thông tin mà không lo khóa thủ công.
FIFO (First In First Out)
import queue
my_queue = queue.Queue()
my_queue.put("Item A")
my_queue.put("Item B")
print(my_queue.get()) # Item A
LIFO (Last In First Out - Stack)
stack_q = queue.LifoQueue()
stack_q.put(1)
stack_q.put(2)
print(stack_q.get()) # 2
Priority Queue
priority_q = queue.PriorityQueue()
priority_q.put((1, "Urgent"))
priority_q.put((5, "Normal"))
print(priority_q.get()) # (1, 'Urgent')
5. Concurrent Futures (Cấp Cao)
Thư viện concurrent.futures trừu tượng hóa việc quản lý hồ luồng/tiến trình, giúp viết mã đồng thuận ngắn gọn hơn.
ThreadPoolExecutor & ProcessPoolExecutor
from concurrent.futures import ThreadPoolExecutor
import time
def calculate_square(n):
time.sleep(1)
return n * n
with ThreadPoolExecutor(max_workers=3) as executor:
# Sử dụng map trả về iterator kết quả
results = list(executor.map(calculate_square, [1, 2, 3, 4, 5]))
print(results)
# Hoặc submit từng任务
future = executor.submit(calculate_square, 10)
print(future.result())
Với ProcessPoolExecutor, cú pháp tương tự nhưng phù hợp cho tác vụ chiếm dụng CPU nhờ vượt qua giới hạn GIL.
from concurrent.futures import ProcessPoolExecutor
def heavy_computation(num):
# Tính toán nặng
return sum(i*i for i in range(1000000))
if __name__ == '__main__':
with ProcessPoolExecutor(max_workers=2) as pool:
for res in pool.map(heavy_computation, range(5)):
print(res)