W3docs

Python Threading

Изучите Python threading: создание потоков, синхронизация с Lock, очереди, ThreadPoolExecutor и роль GIL в многопоточности.

Модуль threading в Python позволяет запускать несколько задач в одном процессе одновременно. Каждая задача выполняется в собственном потоке — лёгкой единице выполнения, разделяющей память процесса. Threading — правильный инструмент, когда программа большую часть времени ожидает (чтение файла, HTTP-запрос, запрос к базе данных), а вы хотите выполнять полезную работу во время этого ожидания вместо блокировки.

В этой главе рассматриваются:

  • Создание и запуск потоков с помощью threading.Thread
  • Ожидание завершения потоков с помощью join
  • Демон-потоки и фоновые задачи
  • Предотвращение гонки данных с помощью Lock и with
  • Координация потоков с помощью Event и Semaphore
  • Потокобезопасная коммуникация через queue.Queue
  • ThreadPoolExecutor для управляемых пулов потоков
  • Глобальная блокировка интерпретатора (GIL) и почему потоки не ускоряют CPU-зависимый код
  • Когда выбирать threading вместо asyncio

Создание и запуск потока

Импортируйте threading и создайте объект Thread, передав запускаемую функцию в качестве target. Вызовите .start() для запуска потока:

import threading
import time

def greet(name):
    time.sleep(0.5)         # simulate some work
    print(f'Hello, {name}!')

t = threading.Thread(target=greet, args=('Alice',))
t.start()
print('Thread started — main continues running')
t.join()                    # wait for the thread to finish
print('Thread finished')
# Thread started — main continues running
# Hello, Alice!
# Thread finished

Ключевые моменты:

  • args — это кортеж позиционных аргументов, передаваемых в target. Используйте kwargs для именованных аргументов.
  • Без .join() главный поток может завершиться раньше, чем дочерний поток закончит работу.
  • .start() возвращает управление немедленно; новый поток выполняется конкурентно.

Передача именованных аргументов

import threading

def connect(host, port=80):
    print(f'Connecting to {host}:{port}')

t = threading.Thread(target=connect, kwargs={'host': 'example.com', 'port': 443})
t.start()
t.join()
# Connecting to example.com:443

Запуск нескольких потоков одновременно

Настоящее преимущество threading — возможность выполнять несколько задач параллельно. Сначала запустите все потоки, затем дождитесь завершения каждого:

import threading
import time

def download(url):
    time.sleep(1)           # simulate a 1-second network request
    print(f'Downloaded: {url}')

urls = [
    'https://example.com/data1',
    'https://example.com/data2',
    'https://example.com/data3',
]

start = time.perf_counter()

threads = [threading.Thread(target=download, args=(url,)) for url in urls]
for t in threads:
    t.start()
for t in threads:
    t.join()

elapsed = time.perf_counter() - start
print(f'All downloads finished in {elapsed:.1f}s')
# Downloaded: https://example.com/data1
# Downloaded: https://example.com/data2
# Downloaded: https://example.com/data3
# All downloads finished in 1.0s

Без потоков это заняло бы 3 секунды (последовательно). С тремя потоками — около 1 секунды, потому что ожидания перекрываются.

Наследование Thread

Для более сложной логики создайте подкласс threading.Thread и переопределите run(). Сохраняйте результаты как атрибуты экземпляра, чтобы вызывающий код мог прочитать их после join():

import threading
import time

class DownloadThread(threading.Thread):
    def __init__(self, url):
        super().__init__()
        self.url = url
        self.result = None

    def run(self):
        time.sleep(0.5)     # simulate download
        self.result = f'Data from {self.url}'

threads = [DownloadThread(f'https://example.com/page{i}') for i in range(3)]
for t in threads:
    t.start()
for t in threads:
    t.join()

for t in threads:
    print(t.result)
# Data from https://example.com/page0
# Data from https://example.com/page1
# Data from https://example.com/page2

Демон-потоки

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

import threading
import time

def heartbeat():
    while True:
        print('♥ still running')
        time.sleep(1)

t = threading.Thread(target=heartbeat, daemon=True)
t.start()

time.sleep(2.5)
print('Main thread exiting — daemon will be killed')
# ♥ still running
# ♥ still running
# Main thread exiting — daemon will be killed

