Работа с потоками в Python: threading на практике

Работа с потоками в Python: threading на практике Полезное

Модуль threading — это стандартный модуль Python для запуска нескольких потоков выполнения внутри одного процесса. Потоки делят общую память процесса, поэтому их удобно использовать там, где программа в основном ждет: сеть, диск, внешние сервисы.

Здесь разберем прикладную сторону: как запустить пачку задач через пул потоков и собрать результаты и ошибки, как передавать данные через очередь, как аккуратно остановить фоновый поток, как не попасть во взаимную блокировку и когда вместо потоков брать процессы или asyncio. Примеры проверены на Python 3.14.7 (официальный образ python:3.14-slim) 24.09.2026.

Три термина, которые нужно развести сразу

  • Поток (thread) — единица выполнения внутри процесса. Все потоки процесса видят одни и те же объекты.
  • Процесс — отдельный экземпляр программы со своей памятью. Данные между процессами передаются явно (очереди, каналы, файлы).
  • GIL — глобальная блокировка интерпретатора в обычной сборке CPython: в каждый момент байткод Python исполняет только один поток процесса. На время ожидания ввода-вывода поток отпускает GIL.

Отсюда практическое правило с границей: в обычной сборке CPython потоки ускоряют задачи, которые ждут (I/O-bound), и почти не ускоряют чистые вычисления на Python (CPU-bound). Подробно механика GIL и замер на вычислениях — в статье Потоки в Python: threading, GIL и multiprocessing. Исключение — free-threaded сборка без GIL (о ней ниже).

Пул потоков: основной способ запускать задачи

В прикладном коде потоки редко создают вручную по одному. Обычно берут ThreadPoolExecutor из concurrent.futures: он держит ограниченное число рабочих потоков, раздает им задачи и возвращает результаты. Минимальный полный пример — три «запроса», каждый ждет 1 секунду:

import time
from concurrent.futures import ThreadPoolExecutor

def fetch(url):
    time.sleep(1)  # имитация сетевого ожидания
    return f"{url}: 200"

urls = ["https://example.com/a", "https://example.com/b", "https://example.com/c"]

start = time.perf_counter()
with ThreadPoolExecutor(max_workers=3) as pool:
    for result in pool.map(fetch, urls):
        print(result)
print(f"Итого: {time.perf_counter() - start:.1f} с")

Вывод прогона:

https://example.com/a: 200
https://example.com/b: 200
https://example.com/c: 200
Итого: 1.0 с

Последовательно три ожидания заняли бы около 3 секунд, в пуле — около 1 секунды, потому что потоки ждут одновременно. Разбор по строкам:

  • with ThreadPoolExecutor(...) — при выходе из блока пул дожидается всех задач и закрывает потоки, отдельный join() не нужен.
  • max_workers — верхняя граница числа потоков. Для сетевых запросов ее выбирают по допустимой нагрузке на сервис, а не по числу ядер.
  • pool.map возвращает результаты в порядке входного списка, даже если задачи завершились в другом порядке.

Ошибки в задачах: submit и as_completed

Если задача в потоке упала, исключение не теряется: future.result() поднимет его в вызывающем коде. as_completed отдает задачи по мере завершения:

from concurrent.futures import ThreadPoolExecutor, as_completed

def parse(n):
    if n == 3:
        raise ValueError(f"плохие данные в задаче {n}")
    return n * 10

with ThreadPoolExecutor(max_workers=4) as pool:
    futures = {pool.submit(parse, n): n for n in range(1, 6)}
    results = {}
    for fut in as_completed(futures):
        n = futures[fut]
        try:
            results[n] = fut.result()
        except ValueError as exc:
            print("Ошибка:", exc)

print(dict(sorted(results.items())))
Ошибка: плохие данные в задаче 3
{1: 10, 2: 20, 4: 40, 5: 50}

Одна упавшая задача не ломает остальные. С pool.map поведение другое: исключение поднимется при получении результата упавшей задачи, и перебор оборвется на ней. Если нужны все результаты и все ошибки — используйте submit + as_completed.

Очередь Queue: передача данных между потоками

Когда одни потоки производят данные, а другие обрабатывают, общий список с ручными блокировками не нужен: queue.Queue уже потокобезопасна. Схема «производитель — потребители» с маркером остановки:

import threading
import queue

tasks = queue.Queue(maxsize=10)
results = queue.Queue()
STOP = object()  # маркер конца работы

