Hướng Dẫn Toàn Diện Về Đa Luồng Và Đồng Bộ Hóa Trong Python

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:

  1. Khoa GIL được thiết lập.
  2. Luồng hiện tại được chọn để chạy.
  3. Luồng thực thi một số lượng bytecode nhất định hoặc tự nhả控制权 (ví dụ: gọi time.sleep()).
  4. Luồng chuyển sang trạng thái chờ.
  5. Mở khóa GIL.
  6. 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)

Thẻ: python threading gil Concurrency multiprocessing

Đăng vào ngày 22 tháng 9 lúc 08:52