Используйте демон-потоки для фонового мониторинга или задач логирования, которые не должны препятствовать завершению программы. Никогда не используйте их для задач, которые должны завершаться корректно (запись файлов, коммиты в базу данных) — они завершаются без какой-либо очистки.

Имена потоков и интроспекция

Каждый поток имеет имя. Вы можете задать его явно или позволить Python присвоить его автоматически. Используйте threading.current_thread() для проверки текущего потока и threading.active_count() для подсчёта активных потоков:

import threading

def worker():
    t = threading.current_thread()
    print(f'Running in thread: {t.name}')

t = threading.Thread(target=worker, name='WorkerThread-1')
t.start()
t.join()

print(f'Active threads: {threading.active_count()}')
# Running in thread: WorkerThread-1
# Active threads: 1

Синхронизация: предотвращение гонки данных

Потоки разделяют память процесса. Когда два потока одновременно читают и записывают одну и ту же переменную, возникает гонка данных — недетерминированное поведение, которое сложно воспроизвести или отладить.

Следующий пример без блокировки выдаёт непредсказуемый итоговый счётчик, потому что инкременты из разных потоков могут перекрываться:

import threading

counter = 0

def unsafe_increment():
    global counter
    for _ in range(100_000):
        counter += 1         # read-modify-write: not atomic!

threads = [threading.Thread(target=unsafe_increment) for _ in range(5)]
for t in threads:
    t.start()
for t in threads:
    t.join()

# counter is somewhere between 100000 and 500000 — unpredictable
print('Final counter:', counter)

Lock

threading.Lock гарантирует, что только один поток выполняет защищённый раздел в одно время. Используйте его как менеджер контекста с with, чтобы блокировка всегда освобождалась, даже при возникновении исключения:

import threading

counter = 0
lock = threading.Lock()

def safe_increment():
    global counter
    for _ in range(100_000):
        with lock:              # acquire before read-modify-write
            counter += 1        # now only one thread at a time can run this

threads = [threading.Thread(target=safe_increment) for _ in range(5)]
for t in threads:
    t.start()
for t in threads:
    t.join()

print('Final counter:', counter)   # always 500000

RLock (повторно входимая блокировка)

Если потоку нужно получить одну и ту же блокировку дважды (например, один метод вызывает другой метод, который тоже получает эту блокировку), используйте threading.RLock. Он позволяет одному и тому же потоку повторно получить блокировку без взаимоблокировки:

import threading

lock = threading.RLock()

def outer():
    with lock:
        print('Outer acquired')
        inner()             # inner also acquires the same lock

def inner():
    with lock:              # works because RLock counts acquisitions
        print('Inner acquired')

t = threading.Thread(target=outer)
t.start()
t.join()
# Outer acquired
# Inner acquired

Координация потоков: Event и Semaphore

Event

threading.Event — простой сигнал. Один поток вызывает .set() для сигнализации; другие потоки вызывают .wait(), чтобы блокироваться до получения сигнала:

import threading
import time

ready = threading.Event()

def worker():
    print('Worker: waiting for signal...')
    ready.wait()            # blocks here until ready.set() is called
    print('Worker: signal received, starting work')

t = threading.Thread(target=worker)
t.start()

time.sleep(0.5)
print('Main: sending signal')
ready.set()
t.join()
# Worker: waiting for signal...
# Main: sending signal
# Worker: signal received, starting work

Используйте Event для координации порядка запуска — например, чтобы задержать рабочие потоки до установления соединения с базой данных.

Semaphore

threading.Semaphore ограничивает количество потоков, которые могут одновременно находиться в разделе кода. Это полезно для ограничения доступа к общему ресурсу, например пулу соединений:

import threading
import time

# Allow at most 2 threads to enter the critical section at once
semaphore = threading.Semaphore(2)

def use_connection(name):
    with semaphore:
        print(f'{name}: using connection')
        time.sleep(0.5)
        print(f'{name}: releasing connection')

threads = [threading.Thread(target=use_connection, args=(f'T{i}',)) for i in range(4)]
for t in threads:
    t.start()
for t in threads:
    t.join()
