Потоки в Python: threading, start, join, Lock и Event

Python Автор: Среда и версия: CPython 3.14.5; Linux

Поток в Python создаётся классом threading.Thread: функция уходит в target, её аргументы кортежем в args, дальше start() и join(). Потоки ускоряют работу, которая простаивает на вводе-выводе: запросы по сети, чтение файлов, ожидание ответа базы. Вычисления в чистом Python они не ускоряют, мешает GIL.

Весь разбор идёт на одной задаче: скачать несколько страниц. Настоящая сеть здесь не нужна, задержку изображает time.sleep.

import threading
import time

def fetch(path):
    time.sleep(0.25)
    print(f"загрузил {path}")
    return f"<html>{path}</html>"

PAGES = ["/pricing", "/docs", "/blog", "/about"]

Всё, что ниже, работает на CPython 3.14.5 под Linux. Выводы в статье получены запуском, а не написаны от руки.

Как создать и запустить поток

Три вызова: конструктор, start(), join().

t = threading.Thread(target=fetch, args=("/pricing",))
t.start()
print("главный поток не ждёт")
t.join()
print("join вернулся")
главный поток не ждёт
загрузил /pricing
join вернулся

start() запускает поток и сразу возвращает управление, поэтому «главный поток не ждёт» печатается первым. join() блокирует вызывающий поток, пока целевой не закончит. Без join() главный поток дойдёт до конца программы и будет ждать дочерние уже на выходе из интерпретатора, но результата ты к тому моменту не получишь.

Второй раз запустить тот же объект нельзя: t.start() после завершения бросит RuntimeError: threads can only be started once. Нужен ещё один запуск, создавай новый Thread.

Вызвать t.run() вместо t.start() можно, и это частая опечатка. run() просто выполнит функцию в текущем потоке: новый поток не создастся, t.is_alive() останется False, параллельности не будет.

Почему args обязан быть кортежем

args разворачивается в позиционные аргументы через *args. Строка тоже итерируемая, поэтому args=("/pricing") без запятой распадается на символы:

t = threading.Thread(target=fetch, args=("/pricing"))
t.start()
Exception in thread Thread-1 (fetch):
Traceback (most recent call last):
  ...
TypeError: fetch() takes 1 positional argument but 8 were given

Восемь аргументов — это восемь символов строки /pricing. Запятая обязательна: args=("/pricing",). Список тоже подойдёт (args=["/pricing"]), а именованные аргументы передаются отдельным словарём: kwargs={"path": "/pricing"}.

TypeError при этом вылетел внутри потока, поэтому программа не упала и вернула код 0. Почему так, разбираем ниже.

Как узнать имя текущего потока

threading.current_thread() возвращает объект потока, из которого ты вызвал функцию. У главного потока имя MainThread, до него же можно добраться через threading.main_thread().

def fetch(path):
    time.sleep(0.25)
    print(f"{threading.current_thread().name}: {path}")
    return f"<html>{path}</html>"

print("стартуем из", threading.current_thread().name)

workers = [threading.Thread(target=fetch, args=(p,), name=f"loader{i}")
           for i, p in enumerate(PAGES)]
for w in workers:
    w.start()
for w in workers:
    w.join()

print("осталось потоков:", threading.active_count())
стартуем из MainThread
loader0: /pricing
loader1: /docs
loader2: /blog
loader3: /about
осталось потоков: 1

Порядок строк совпал с порядком старта, но полагаться на это нельзя. На сотне прогонов он сломался тринадцать раз, чаще всего loader3 обгонял loader2. Планировщик ничего не обещает.

Имя задаётся параметром name. Не задал, и Python соберёт его сам из счётчика и имени целевой функции: Thread-1 (fetch). Имя нужно только людям. На поведение оно не влияет и уникальным быть не обязано, а машина различает потоки по t.ident и t.native_id (последний — тот же TID, что показывает htop).

active_count() в конце равен единице: остался один MainThread, четыре загрузчика отработали и умерли.

Чем daemon-поток отличается от обычного

