Bỏ qua để đến nội dung

Multiprocessing trong Python: dùng hết mọi lõi CPU

Ở bài GIL và Threading, bạn đã thấy thread không giúp tăng tốc tính toán bằng Python thuần vì GIL. Giải pháp tiêu chuẩn là chạy nhiều tiến trình (process): mỗi tiến trình có interpreter riêng, bộ nhớ riêng và GIL riêng, nên chúng thực sự chạy song song trên nhiều lõi.

Nhưng tiến trình không miễn phí. Chúng không dùng chung bộ nhớ, mọi dữ liệu trao đổi phải được tuần tự hoá (pickle) và gửi qua đường ống. Dùng sai cách, chương trình đa tiến trình có thể chậm hơn bản tuần tự hàng nghìn lần - bạn sẽ thấy ví dụ thật trong bài.

Trong bài này, bạn sẽ học:

  • Vì sao tiến trình vượt qua được GIL và cái giá phải trả
  • Ba cách khởi tạo tiến trình: fork, spawn, forkserver và vì sao cần if __name__ == "__main__"
  • Dùng ProcessPoolExecutormultiprocessing.Pool hiệu quả
  • chunksize - tham số có thể làm chương trình nhanh gấp 500 lần
  • Giới hạn của pickle và cách xử lý
  • Chia sẻ dữ liệu: Queue, Value/Array, Manager, và shared_memory zero-copy
  • Các mẫu thiết kế và bẫy thường gặp

Số liệu đo trên máy 8 lõi, CPython 3.13, macOS.

Một tiến trình, nhiều thread Nhiều tiến trình
┌──────────────────────────────────┐ ┌────────────┐ ┌────────────┐ ┌────────────┐
│ Bộ nhớ chung (mọi object) │ │ Bộ nhớ A │ │ Bộ nhớ B │ │ Bộ nhớ C │
│ Một GIL │ │ GIL A │ │ GIL B │ │ GIL C │
│ ┌────┐ ┌────┐ ┌────┐ │ │ ┌────┐ │ │ ┌────┐ │ │ ┌────┐ │
│ │ T1 │ │ T2 │ │ T3 │ │ │ │main│ │ │ │main│ │ │ │main│ │
│ └────┘ └────┘ └────┘ │ │ └────┘ │ │ └────┘ │ │ └────┘ │
└──────────────────────────────────┘ └─────┬──────┘ └─────┬──────┘ └─────┬──────┘
└──── pipe / queue (pickle) ───┘
Thread Process
Song song với code Python thuần ❌ (có GIL)
Chi phí khởi tạo rất nhỏ lớn (tạo interpreter mới, import lại module)
Chia sẻ dữ liệu trực tiếp, cần khoá phải gửi qua pickle, hoặc dùng shared memory
Bộ nhớ dùng chung mỗi tiến trình một bản
Một worker crash có thể kéo sập cả chương trình các tiến trình khác vẫn sống

Ví dụ đầu tiên: tăng tốc tác vụ CPU-bound

Phần tiêu đề “Ví dụ đầu tiên: tăng tốc tác vụ CPU-bound”
import time
from concurrent.futures import ProcessPoolExecutor
def count(n):
total = 0
for i in range(n):
total += i
return total
if __name__ == "__main__":
N = 20_000_000
start = time.perf_counter()
results = [count(N) for _ in range(4)]
print(f"Tuần tự: {time.perf_counter() - start:.2f}s") # 2.68s
start = time.perf_counter()
with ProcessPoolExecutor(max_workers=4) as pool:
results = list(pool.map(count, [N] * 4))
print(f"4 tiến trình: {time.perf_counter() - start:.2f}s") # 0.78s

Nhanh gấp ~3.4 lần với 4 tiến trình. Không đạt đúng 4 lần vì mất thời gian khởi động các tiến trình con.

API ProcessPoolExecutor giống hệt ThreadPoolExecutor - bạn có thể đổi qua lại chỉ bằng một dòng. Đây là lý do nên học concurrent.futures trước.

Vì sao bắt buộc phải có if __name__ == "__main__":?

Phần tiêu đề “Vì sao bắt buộc phải có if __name__ == "__main__":?”

Hãy tìm hiểu cách tiến trình con được tạo ra. Python có ba start method:

Hệ điều hành nhân bản nguyên tiến trình cha, bao gồm toàn bộ bộ nhớ (theo cơ chế copy-on-write). Tiến trình con có sẵn mọi biến, mọi module đã import.

  • ✅ Khởi tạo rất nhanh.
  • Không an toàn khi tiến trình cha có nhiều thread: chỉ thread gọi fork được nhân bản, nhưng mọi khoá đang bị thread khác giữ cũng bị copy ở trạng thái “đang khoá” → tiến trình con có thể treo vĩnh viễn. Nhiều thư viện (logging, thư viện mạng, CUDA) dùng thread ngầm.
  • ❌ Không có trên Windows; không an toàn trên macOS.