# T0: using connection
# T1: using connection          <- only 2 at a time
# T0: releasing connection
# T2: using connection
# T1: releasing connection
# T3: using connection
# T2: releasing connection
# T3: releasing connection

Локальные данные потока

threading.local() создаёт объект, хранящий отдельные значения для каждого потока. Это полезно для кешей или курсоров базы данных на уровне потока:

import threading

local_data = threading.local()

def set_user(name):
    local_data.user = name       # each thread writes its own copy
    print(f'{threading.current_thread().name}: user = {local_data.user}')

threads = [
    threading.Thread(target=set_user, args=(f'user{i}',), name=f'Thread-{i}')
    for i in range(3)
]
for t in threads:
    t.start()
for t in threads:
    t.join()
# Thread-0: user = user0
# Thread-1: user = user1
# Thread-2: user = user2

Обращение к local_data.user в потоке, где это значение никогда не задавалось, вызовет AttributeError, как и при обычном обращении к несуществующему атрибуту.

Потокобезопасные очереди

Класс queue.Queue (из стандартного модуля queue, не asyncio) — это потокобезопасная очередь FIFO. Потоки могут вызывать put и get без блокировки — вся синхронизация обрабатывается внутренне.

Классический паттерн — производитель-потребитель: один или несколько потоков-производителей генерируют работу, а потоки-потребители обрабатывают её:

import threading
import queue
import time

q = queue.Queue(maxsize=5)

def producer():
    for i in range(1, 5):
        q.put(f'item-{i}')
        print(f'Produced item-{i}')
        time.sleep(0.05)

def consumer():
    while True:
        item = q.get()
        if item is None:        # sentinel: stop when None is received
            break
        print(f'Consumed {item}')
        q.task_done()

prod = threading.Thread(target=producer)
cons = threading.Thread(target=consumer)

cons.start()
prod.start()
prod.join()
q.put(None)     # signal consumer to stop
cons.join()
# Produced item-1
# Consumed item-1
# Produced item-2
# Consumed item-2
# Produced item-3
# Consumed item-3
# Produced item-4
# Consumed item-4

queue.Queue также предоставляет task_done() и join() для отслеживания обработки всех элементов очереди, а также queue.LifoQueue / queue.PriorityQueue для альтернативных порядков.

ThreadPoolExecutor: управляемые пулы потоков

Создавать новый объект Thread для каждой задачи расточительно при большом количестве короткоживущих задач. concurrent.futures.ThreadPoolExecutor управляет пулом переиспользуемых рабочих потоков и возвращает объекты Future для каждой отправленной задачи:

import concurrent.futures
import time

def fetch_url(url):
    time.sleep(0.5)     # simulate network I/O
    return f'Response from {url}'

urls = [
    'https://api.example.com/users',
    'https://api.example.com/posts',
    'https://api.example.com/comments',
]

with concurrent.futures.ThreadPoolExecutor(max_workers=3) as executor:
    # submit all tasks and get Future objects
    futures = {executor.submit(fetch_url, url): url for url in urls}

    for future in concurrent.futures.as_completed(futures):
        url = futures[future]
        print(future.result())

# Response from https://api.example.com/users    (order may vary)
# Response from https://api.example.com/comments
# Response from https://api.example.com/posts

executor.map(fn, iterable) — более краткая форма, когда не нужны отдельные объекты Future:

import concurrent.futures
import time

def square(n):
    time.sleep(0.01)
    return n * n

with concurrent.futures.ThreadPoolExecutor(max_workers=4) as executor:
    results = list(executor.map(square, range(10)))

print(results)
# [0, 1, 4, 9, 16, 25, 36, 49, 64, 81]

executor.map сохраняет порядок входных данных в выводе, в отличие от as_completed, который выдаёт результаты в порядке завершения.

Глобальная блокировка интерпретатора (GIL)

CPython (стандартный интерпретатор Python) имеет Global Interpreter Lock — мьютекс, который позволяет только одному потоку выполнять байткод Python за раз. Это означает, что потоки в CPython не могут выполнять Python-код в истинном параллелизме на нескольких ядрах CPU.