def worker():
    while True:
        item = tasks.get()
        try:
            if item is STOP:
                return
            results.put(item * item)
        finally:
            tasks.task_done()

workers = [threading.Thread(target=worker) for _ in range(3)]
for w in workers:
    w.start()

for n in range(1, 7):
    tasks.put(n)          # блокируется, если очередь заполнена
for _ in workers:
    tasks.put(STOP)       # по одному маркеру на каждого рабочего

tasks.join()              # ждем, пока все элементы обработаны
for w in workers:
    w.join()

print(sorted(results.get() for _ in range(results.qsize())))
[1, 4, 9, 16, 25, 36]

Что здесь важно:

  • maxsize=10 ограничивает очередь: быстрый производитель будет ждать, а не заполнит память.
  • task_done() вызывается в finally на каждый get(), иначе tasks.join() не дождется завершения.
  • Маркер STOP кладется по одному на каждый поток: один маркер завершит только одного рабочего.
  • qsize() — приблизительный размер. Здесь он точен только потому, что все потоки уже завершены.

Как остановить поток

Принудительно «убить» поток из другого потока в threading нельзя. Поток завершается сам, когда его функция возвращает управление, поэтому остановку нужно заложить в код. Для этого удобнее threading.Event, чем глобальный флаг: Event.wait(timeout) одновременно служит паузой и просыпается сразу по сигналу.

import threading
import time

stop = threading.Event()

def heartbeat():
    n = 0
    while not stop.is_set():
        n += 1
        print("тик", n)
        stop.wait(0.3)   # спит до 0.3 с, но просыпается сразу по stop.set()
    print("поток завершился аккуратно")

t = threading.Thread(target=heartbeat)
t.start()
time.sleep(1)
stop.set()
t.join()
тик 1
тик 2
тик 3
тик 4
поток завершился аккуратно

Число тиков зависит от таймингов, в нашем прогоне их было 4. С time.sleep(0.3) поток реагировал бы на остановку с задержкой до конца паузы; с stop.wait — сразу.

Демон-поток (daemon=True) тоже не ждет остановки, но по-другому: он обрывается при выходе интерпретатора в произвольной точке, без finally и закрытия файлов. Поэтому демон подходит для фоновых задач без важного состояния, а запись в файл или базу лучше останавливать через Event и join().

Lock, RLock и взаимная блокировка

Lock дает одному потоку монопольный доступ к общему ресурсу. Захват через with lock: гарантирует освобождение даже при исключении. Типовая ловушка — повторный захват того же Lock тем же потоком, например когда одна функция под блокировкой вызывает другую, которая тоже берет блокировку. Обычный acquire() в такой ситуации ждал бы вечно. Чтобы пример не завис, покажем это через timeout:

import threading

lock = threading.Lock()

def outer():
    with lock:
        inner()

def inner():
    ok = lock.acquire(timeout=1)   # тот же поток просит тот же Lock
    print("inner получил Lock:", ok)
    if ok:
        lock.release()

outer()
inner получил Lock: False

Через секунду acquire вернул False: Lock не помнит владельца, и поток заблокировал сам себя. Исправление для вложенных вызовов — RLock, который владелец может захватывать повторно (освобождать нужно столько же раз, with делает это сам):

import threading

lock = threading.RLock()

def outer():
    with lock:
        inner()

def inner():
    with lock:        # владелец может войти повторно
        print("inner внутри RLock")

outer()
inner внутри RLock

RLock решает только самоблокировку одного потока. Взаимная блокировка двух потоков (первый держит A и ждет B, второй держит B и ждет A) им не лечится. Правило для этого случая: все потоки берут несколько блокировок в одном и том же порядке, а в спорных местах используют acquire(timeout=...) с обработкой отказа.

Еще одна частая ошибка — освободить блокировку, которую никто не захватил:

import threading

lock = threading.Lock()
lock.release()
RuntimeError: release unlocked lock

Исправление — не вызывать acquire/release вручную там, где подходит with lock:.

Semaphore: ограничить число одновременных операций

Semaphore(n) пропускает внутрь не больше n потоков одновременно. Это удобно, когда потоков много, а внешний ресурс (API, база, диск) выдерживает ограниченную параллельность:

import threading
import time

limit = threading.Semaphore(2)   # не больше двух одновременно
active = 0
peak = 0
guard = threading.Lock()

