Keyboard shortcuts

Press ← or → to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

Python: chọn async, thread hay process từ workload

Câu hỏi: chờ I/O và tính CPU có cùng hưởng lợi từ concurrency không, và pool startup tốn gì?

Cần biết trước: Python, async/await và process. Lab dùng CPython3.14.4, thư viện chuẩn, server HTTP chỉ nghe loopback với dữ liệu giả. Không cần database. Xác nhận build/GIL lúc chạy; không lấy kết quả build này để chứng nhận free-threaded.

Giữ lượng công việc và đầu ra giống nhau

I/O gồm8request mới tới server local, mỗi request chờ0.03s rồi trả số từ task ID. CPU gồm8tác vụ, mỗi tác vụ chạy600000vòng phép tính integer Python thuần. Thread/process giới hạn4worker; async I/O có semaphore4. CPU async không offload: coroutine chạy tính toán trước khi trả, không tự làm Python bytecode song song.

Timer tính toàn vòng create pool, execute, collect và shutdown, gồm IPC ở process. Server startup, expected-result và validation nằm ngoài timer. Chạy3lượt mỗi tổ hợp, luân phiên thứ tự; không ép cold cache. HTTP client blocking và async có implementation khác, vì vậy không quy toàn bộ chênh lệch cho scheduler.

Trên GIL build, nhiều thread không đồng thời chạy bytecode thuần trong một interpreter; I/O và extension có thể nhả GIL. Process có interpreter riêng nhưng phải chuyển input/ output, khởi tạo và quản lifecycle. GIL/build, executor.

Lưu concurrency.py trong thư mục trống

from __future__ import annotations

import asyncio
import json
import multiprocessing
import os
import platform
import statistics
import sys
import sysconfig
import threading
import time
from concurrent.futures import ProcessPoolExecutor, ThreadPoolExecutor
from datetime import UTC, datetime
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from pathlib import Path
from typing import TypedDict
from urllib.request import urlopen

MODES = ("seq", "async", "thread", "process")
IDS = tuple(range(8))


class Sample(TypedDict):
    work: str
    mode: str
    repeat: int
    seconds: float


class Handler(BaseHTTPRequestHandler):
    def do_GET(self) -> None:
        task_id = int(self.path.removeprefix("/"))
        time.sleep(0.2 if task_id == 999 else 0.03)
        body = str(task_id * 17).encode()
        self.send_response(200)
        self.send_header("Content-Length", str(len(body)))
        self.end_headers()
        try:
            self.wfile.write(body)
        except BrokenPipeError, ConnectionResetError:
            pass

    def log_message(self, format: str, *args: object) -> None:
        pass


def io_work(item: tuple[int, int]) -> int:
    port, task_id = item
    with urlopen(f"http://127.0.0.1:{port}/{task_id}", timeout=2) as response:
        return int(response.read())


def cpu_work(task_id: int) -> int:
    return sum((i * 17 + task_id) % 997 for i in range(600000))


async def async_io(port: int, task_id: int, limit: asyncio.Semaphore) -> int:
    async with limit:
        reader, writer = await asyncio.open_connection("127.0.0.1", port)
        try:
            writer.write(f"GET /{task_id} HTTP/1.0\r\nHost: localhost\r\n\r\n".encode())
            await writer.drain()
            async with asyncio.timeout(2):
                data = await reader.read()
            headers, body = data.split(b"\r\n\r\n", 1)
            assert headers.startswith(b"HTTP/1.0 200")
            return int(body)
        finally:
            writer.close()
            await writer.wait_closed()


async def async_cpu(task_id: int) -> int:
    return cpu_work(task_id)


async def async_run(work: str, port: int) -> list[int]:
    limit = asyncio.Semaphore(4)
    async with asyncio.TaskGroup() as group:
        tasks = [
            group.create_task(
                async_io(port, i, limit) if work == "io" else async_cpu(i)
            )
            for i in IDS
        ]
    return [task.result() for task in tasks]