Интерпретатор на выходе ждёт только не-daemon потоки. Daemon-поток при выходе просто прекращает существовать.

def flush_logs():
    try:
        while True:
            print("сбрасываю буфер")
            time.sleep(0.2)
    finally:
        print("finally: закрываю файл")

t = threading.Thread(target=flush_logs, daemon=True)
t.start()
time.sleep(0.5)
print("главный поток закончил")
сбрасываю буфер
сбрасываю буфер
сбрасываю буфер
главный поток закончил

Бесконечный цикл не помешал программе завершиться, и код возврата остался нулевым. Но finally не выполнился. Daemon-поток не разматывает стек и не доводит до конца with-блоки: файл не закроется, транзакция не откатится, буфер не долетит до диска.

Отсюда правило. Фоновые задачи без побочных эффектов (heartbeat, сбор метрик, прогрев кэша) — daemon. Всё, что пишет данные, запускай обычным потоком и на выходе дожидайся его явным join().

daemon меняется только до start(); на живом потоке присвоение бросит RuntimeError. Значение по умолчанию наследуется от родителя, поэтому поток, запущенный из daemon-потока, тоже daemon.

Почему исключение внутри потока не роняет программу

Исключение живёт в стеке своего потока и в главный не всплывает. try/except вокруг start() и join() его не поймает.

def fetch(path):
    time.sleep(0.1)
    raise ValueError(f"404 на {path}")

try:
    t = threading.Thread(target=fetch, args=("/docs",))
    t.start()
    t.join()
except ValueError:
    print("сюда мы не попадём")

print("главный поток продолжает работу")
Exception in thread Thread-1 (fetch):
Traceback (most recent call last):
  ...
  File "pages.py", line 6, in fetch
    raise ValueError(f"404 на {path}")
ValueError: 404 на /docs
главный поток продолжает работу

Трейсбек напечатан, ветка except пропущена, программа завершилась с кодом 0. В продакшене это выглядит так: воркер тихо умер, очередь не разбирается, мониторинг зелёный.

Ловить исключения потоков глобально умеет threading.excepthook (появился в Python 3.8). Он получает объект с полями thread, exc_type, exc_value, exc_traceback:

errors = []

def on_thread_error(args):
    errors.append((args.thread.name, args.exc_value))

threading.excepthook = on_thread_error

t = threading.Thread(target=fetch, args=("/docs",), name="loader1")
t.start()
t.join()
print(errors)
[('loader1', ValueError('404 на /docs'))]

Хук глобальный и годится для логирования. Если исключение нужно получить в вызывающем коде, бери concurrent.futures.ThreadPoolExecutor: он складывает результат в Future, и future.result() перебрасывает исключение туда, где ты его ждёшь.

Зачем нужен threading.local

Объект threading.local() выглядит как обычный объект с атрибутами, но каждый поток видит собственный набор значений. Записал state.session в одном потоке — в другом этого атрибута нет.

Типовой повод: ресурс, который нельзя делить между потоками, а таскать параметром через десять вызовов не хочется. HTTP-сессия, курсор к базе, буфер парсера.

Четыре потока разбирают 40 страниц, PAGES = [f"/page/{i}" for i in range(40)]:

state = threading.local()
sessions = []
lock = threading.Lock()

def fetch(path):
    if not hasattr(state, "session"):
        state.session = f"session-{threading.current_thread().name}"
        with lock:
            sessions.append(state.session)
    time.sleep(0.01)
    return f"<html>{path}</html>"

def worker(paths):
    for path in paths:
        fetch(path)

threads = [threading.Thread(target=worker, args=(PAGES[i::4],), name=f"loader{i}")
           for i in range(4)]
for t in threads:
    t.start()
for t in threads:
    t.join()

print(f"страниц: {len(PAGES)}, сессий открыто: {len(sessions)}")
print(sorted(sessions))
страниц: 40, сессий открыто: 4
['session-loader0', 'session-loader1', 'session-loader2', 'session-loader3']

Сорок вызовов fetch, четыре сессии. С обычной глобальной переменной сессия открылась бы одна на всех, и потоки затирали бы друг другу состояние. Хранить в local() то, что потом нужно собрать по всем потокам, смысла нет: это ровно тот случай, когда данные общие и нужен Lock.