Khởi động một interpreter Python hoàn toàn mới, rồi import lại module chính của bạn để lấy các hàm cần chạy.

  • ✅ Sạch sẽ, an toàn, chạy mọi nơi.
  • ❌ Chậm hơn (mỗi tiến trình con phải import lại mọi thứ).

Chính vì spawn import lại file của bạn, nếu code tạo tiến trình nằm ở cấp module, tiến trình con khi import sẽ lại tạo tiến trình con… vô hạn. if __name__ == "__main__": đảm bảo đoạn code đó chỉ chạy trong tiến trình gốc (trong tiến trình con, __name__"__mp_main__").

forkserver (mặc định trên Linux từ Python 3.14)

Phần tiêu đề “forkserver (mặc định trên Linux từ Python 3.14)”

Lúc đầu, Python khởi động một tiến trình “máy chủ” đơn luồng sạch sẽ. Mỗi khi cần tiến trình con, máy chủ này fork chính nó. Vừa nhanh gần như fork, vừa an toàn vì máy chủ không có thread nào.

import multiprocessing as mp
if __name__ == "__main__":
# Cách 1: đặt MỘT lần, ngay đầu chương trình (trước khi tạo bất kỳ tiến trình nào)
mp.set_start_method("spawn")
print(mp.get_start_method()) # 'spawn'
# Cách 2: dùng context riêng, không ảnh hưởng toàn cục
ctx = mp.get_context("forkserver") # forkserver không có trên Windows
q = ctx.Queue()
# ctx.Process(...), ctx.Pool(...)

Khi bạn gọi pool.map(func, items), điều thực sự xảy ra là:

Tiến trình cha Tiến trình con
pickle(func, item) ──► pipe ──► unpickle ──► func(item) ──► pickle(kết quả)
unpickle(kết quả) ◄──────────────── pipe ◄───────────────────────┘

Hệ quả thứ nhất: hàm và dữ liệu phải pickle được.

from concurrent.futures import ProcessPoolExecutor
if __name__ == "__main__":
with ProcessPoolExecutor(2) as pool:
list(pool.map(lambda x: x * 2, [1, 2, 3])) # lỗi lộ ra khi đọc kết quả
# PicklingError: Can't pickle <function <lambda> ...>

Những thứ không pickle được: lambda, hàm lồng trong hàm, generator, file đang mở, socket, lock thường, kết nối database. Cách xử lý:

  • Định nghĩa hàm ở cấp module thay cho lambda.
  • Cần “gắn sẵn” tham số? Dùng functools.partial(func, arg) - partial của hàm cấp module pickle được.
  • Tài nguyên như kết nối database: tạo trong tiến trình con, dùng initializer (xem bên dưới).

Hệ quả thứ hai: dữ liệu bị sao chép. Truyền một list 1 GB vào mỗi task nghĩa là pickle và gửi 1 GB mỗi lần.

chunksize: tham số bị bỏ quên đắt giá nhất

Phần tiêu đề “chunksize: tham số bị bỏ quên đắt giá nhất”

Hãy thử một tác vụ rất nhẹ - bình phương 200.000 số:

import time
from concurrent.futures import ProcessPoolExecutor
def square(x):
return x * x
if __name__ == "__main__":
data = range(200_000)
start = time.perf_counter()
[square(x) for x in data]
print(f"Tuần tự: {time.perf_counter() - start:.3f}s") # 0.012s
with ProcessPoolExecutor(4) as pool:
start = time.perf_counter()
list(pool.map(square, data))
print(f"chunksize=1: {time.perf_counter() - start:.3f}s") # 18.809s (!!!)
start = time.perf_counter()
list(pool.map(square, data, chunksize=5000))
print(f"chunksize=5000: {time.perf_counter() - start:.3f}s") # 0.036s

18.8 giây so với 0.012 giây - bản đa tiến trình chậm hơn 1500 lần! Với chunksize=1 (mặc định của ProcessPoolExecutor.map), mỗi phần tử là một lần gửi/nhận qua pipe. Chi phí giao tiếp lớn hơn công việc thật hàng nghìn lần.

chunksize=5000 gom 5000 phần tử thành một gói, giảm số lần giao tiếp từ 200.000 xuống 40. Dù vậy, nó vẫn chậm hơn bản tuần tự - vì bình phương một số quá rẻ, không đáng để song song hoá.

multiprocessing.Pool.map thì tự tính chunksize (khoảng len / (4 × số tiến trình)), nên cùng ví dụ này chỉ mất 0.06s.