def run(work: str, mode: str, port: int) -> list[int]:
    if mode == "seq":
        return [io_work((port, i)) if work == "io" else cpu_work(i) for i in IDS]
    if mode == "async":
        return asyncio.run(async_run(work, port))
    pool = (
        ThreadPoolExecutor(max_workers=4)
        if mode == "thread"
        else ProcessPoolExecutor(
            max_workers=4, mp_context=multiprocessing.get_context("spawn")
        )
    )
    with pool:
        if work == "io":
            return list(pool.map(io_work, ((port, i) for i in IDS)))
        return list(pool.map(cpu_work, IDS))


def ping(payload: bytes) -> tuple[int, int]:
    return len(payload), os.getpid()


def finite_work() -> int:
    time.sleep(0.15)
    return 1


async def cancellation(port: int) -> None:
    tasks: list[asyncio.Task[int]] = []
    try:
        async with asyncio.timeout(0.02):
            async with asyncio.TaskGroup() as group:
                tasks = [
                    group.create_task(async_io(port, 999, asyncio.Semaphore(4)))
                    for _ in range(2)
                ]
    except TimeoutError:
        assert len(tasks) == 2 and all(
            task.done() and task.cancelled() for task in tasks
        )
    else:
        raise AssertionError("timeout không xảy ra")
    assert asyncio.all_tasks() == {asyncio.current_task()}


def main() -> None:
    assert platform.python_version() == "3.14.4"
    gil = sys._is_gil_enabled()
    assert gil and not sysconfig.get_config_var("Py_GIL_DISABLED")
    server = ThreadingHTTPServer(("127.0.0.1", 0), Handler)
    server.daemon_threads = False
    server.block_on_close = True
    thread = threading.Thread(target=server.serve_forever)
    thread.start()
    port = server.server_port
    samples: list[Sample] = []
    ipc_samples = []
    try:
        expected = {"io": [i * 17 for i in IDS], "cpu": [cpu_work(i) for i in IDS]}
        for work in ("io", "cpu"):
            for repeat in range(3):
                for mode in MODES[repeat:] + MODES[:repeat]:
                    started = time.perf_counter()
                    result = run(work, mode, port)
                    elapsed = time.perf_counter() - started
                    assert result == expected[work]
                    samples.append(
                        {
                            "work": work,
                            "mode": mode,
                            "repeat": repeat + 1,
                            "seconds": elapsed,
                        }
                    )
        for size in (16, 1048576):
            payload = b"x" * size
            for repeat in range(3):
                started = time.perf_counter()
                with ProcessPoolExecutor(
                    max_workers=4, mp_context=multiprocessing.get_context("spawn")
                ) as pool:
                    ipc_result = list(pool.map(ping, [payload] * 8))
                elapsed = time.perf_counter() - started
                assert all(length == size for length, _ in ipc_result)
                ipc_samples.append(
                    {
                        "bytes_per_task": size,
                        "repeat": repeat + 1,
                        "seconds": elapsed,
                        "workers_observed": len({pid for _, pid in ipc_result}),
                    }
                )
        asyncio.run(cancellation(port))
        with ProcessPoolExecutor(
            max_workers=1, mp_context=multiprocessing.get_context("spawn")
        ) as pool:
            future = pool.submit(finite_work)
            try:
                future.result(timeout=0.01)
            except TimeoutError:
                cancelled = future.cancel()
            else:
                raise AssertionError("future timeout không xảy ra")
        assert future.done() and (future.cancelled() or future.result() == 1)
        assert not multiprocessing.active_children()
        summaries = []
        for work in ("io", "cpu"):
            for mode in MODES:
                times = [
                    s["seconds"]
                    for s in samples
                    if s["work"] == work and s["mode"] == mode
                ]
                assert len(times) == 3
                summaries.append(
                    {
                        "work": work,
                        "mode": mode,
                        "mean_s": statistics.mean(times),
                        "min_s": min(times),
                        "max_s": max(times),
                        "stdev_s": statistics.stdev(times),
                    }
                )
        metadata = {
            "checked_at": datetime.now(UTC).isoformat(),
            "python": platform.python_version(),
            "platform": platform.platform(),
            "logical_cpus": os.cpu_count(),
            "gil_enabled": gil,
            "gil_disabled_build": sysconfig.get_config_var("Py_GIL_DISABLED"),
            "start_method": "spawn",
            "tasks": 8,
            "workers_max": 4,
            "cpu_iterations": 600000,
            "io_delay_s": 0.03,
            "cache": "no eviction; rotated mode order",
            "timer": "create/execute/collect/shutdown including IPC",
            "future_cancel_succeeded": cancelled,
        }
        payload_json = {
            "metadata": metadata,
            "samples": samples,
            "summaries": summaries,
            "ipc_samples": ipc_samples,
        }
        Path("concurrency-results.json").write_text(
            json.dumps(payload_json, indent=2) + "\n"
        )
        sys.stdout.write("CONC_RESULT " + json.dumps(payload_json) + "\n")
    finally:
        server.shutdown()
        server.server_close()
        thread.join(timeout=2)
        assert not thread.is_alive() and not multiprocessing.active_children()
    sys.stdout.write(
        "samples=24 outputs=equal timeout/cancel=ok server/tasks/processes=closed\n"
    )


