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 500000RLock (повторно входимая блокировка)
Если потоку нужно получить одну и ту же блокировку дважды (например, один метод вызывает другой метод, который тоже получает эту блокировку), используйте 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-4queue.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/postsexecutor.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-зависимые программы, но работают по-разному:
threading | asyncio | |
|---|---|---|
| Модель параллелизма | Вытесняющая — ОС переключает потоки | Кооперативная — корутины передают управление в 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) | Управляемый пул переиспользуемых рабочих потоков |