Quy tắc thực hành:

  1. Chỉ song song hoá khi mỗi đơn vị công việc đủ nặng (tối thiểu vài mili-giây).
  2. Với nhiều phần tử nhẹ, luôn đặt chunksize khi dùng ProcessPoolExecutor.map, hoặc tự chia dữ liệu thành các lô.
  3. Truyền ít dữ liệu, trả về ít dữ liệu. Truyền đường dẫn file thay vì nội dung file; trả về kết quả tổng hợp thay vì dữ liệu thô.

initializer: chuẩn bị tài nguyên cho mỗi tiến trình

Phần tiêu đề “initializer: chuẩn bị tài nguyên cho mỗi tiến trình”

Giả sử mỗi task cần một model máy học nặng hoặc một kết nối database. Tạo lại mỗi task thì quá chậm, truyền qua pickle thì không được. Giải pháp: khởi tạo một lần cho mỗi tiến trình con:

import sqlite3
from concurrent.futures import ProcessPoolExecutor
_conn = None # biến toàn cục RIÊNG của từng tiến trình con
def init_worker(db_path):
global _conn
_conn = sqlite3.connect(db_path)
def lookup(user_id):
row = _conn.execute("SELECT ? * 2", (user_id,)).fetchone()
return row[0]
if __name__ == "__main__":
with ProcessPoolExecutor(
max_workers=4,
initializer=init_worker,
initargs=("app.db",),
) as pool:
print(list(pool.map(lookup, range(10), chunksize=3)))

Tham số liên quan: max_tasks_per_child (3.11+) cho phép tái tạo tiến trình con sau N task - hữu ích khi thư viện bạn dùng bị rò rỉ bộ nhớ.

from multiprocessing import Pool
counter = 0
def increment(_):
global counter
counter += 1
return counter
if __name__ == "__main__":
with Pool(4) as pool:
print(pool.map(increment, range(8))) # ví dụ: [1, 1, 1, 1, 2, 2, 2, 3]
print(counter) # 0 - tiến trình cha không hề thay đổi

Mỗi tiến trình có bản counter riêng. Kết quả phụ thuộc vào việc task nào rơi vào tiến trình nào. Muốn chia sẻ trạng thái, bạn phải dùng các công cụ dưới đây.

import multiprocessing as mp
def producer(q):
for i in range(5):
q.put(i * i)
q.put(None) # sentinel báo hết
def consumer(q):
while (item := q.get()) is not None:
print("nhận", item)
if __name__ == "__main__":
q = mp.Queue()
p1 = mp.Process(target=producer, args=(q,))
p2 = mp.Process(target=consumer, args=(q,))
p1.start(); p2.start()
p1.join(); p2.join()

mp.Queue giống queue.Queue nhưng hoạt động giữa các tiến trình (bên dưới là pipe + pickle + một thread nền). mp.Pipe() nhanh hơn nhưng chỉ nối hai đầu.

ValueArray - số và mảng trong bộ nhớ chia sẻ

Phần tiêu đề “Value và Array - số và mảng trong bộ nhớ chia sẻ”
import multiprocessing as mp
def add(counter, n):
for _ in range(n):
with counter.get_lock(): # Value có sẵn một khoá
counter.value += 1
if __name__ == "__main__":
counter = mp.Value("i", 0) # 'i' = int C, giống module array
procs = [mp.Process(target=add, args=(counter, 10_000)) for _ in range(4)]
for p in procs: p.start()
for p in procs: p.join()
print(counter.value) # 40000

Giống như thread, counter.value += 1 không nguyên tử - bỏ get_lock() sẽ mất dữ liệu.

import multiprocessing as mp
def record(shared, key):
shared[key] = key ** 2
if __name__ == "__main__":
with mp.Manager() as manager:
shared = manager.dict()
with mp.Pool(4) as pool:
pool.starmap(record, [(shared, i) for i in range(5)])
print(dict(shared)) # {0: 0, 1: 1, 2: 4, 3: 9, 4: 16}

Manager chạy một tiến trình máy chủ giữ object thật; các tiến trình khác chỉ giữ proxy, mỗi thao tác là một lời gọi qua mạng nội bộ. Tiện nhưng chậm - không dùng cho vòng lặp nóng.

shared_memory - chia sẻ vùng nhớ thật, không copy (3.8+)

Phần tiêu đề “shared_memory - chia sẻ vùng nhớ thật, không copy (3.8+)”

Khi cần nhiều tiến trình cùng làm việc trên một mảng số lớn (ảnh, ma trận), multiprocessing.shared_memory cấp một vùng nhớ mà mọi tiến trình ánh xạ trực tiếp. Kết hợp với NumPy qua buffer protocol:

import numpy as np
from multiprocessing import shared_memory
from concurrent.futures import ProcessPoolExecutor
SHAPE = (4, 1_000_000)
def process_row(args):
shm_name, row = args
shm = shared_memory.SharedMemory(name=shm_name) # gắn vào vùng nhớ đã có
try:
data = np.ndarray(SHAPE, dtype=np.float64, buffer=shm.buf)
data[row] *= 2 # sửa TRỰC TIẾP, không copy
result = float(data[row].sum())
del data # bỏ view trước khi close
return result
finally:
shm.close()
if __name__ == "__main__":
shm = shared_memory.SharedMemory(create=True, size=np.zeros(SHAPE).nbytes)
try:
data = np.ndarray(SHAPE, dtype=np.float64, buffer=shm.buf)
data[:] = 1.0
with ProcessPoolExecutor(4) as pool:
sums = list(pool.map(process_row, [(shm.name, r) for r in range(SHAPE[0])]))
print(sums) # [2000000.0, 2000000.0, 2000000.0, 2000000.0]
print(data[:, :3]) # mọi hàng đã được nhân đôi - tiến trình cha thấy ngay
del data
finally:
shm.close()
shm.unlink() # giải phóng vùng nhớ - CHỈ tiến trình tạo ra gọi

Chỉ có tên vùng nhớ (một chuỗi ngắn) được gửi qua pickle, 32 MB dữ liệu không bị sao chép lần nào. Nhớ:

  • Mỗi tiến trình gọi close() khi xong.
  • Đúng một tiến trình gọi unlink(), nếu không vùng nhớ bị rò rỉ tới khi khởi động lại máy (trên Linux).
  • Nếu nhiều tiến trình ghi cùng một vị trí, bạn vẫn cần khoá. Thiết kế tốt nhất là chia vùng: mỗi tiến trình một phần riêng, như ví dụ trên chia theo hàng.

ShareableList trong cùng module cho phép chia sẻ một list nhỏ các giá trị đơn giản.

  • Exception trong tiến trình con được pickle và ném lại khi bạn gọi future.result(), kèm traceback gốc dạng chuỗi.
  • Nếu một tiến trình con chết đột ngột (bị OOM killer, segfault trong extension C), ProcessPoolExecutor báo BrokenProcessPool và mọi task đang chờ đều thất bại. Hãy bắt lỗi này và tạo pool mới nếu cần.
  • Nhấn Ctrl+C khi dùng pool có thể để lại tiến trình “mồ côi”. Dùng with để pool được dọn dẹp, và cân nhắc pool.shutdown(cancel_futures=True) (3.9+).

Trước khi chuyển sang đa tiến trình, hãy tự hỏi:

  1. Có thể dùng thư viện đã song song hoá sẵn không? NumPy, Pandas, Polars, PyTorch thường đã chạy code C đa luồng. Vector hoá bằng NumPy thường nhanh hơn cả 8 tiến trình Python thuần.
  2. Công việc có đủ nặng không? Nếu mỗi task dưới 1ms, chi phí pickle sẽ nuốt hết lợi ích.
  3. Dữ liệu có lớn không? Nếu có, hãy để tiến trình con tự đọc từ file/database, hoặc dùng shared_memory.
  4. Bộ nhớ có đủ không? 8 tiến trình, mỗi cái nạp một model 2 GB = 16 GB RAM.
  1. Viết hàm count_primes(start, end) rồi đếm số nguyên tố dưới 5.000.000 bằng cách chia thành 8 đoạn chạy trên ProcessPoolExecutor. So sánh thời gian với bản tuần tự.
  2. Thử lại bài 1 với 1000 đoạn nhỏ, rồi với chunksize khác nhau. Giải thích kết quả.
  3. Viết chương trình tính kích thước trung bình của mọi file .py trong một thư mục lớn: truyền đường dẫn vào tiến trình con, trả về (số file, tổng byte), cộng dồn ở tiến trình cha.
  • Multiprocessing vượt qua GIL bằng cách cho mỗi tiến trình một interpreter và GIL riêng - phù hợp với tính toán CPU-bound bằng Python thuần.
  • Luôn bọc code khởi chạy trong if __name__ == "__main__": và viết code chạy đúng với spawn.
  • Mọi thứ trao đổi đều qua pickle → truyền ít, trả về ít, dùng hàm cấp module.
  • chunksize quyết định sống còn khi có nhiều task nhỏ.
  • Dùng initializer cho tài nguyên nặng, Queue cho thông điệp, Value/Array cho số đơn giản, shared_memory cho mảng lớn.

Bài tiếp theo: asyncio chuyên sâu - khi bạn cần xử lý hàng nghìn kết nối đồng thời.