def job():
    global active, peak
    with limit:
        with guard:
            active += 1
            peak = max(peak, active)
        time.sleep(0.2)
        with guard:
            active -= 1

threads = [threading.Thread(target=job) for _ in range(6)]
for t in threads:
    t.start()
for t in threads:
    t.join()
print("максимум одновременно:", peak)
максимум одновременно: 2

Шесть потоков запущены, но в защищенный участок одновременно попадали не больше двух. Счетчики active и peak меняются под отдельным Lock: сам семафор не защищает общие переменные внутри участка.

Потоки, процессы или asyncio: как выбрать

Задача Что взять Почему Когда иначе
Много сетевых запросов, чтение файлов, вызовы API на синхронных библиотеках ThreadPoolExecutor / threading Ожидание идет параллельно, код остается обычным синхронным Тысячи одновременных соединений — смотреть на asyncio
Тысячи соединений, код на async-библиотеках asyncio Один поток переключает задачи в точках await, дешевле тысяч потоков Нужна блокирующая библиотека — вынести ее в asyncio.to_thread
Вычисления на чистом Python (парсинг, расчеты) ProcessPoolExecutor / multiprocessing Каждый процесс со своим интерпретатором и GIL, процессы ОС распределяет по разным ядрам Данные большие и дорого передаются между процессами — сначала профилировать
Вычисления в NumPy, сжатие, хеширование больших буферов Часто подходят и потоки Многие C-расширения отпускают GIL на время работы Зависит от конкретной библиотеки и функции — измерять
Фоновая периодическая задача в приложении threading.Thread + Event Простая остановка и join() В async-приложении — задача asyncio

Отдельная оговорка про free-threaded сборку CPython (вариант интерпретатора без GIL, в Python 3.14 обозначается python3.14t). С Python 3.14 она официально поддерживается (PEP 779) и больше не считается экспериментальной. В ней потоки Python могут выполнять байткод параллельно, но это не сборка по умолчанию, и не все C-расширения ее поддерживают. Проверить, включен ли GIL в текущем интерпретаторе, можно через sys._is_gil_enabled() (доступно с Python 3.13).

Выводы

  • threading запускает потоки в одном процессе с общей памятью. В обычной сборке CPython они выигрывают на ожидании (сеть, диск), а не на вычислениях на Python.
  • Для прикладных задач основной инструмент — ThreadPoolExecutor: map для результатов по порядку, submit + as_completed для сбора ошибок без обрыва перебора.
  • Данные между потоками удобнее передавать через queue.Queue с maxsize, task_done() и маркером остановки на каждый рабочий поток.
  • Поток нельзя остановить извне принудительно: закладывайте остановку через Event, демон-потоки оставляйте для задач без важного состояния.
  • RLock лечит повторный захват в одном потоке, но не взаимную блокировку двух потоков: там нужен единый порядок захвата и timeout.

Где применяется / связь с практикой

Освойте тему на практике

Пулы потоков и очереди встречаются в любом бэкенде и автоматизации: параллельные запросы к API, загрузка и обработка файлов, фоновые задачи в веб-сервисе, парсеры. На собеседованиях Python-разработчика обязательно спрашивают про GIL, выбор между потоками, процессами и asyncio и про гонки данных. Эти темы вместе с асинхронным программированием и профилированием системно разбирают на курсе «Python-разработчик. Продвинутый уровень». Попробовать формат можно на бесплатных открытых уроках.

FAQ

Сколько потоков ставить в max_workers?
Для I/O-задач число определяет допустимая нагрузка на внешний ресурс и лимиты сервиса, а не число ядер; для начала часто берут 5-20 и подбирают замером. Если параметр не указан, Python выбирает значение сам, и формула менялась между версиями.

Потокобезопасны ли list и dict в Python?
Отдельные операции вроде append в обычной сборке CPython выполняются без повреждения структуры, но составные действия («проверить, потом изменить», x += 1) атомарными не являются. Для общих данных надежнее Lock или queue.Queue.

Можно ли получить значение, которое вернула функция потока, из threading.Thread?
Напрямую нет: Thread не хранит возвращаемое значение. Результат передают через очередь или общий объект под блокировкой, либо сразу используют ThreadPoolExecutor, где future.result() возвращает значение.

OTUS Журнал
Бесплатные открытые уроки (поп-ап)