Почему счётчик теряет значения

downloaded += len(...) выглядит как одна операция, а на самом деле их три: прочитать текущее значение, посчитать новое, записать обратно. Если между чтением и записью успел поработать другой поток, его результат затрётся.

Те же 40 страниц и четыре потока, fetch со sleep(0.01). Загрузчики складывают размер каждой страницы в общий счётчик:

PAGES = [f"/page/{i}" for i in range(40)]
downloaded = 0

def worker(paths):
    global downloaded
    for path in paths:
        downloaded += len(fetch(path))

threads = [threading.Thread(target=worker, args=(PAGES[i::4],)) for i in range(4)]
for t in threads:
    t.start()
for t in threads:
    t.join()

print("ожидали 830, получили", downloaded)
ожидали 830, получили 208

208 вместо 830. На пятидесяти прогонах выпадало то 208, то 207, и это ровно суммы отдельных потоков: 207, 207, 208, 208. Потеряно три четверти данных: каждый поток прочитал downloaded до того, как соседи записали своё, и уцелела работа одного потока из четырёх. Какого именно, решает планировщик.

Разгадка в порядке вычислений. Python сначала кладёт старое значение downloaded на стек, потом идёт вызывать fetch(path), и вот на этом вызове поток засыпает на time.sleep и отдаёт GIL другим. Все четверо просыпаются со своим старым значением на руках и записывают его поверх чужого.

Теперь то же самое, но вызов вынесен из выражения:

body = fetch(path)
downloaded += len(body)
ожидали 830, получили 830

Сходится, и так на всех пятидесяти прогонах. Между чтением и записью не осталось ничего, что отдаёт GIL, и переключение в этот зазор не попадает. Считать такой код корректным нельзя: гонка никуда не делась, просто окно стало слишком узким, чтобы в него попасть. Ширина окна зависит от того, как собран конкретный CPython, и в сборке без GIL (в 3.14 она уже официально поддерживается) вся защита исчезает.

Единственный надёжный ответ — блокировка. Lock пропускает в защищённый участок ровно один поток:

lock = threading.Lock()

def worker(paths):
    global downloaded
    for path in paths:
        body = fetch(path)
        with lock:
            downloaded += len(body)
ожидали 830, получили 830

Под lock держи только чтение и запись общей переменной. Загрузку внутрь with не затаскивай: пока один поток спит в fetch, остальные будут стоять в очереди, и весь смысл потоков пропадёт.

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

Убить поток снаружи нельзя, метода stop() у Thread нет и не будет: поток может держать блокировку или недописанный файл. Останавливается поток только сам, увидев сигнал.

Флаг в виде bool работает, но заставляет крутить цикл с time.sleep и реагировать с задержкой. threading.Event решает обе задачи сразу: wait(timeout) спит до таймаута и просыпается мгновенно, как только флаг взведён.

stop = threading.Event()

def poll_queue():
    n = 0
    while not stop.is_set():
        n += 1
        print(f"проверка очереди #{n}")
        stop.wait(0.2)
    print("поток вышел из цикла сам")

t = threading.Thread(target=poll_queue)
t.start()
time.sleep(0.5)
stop.set()
t.join()
print("join вернулся")
проверка очереди #1
проверка очереди #2
проверка очереди #3
поток вышел из цикла сам
join вернулся

stop.set() взводит флаг, поток досматривает итерацию и выходит, join() дожидается выхода. Если поток заблокирован на чтении сокета, Event его не разбудит: там нужен таймаут на самой операции ввода-вывода.

Когда потоки ускоряют, а когда нет

Восемь страниц по 0.25 секунды, каждая в своём потоке:

последовательно: 2.00 c
в 8 потоков:      0.25 c

Восьмикратно, потому что всё это время потоки спали. Ждущий поток GIL не держит, и ожидания складываются друг с другом вместо того, чтобы выстраиваться в очередь.

Теперь вместо сна восемь раз считается вот это, n = 5_000_000. Машина двадцатиядерная:

def checksum(n):
    total = 0
    for i in range(n):
        total += i * i
    return total
последовательно: 1.28 c
в 8 потоков:      1.37 c

Не быстрее, а чуть медленнее: добавились накладные расходы на переключение. Байткод CPython в сборке с GIL выполняет по одному потоку за раз, независимо от числа ядер, а такую сборку ты и получаешь по умолчанию. Механизм разобран в статье про GIL.

Практический вывод короткий. Ждёшь сеть, диск или базу, бери потоки. Считаешь в чистом Python, бери multiprocessing или библиотеку, которая на время расчёта отпускает GIL (NumPy так делает). Если ждущих соединений тысячи, потоки станут дороже полезной работы, и туда лучше идёт asyncio. Развёрнутое сравнение трёх подходов — в разборе потоки, процессы или asyncio.

На чём ловят на собеседовании

«Обернул join() в try/except, значит ошибки поймаю». Не поймаешь. Исключение остаётся в своём потоке, трейсбек уходит в stderr, код возврата нулевой. Собирай ошибки через threading.excepthook или бери ThreadPoolExecutor с future.result().

join(timeout=...) не останавливает поток. Таймаут ограничивает ожидание, а не работу. После выхода по таймауту поток продолжает крутиться, t.is_alive() вернёт True, и узнать это можно только так: возвращаемого значения у join() нет вовсе.

«GIL защищает от гонок». GIL защищает внутренние структуры интерпретатора, а не твою логику. Любое «прочитал, изменил, записал» через две строки кода состоит из нескольких шагов, и между ними поток могут прервать. Пример с 207–208 вместо 830 выше.

Счётчик, который «работает». Тот же += в другом порядке вычислений выдал 830 на всех пятидесяти прогонах. Проверка запуском тут ничего не доказывает: гонка воспроизводится не всегда, и на проде она проявится в другой момент.

t.run() вместо t.start(). Функция выполнится в текущем потоке. Программа отработает и даст верный результат, просто последовательно и без единого нового потока.

Частые вопросы

Как получить результат из потока

Thread возвращаемое значение выбрасывает. Либо клади результат в общую структуру (list.append и dict[key] = value потокобезопасны, отдельный Lock под них не нужен), либо используй пул:

from concurrent.futures import ThreadPoolExecutor

with ThreadPoolExecutor(max_workers=4) as pool:
    futures = {pool.submit(fetch, p): p for p in PAGES}
    for future, path in futures.items():
        print(path, "->", future.result())

Сорок страниц по 0.25 секунды на четырёх воркерах отрабатывают за 2.5 секунды вместо десяти. future.result() отдаёт значение или бросает то исключение, которое случилось внутри потока.

Чем Thread(target=) отличается от наследования Thread

Ничем по возможностям. Наследуешься — переопределяешь run() вместо передачи target. Наследование берут, когда потоку нужно своё состояние и несколько методов. В остальных случаях target= короче. __init__ в наследнике обязан вызвать super().__init__(), иначе start() бросит RuntimeError: thread.__init__() not called.

Сколько потоков можно создать

Технического лимита в Python нет, упрёшься в память под стеки и в лимиты ОС. Ориентир практический: потоки под ввод-вывод держат десятками, а не тысячами, и почти всегда через ThreadPoolExecutor(max_workers=N). Пул переиспользует потоки, ограничивает нагрузку на чужой сервер и не даёт создать тысячу объектов там, где хватит восьми.

Потокобезопасен ли print

Не полностью. print("a", "b", "c") делает шесть записей в поток вывода: каждый аргумент и каждый пробел отдельно, между ними может вклиниться чужая строка. F-string сводит это к двум записям (текст и перевод строки), так что сама строка уцелеет, а перенос всё ещё может уехать. Нужен честно синхронизированный вывод, бери logging.

Что учить дальше

Потоки упираются в GIL, и это следующая тема: как он устроен и почему не мешает ждать сеть. Выбор инструмента под конкретную задачу разобран в сравнении потоков, процессов и asyncio.

Источники