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,forkservervà vì sao cầnif __name__ == "__main__" - Dùng
ProcessPoolExecutorvàmultiprocessing.Poolhiệ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_memoryzero-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.
Thread vs Process: khác nhau ở đâu?
Phần tiêu đề “Thread vs Process: khác nhau ở đâu?” 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 timefrom 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.78sNhanh 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:
fork (mặc định cũ trên Linux)
Phần tiêu đề “fork (mặc định cũ trên Linux)”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.
spawn (mặc định trên Windows và macOS)
Phần tiêu đề “spawn (mặc định trên Windows và 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__ là "__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(...)Mọi thứ đi qua pickle
Phần tiêu đề “Mọi thứ đi qua pickle”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 conpickle(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 timefrom 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.036s18.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:
- Chỉ song song hoá khi mỗi đơn vị công việc đủ nặng (tối thiểu vài mili-giây).
- Với nhiều phần tử nhẹ, luôn đặt
chunksizekhi dùngProcessPoolExecutor.map, hoặc tự chia dữ liệu thành các lô. - 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 sqlite3from 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ớ.
Biến toàn cục KHÔNG được chia sẻ
Phần tiêu đề “Biến toàn cục KHÔNG được chia sẻ”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 đổiMỗ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.
Giao tiếp và chia sẻ dữ liệu
Phần tiêu đề “Giao tiếp và chia sẻ dữ liệu”Queue và Pipe - gửi thông điệp
Phần tiêu đề “Queue và Pipe - gửi thông điệp”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.
Value và Array - 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) # 40000Giống như thread, counter.value += 1 không nguyên tử - bỏ get_lock() sẽ mất dữ liệu.
Manager - object Python chia sẻ qua proxy
Phần tiêu đề “Manager - object Python chia sẻ qua proxy”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 npfrom multiprocessing import shared_memoryfrom 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ọiChỉ 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.
Xử lý lỗi và dừng chương trình
Phần tiêu đề “Xử lý lỗi và dừng chương trình”- 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),
ProcessPoolExecutorbáoBrokenProcessPoolvà 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ắcpool.shutdown(cancel_futures=True)(3.9+).
Cân nhắc trước khi dùng multiprocessing
Phần tiêu đề “Cân nhắc trước khi dùng multiprocessing”Trước khi chuyển sang đa tiến trình, hãy tự hỏi:
- 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.
- 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.
- 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. - Bộ nhớ có đủ không? 8 tiến trình, mỗi cái nạp một model 2 GB = 16 GB RAM.
Bài tập
Phần tiêu đề “Bài tập”- 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ênProcessPoolExecutor. So sánh thời gian với bản tuần tự. - Thử lại bài 1 với 1000 đoạn nhỏ, rồi với
chunksizekhác nhau. Giải thích kết quả. - Viết chương trình tính kích thước trung bình của mọi file
.pytrong 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.
Kết luận
Phần tiêu đề “Kết luận”- 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ớispawn. - Mọi thứ trao đổi đều qua pickle → truyền ít, trả về ít, dùng hàm cấp module.
chunksizequyết định sống còn khi có nhiều task nhỏ.- Dùng
initializercho tài nguyên nặng,Queuecho thông điệp,Value/Arraycho số đơn giản,shared_memorycho 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.