Практическое следствие:

  • I/O-зависимые задачи: потоки действительно ускоряют программу. Пока один поток ожидает сетевого ответа, GIL освобождается и другой поток получает управление. Все приведённые выше примеры демонстрируют это поведение.
  • CPU-зависимые задачи: потоки не ускоряют работу и могут даже немного замедлить её из-за накладных расходов на переключение контекста.
import threading
import time

def cpu_bound(n):
    total = 0
    for i in range(n):
        total += i
    return total

# Sequential
start = time.perf_counter()
cpu_bound(5_000_000)
cpu_bound(5_000_000)
single = time.perf_counter() - start

# Two threads — GIL prevents true parallelism
start = time.perf_counter()
t1 = threading.Thread(target=cpu_bound, args=(5_000_000,))
t2 = threading.Thread(target=cpu_bound, args=(5_000_000,))
t1.start(); t2.start()
t1.join(); t2.join()
threaded = time.perf_counter() - start

print(f'Single-threaded: {single:.2f}s')
print(f'Two threads:     {threaded:.2f}s')
# Two threads are NOT faster (similar elapsed time)

Для истинного параллелизма CPU в Python используйте multiprocessing или concurrent.futures.ProcessPoolExecutor — каждый процесс имеет собственный GIL.

Threading vs. asyncio

И threading, и asyncio ускоряют I/O-зависимые программы, но работают по-разному:

threadingasyncio
Модель параллелизмаВытесняющая — ОС переключает потокиКооперативная — корутины передают управление в await
Лучше всего подходит дляБлокирующих сторонних библиотекБиблиотек с поддержкой async (aiohttp, asyncpg)
Общее состояниеТребует явных блокировокБезопасно в рамках одного цикла событий
Накладные расходыОдин поток ОС на задачуОчень низкие — тысячи корутин в одном потоке
Кривая обученияПривычная (синхронный стиль кода)Требует async/await повсюду

Правило большого пальца: если вы используете библиотеку, у которой есть async-совместимая версия (например, aiohttp вместо requests), выбирайте asyncio. Если вы вынуждены работать с синхронными блокирующими библиотеками, используйте threading. Для CPU-зависимой работы используйте multiprocessing.

Распространённые ошибки

Запуск потока дважды. Вызов .start() на одном объекте Thread более одного раза вызывает RuntimeError. Создавайте новый экземпляр Thread для каждого выполнения.

Забытый join. Поток без join может всё ещё выполняться, когда программа завершается. Всегда вызывайте join для потоков, завершение которых важно, или делайте их демонами, если они действительно выполняются в режиме «запустил и забыл».

Удержание блокировки слишком долго. Блокировка большого блока кода сводит на нет смысл параллелизма. Держите заблокированные секции как можно короче — защищайте только операцию чтения-изменения-записи.

Взаимоблокировка (Deadlock). Взаимоблокировка происходит, когда два потока держат блокировку, которую ожидает другой. Предотвратите это, всегда получая несколько блокировок в одном порядке во всех потоках.

import threading

lock_a = threading.Lock()
lock_b = threading.Lock()

# DEADLOCK: Thread 1 holds lock_a, waits for lock_b
#           Thread 2 holds lock_b, waits for lock_a

# FIX: always acquire locks in the same order (lock_a then lock_b) in every thread

Изменение списка во время итерации в другом потоке. Оберните весь доступ (чтение и запись) к общим коллекциям блокировкой, чтобы избежать RuntimeError: list changed size during iteration.

Краткий справочник

ИнструментНазначение
threading.Thread(target=fn, args=(...))Создать новый поток
t.start()Запустить поток
t.join()Ожидать завершения потока
t.daemon = TrueПометить как фоновый поток (завершается при выходе)
threading.Lock()Взаимное исключение — только один поток за раз
threading.RLock()Повторно входимая блокировка — один поток может получить несколько раз
threading.Event()Одноразовый сигнал между потоками
threading.Semaphore(n)Ограничить до n одновременных потоков в секции
threading.local()Хранилище данных уровня потока
queue.QueueПотокобезопасная FIFO для паттерна производитель-потребитель
ThreadPoolExecutor(max_workers=n)Управляемый пул переиспользуемых рабочих потоков

Практика

Практика
Which of the following tasks would benefit most from Python threading?
Which of the following tasks would benefit most from Python threading?
Was this page helpful?