if __name__ == "__main__":
    main()
set -euo pipefail
python3 -B concurrency.py
samples=24 outputs=equal timeout/cancel=ok server/tasks/processes=closed
set -euo pipefail
test -f concurrency-results.json

Kết quả và startup/IPC

Đo2026-10-03 lúc07:59:28UTC, macOS27.0.1arm64/10logicalCPU, Python3.14.4, GILenabled=true, Py_GIL_DISABLED=0, process dùngspawn. Thời gian giây, mỗi dòng3sample; không evict cache, thứ tự cách chạy luân phiên.

WorkloadCáchMean (s)Min (s)Max (s)Stddev (s)
ioseq0.3199600.3167770.3239470.003652
ioasync0.0887060.0851600.0928300.003868
iothread0.0778330.0710860.0813320.005844
ioprocess0.1686880.1569170.1817910.012491
cpuseq0.1665120.1574630.1830310.014328
cpuasync0.1543480.1525830.1578630.003044
cputhread0.1523410.1520200.1526380.000310
cpuprocess0.1182920.1168020.1204410.001907

I/O thread/async thấp hơn seq trong lần đo này, còn process trả thêm startup. CPU process có mean thấp hơn ba cách còn lại ở lượng việc này; vẫn gồm create/ IPC/shutdown. CPU thread/async gần thời gian seq, không chứng minh bytecode chạy parallel: CPU async không có điểm await, và build đang bật GIL. Biến thiên và khác biệt overhead có thể làm mean nhích; cần profile nếu muốn quy nguyên nhân.

Ping16B có ba thời gian0.073312/0.073714/0.071822s;1MiB là 0.073671/0.073873/0.075469s. Ping nhỏ quan sát2–4PID thực sự nhận việc, ping lớn4PID; max_workers4 không có nghĩa mọi worker nhận cùng số task. Ping16B/1MiB đo create pool + input serialization + IPC + length computation + collect/shutdown; không phải chi phí IPC thuần. Khoảng thời gian chồng nhau; không trừ hai mean để gán một chi phí transfer chắc chắn. JSON giữ24sample,8summary và6ping sample để đọc lại từng lượt; lần chạy lại có thể khác số. Tất cả kết quả workload bằng expected, timeout/cancel/cleanup đạt.

Timeout không tự giết công việc đang chạy

TaskGroup gom task và await cleanup. asyncio.timeout hủy ở điểm coroutine có thể nhường quyền; finally đóng connection và CancelledError được truyền tiếp. CPU thuần không yield có thể giữ loop qua deadline; đổi thành asyncdef không giải quyết. Task/cancellation.

Future timeout chỉ dừng chờ; cancel có thể thất bại nếu đã running. Lab ghi kết quả cancel, chờ công việc hữu hạn0.15s hoàn thành khi đóng pool, rồi kiểm active_children rỗng. Một công việc treo vô hạn cần policy khác; không gọi timeout là kill. Server vẫn có thể xử lý request client đã hủy; fixture hữu hạn và server_close chờ handlers.

Free-threaded là build/runtime khác; extension có thể làm GIL bật lại. Kiểm sys._is_gil_enabled() bên cạnh build flag rồi benchmark trên build ấy. sys runtime. Chọn từ thời gian và lượng công việc thật, chi phí quản pool và giới hạn tài nguyên; không chọn process chỉ vì có chữ CPU trong tên task.