Перейти к содержимому
k10

k10. Потоки, asyncio и очередь сообщений

k10. Потоки, asyncio и очередь сообщений

UAV Python Pro · Конспект K10 · Модуль M1. Python intensive: потоки, asyncio и очередь сообщений · Редакция 1.0 · После CHECKPOINT-K09

До сих пор учебные программы выполнялись в одном потоке управления: строка за строкой, без параллельной работы. В программах для БПЛА часто нужно одновременно принимать поток телеметрии и принимать решения, ставить команды в очередь и не блокировать весь процесс навечно, если сообщение задерживается. На этом уроке вы познакомитесь с тремя согласованными идеями: поток threading, очередь queue.Queue между производителем и потребителем, кооперативные задачи asyncio. Это учебные модели. Они готовят к архитектуре скрипта на pymavlink, но сами по себе ещё не открывают порт автопилота.

Как устроены примеры. В теории каждый важный фрагмент кода дан дважды: сначала блок с подписью «Чистый код», затем раскрывающийся построчный разбор. В практике для каждой задачи в аккордеоне приведён полный исходный текст файла k10_task*.py.

Если один участок кода долго ждёт ввод-вывод, остальная логика не должна «замирать» без явной необходимости. Поток позволяет выполнять функцию параллельно с main. Очередь безопасно передаёт сообщения между потоками. asyncio позволяет на одном потоке событий кооперативно переключаться между задачами во время ожидания. Для данных, которыми делятся потоки, нужна синхронизация, например Lock. Для ожидания сообщения нужен таймаут, иначе программа может ждать бесконечно.

Содержание

1. Цели урока

После изучения этого урока вы научитесь:

  • объяснять, зачем в программах для БПЛА разделяют приём данных и принятие решений.
  • запускать фоновый поток threading.Thread и останавливать его через флаг.
  • защищать общие данные с помощью threading.Lock.
  • передавать сообщения между потоками через queue.Queue.
  • запускать две кооперативные async-задачи через asyncio.run и asyncio.gather.
  • использовать queue.get с timeout и обрабатывать queue.Empty.
  • формулировать, когда на курсе уместнее поток с очередью, а когда достаточно asyncio как модели ожидания.
Связь с разработкой программного обеспечения для БПЛА. Скрипт наземной станции или бортового компьютера-компаньона редко работает как один короткий калькулятор. Он принимает поток сообщений, обновляет состояние, ставит команды в очередь и пишет журнал. Если всё делать строго «в лоб» в одном блокирующем цикле без структуры, код быстро становится хрупким. Этот урок даёт каркас, на который позже ляжет работа с pymavlink.
↑ К оглавлению

2. Словарь урока

ТерминПростыми словами
Поток (thread) Отдельная линия выполнения кода внутри одного процесса программы.
threading Стандартный модуль Python для создания и управления потоками.
Гонка данных Ситуация, когда два потока одновременно читают и пишут одни и те же данные без согласования, из-за чего результат становится непредсказуемым.
Lock Блокировка: в каждый момент времени критический участок кода выполняет только один поток.
Очередь Queue Структура «первым пришёл — первым ушёл» для безопасной передачи объектов между потоками.
Производитель и потребитель Роли: один код публикует сообщения, другой их обрабатывает.
asyncio Модуль кооперативной многозадачности на цикле событий: задачи сами отдают управление во время ожидания.
Корутина Функция async def, которую можно приостанавливать на await и затем продолжать.
Таймаут Предельное время ожидания события, после которого ожидание прекращают.
GIL Global Interpreter Lock в CPython: ограничение на одновременное исполнение байткода Python в нескольких потоках; для ожидания ввода-вывода потоки всё равно полезны.
↑ К оглавлению

3. Зачем параллельность в программах для БПЛА

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

Поэтому в инженерии разделяют роли. Один участок кода отвечает за приём. Другой участок отвечает за решение. Между ними передают сообщения через очередь или через аккуратно защищённое состояние. Настоящий MAVLink появится позже. Сейчас вы строите ту же схему на учебных числах и словарях.

↑ К оглавлению

4. Поток threading и защита данных

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

Если два потока читают и пишут одну и ту же переменную без согласования, возникает гонка данных. Чтобы её избежать, критический участок окружают Lock: в каждый момент запись или согласованное чтение выполняет только один поток.

Чистый код

# Пример: фоновый поток публикует данные, основной поток читает снимок
import threading
import time
from dataclasses import dataclass

@dataclass
class TelemetrySample:
    t_s: float
    altitude_m: float

class TelemetryHub:
    def __init__(self) -> None:
        self._lock = threading.Lock()
        self._latest: TelemetrySample | None = None
        self._stop = False

    def publish(self, sample: TelemetrySample) -> None:
        with self._lock:
            self._latest = sample

    def snapshot(self) -> TelemetrySample | None:
        with self._lock:
            return self._latest

    def request_stop(self) -> None:
        self._stop = True

    @property
    def stop(self) -> bool:
        return self._stop

def producer(hub: TelemetryHub) -> None:
    t = 0.0
    alt = 0.0
    while not hub.stop:
        hub.publish(TelemetrySample(t_s=t, altitude_m=alt))
        t += 0.2
        alt += 1.5
        time.sleep(0.2)

hub = TelemetryHub()
worker = threading.Thread(target=producer, args=(hub,), daemon=True)
worker.start()
time.sleep(0.5)
print(hub.snapshot())
hub.request_stop()
worker.join(timeout=1.0)
  • Построчный разбор

    Ниже приведён тот же фрагмент с пояснением после каждой смысловой строки. Строки, которые начинаются с последовательности символов «# →», содержат комментарий преподавателя. Эти строки не являются обязательной частью программы.

    import threading
    # → Модуль стандартной библиотеки для работы с потоками.
    
    import time
    # → Нужен для пауз sleep.
    
    from dataclasses import dataclass
    # → Учебный снимок телеметрии как dataclass.
    
    @dataclass
    # → Декоратор для структуры данных.
    
    class TelemetrySample:
    # → Один отсчёт: время и высота.
    
        t_s: float
    # → Время в секундах.
    
        altitude_m: float
    # → Высота в метрах.
    
    class TelemetryHub:
    # → Общее хранилище последнего отсчёта для двух потоков.
    
        def __init__(self) -> None:
    # → Создаём блокировку и пустое состояние.
    
            self._lock = threading.Lock()
    # → Lock защищает чтение и запись _latest от одновременного доступа.
    
            self._latest: TelemetrySample | None = None
    # → Пока данных нет, хранится None.
    
            self._stop = False
    # → Внутренний флаг остановки фонового цикла.
    
        def publish(self, sample: TelemetrySample) -> None:
    # → Публикация нового отсчёта из фонового потока.
    
            with self._lock:
    # → Захватываем lock на время записи.
    
                self._latest = sample
    # → Обновляем последний снимок.
    
        def snapshot(self) -> TelemetrySample | None:
    # → Чтение снимка из основного потока.
    
            with self._lock:
    # → Та же блокировка, чтобы не прочитать противоречивое состояние.
    
                return self._latest
    # → Возвращаем текущее значение или None.
    
        def request_stop(self) -> None:
    # → Публичный способ попросить фоновый цикл завершиться.
    
            self._stop = True
    # → Устанавливаем флаг остановки.
    
        @property
    # → Свойство stop даёт читать флаг без прямого доступа к _stop снаружи.
    
        def stop(self) -> bool:
    # → Возвращает True, если работа должна быть остановлена.
    
            return self._stop
    # → Текущее значение флага.
    
    def producer(hub: TelemetryHub) -> None:
    # → Функция, которую выполнит фоновый поток.
    
        t = 0.0
    # → Учебное время.
    
        alt = 0.0
    # → Учебная высота.
    
        while not hub.stop:
    # → Цикл, пока основной поток не запросит остановку.
    
            hub.publish(TelemetrySample(t_s=t, altitude_m=alt))
    # → Кладём новый отсчёт в общее хранилище.
    
            t += 0.2
    # → Продвигаем учебное время.
    
            alt += 1.5
    # → Продвигаем учебную высоту.
    
            time.sleep(0.2)
    # → Пауза: имитация периода приёма сообщений.
    
    hub = TelemetryHub()
    # → Создаём общее хранилище.
    
    worker = threading.Thread(target=producer, args=(hub,), daemon=True)
    # → Создаём поток. Аргумент daemon=True означает, что поток не удерживает процесс после выхода main в этом учебном сценарии.
    
    worker.start()
    # → Запускаем фоновую работу.
    
    time.sleep(0.5)
    # → Основной поток коротко ждёт появления данных.
    
    print(hub.snapshot())
    # → Читаем последний снимок телеметрии.
    
    hub.request_stop()
    # → Просим производителя остановиться.
    
    worker.join(timeout=1.0)
    # → Ждём завершения потока не дольше одной секунды.
О GIL. В реализации CPython существует Global Interpreter Lock: несколько потоков не исполняют байткод Python абсолютно одновременно на разных ядрах так, как это бывает в некоторых других языках. Тем не менее для задач ожидания ввода-вывода, пауз и параллельной структуры программы потоки остаются практичным инструментом. Для тяжёлой математики на нескольких ядрах позже могут понадобиться другие средства. На этом уроке достаточно понимать, что поток — это прежде всего способ не блокировать всю логику на одном ожидании.
Связь с практикой. Полный текст файла k10_task1_thread_worker.py приведён в задаче 10.1. В нём основной поток несколько раз читает снимок телеметрии, затем запрашивает остановку фонового потока.
↑ К оглавлению

5. Очередь Queue между производителем и потребителем

Очередь Queue из модуля queue предназначена для передачи объектов между потоками. Один код выступает в роли производителя и вызывает put. Другой код выступает в роли потребителя и вызывает get. Очередь сама содержит необходимую синхронизацию для этих операций.

Такая схема хорошо ложится на контур «событие телеметрии → решение → команда». Производитель не обязан знать, как именно будет обработано событие. Потребитель не обязан знать, откуда событие взялось.

Чистый код

# Пример: две очереди — события телеметрии и команды
import queue
import threading
import time
from dataclasses import dataclass

@dataclass(frozen=True)
class TelemetryEvent:
    altitude_m: float

@dataclass(frozen=True)
class Command:
    name: str
    detail: str

def reader(events: queue.Queue) -> None:
    for alt in (0.5, 8.0, 15.0):
        events.put(TelemetryEvent(altitude_m=alt))
        time.sleep(0.05)
    events.put(None)

def worker(events: queue.Queue, commands: queue.Queue) -> None:
    while True:
        item = events.get()
        if item is None:
            commands.put(None)
            break
        if item.altitude_m >= 10.0:
            commands.put(Command("HOLD", f"alt={item.altitude_m}"))
        else:
            commands.put(Command("CLIMB", f"alt={item.altitude_m}"))

events: queue.Queue = queue.Queue()
commands: queue.Queue = queue.Queue()
threading.Thread(target=reader, args=(events,), daemon=True).start()
threading.Thread(target=worker, args=(events, commands), daemon=True).start()
while True:
    cmd = commands.get()
    if cmd is None:
        break
    print(cmd)
  • Построчный разбор

    Ниже приведён тот же фрагмент с пояснением после каждой смысловой строки. Строки, которые начинаются с последовательности символов «# →», содержат комментарий преподавателя. Эти строки не являются обязательной частью программы.

    import queue
    # → Очередь Queue безопасна для передачи объектов между потоками.
    
    import threading
    # → Запуск потоков reader и worker.
    
    import time
    # → Небольшая пауза в учебном издателе.
    
    from dataclasses import dataclass
    # → Структуры события и команды.
    
    @dataclass(frozen=True)
    # → Неизменяемое событие телеметрии.
    
    class TelemetryEvent:
    # → Сообщение «пришла высота».
    
        altitude_m: float
    # → Высота в метрах.
    
    @dataclass(frozen=True)
    # → Неизменяемая команда-решение.
    
    class Command:
    # → Учебная команда для печати (не MAVLink).
    
        name: str
    # → Имя команды, например HOLD или CLIMB.
    
        detail: str
    # → Пояснение для журнала.
    
    def reader(events: queue.Queue) -> None:
    # → Поток-издатель кладёт события в очередь events.
    
        for alt in (0.5, 8.0, 15.0):
    # → Три учебных отсчёта высоты.
    
            events.put(TelemetryEvent(altitude_m=alt))
    # → put добавляет элемент в конец очереди.
    
            time.sleep(0.05)
    # → Короткая пауза между событиями.
    
        events.put(None)
    # → Специальный маркер конца потока данных.
    
    def worker(events: queue.Queue, commands: queue.Queue) -> None:
    # → Поток-обработчик читает события и пишет команды.
    
        while True:
    # → Цикл до маркера конца.
    
            item = events.get()
    # → get ждёт и забирает следующий элемент из очереди.
    
            if item is None:
    # → Получен маркер конца.
    
                commands.put(None)
    # → Передаём конец дальше по конвейеру.
    
                break
    # → Выходим из цикла обработчика.
    
            if item.altitude_m >= 10.0:
    # → Учебное правило принятия решения.
    
                commands.put(Command("HOLD", f"alt={item.altitude_m}"))
    # → Высота достаточна — команда удержания.
    
            else:
    # → Иначе — набор высоты.
    
                commands.put(Command("CLIMB", f"alt={item.altitude_m}"))
    # → Команда набора.
    
    events: queue.Queue = queue.Queue()
    # → Очередь событий телеметрии.
    
    commands: queue.Queue = queue.Queue()
    # → Очередь команд-решений.
    
    threading.Thread(target=reader, args=(events,), daemon=True).start()
    # → Запуск издателя.
    
    threading.Thread(target=worker, args=(events, commands), daemon=True).start()
    # → Запуск обработчика.
    
    while True:
    # → Основной поток читает готовые команды.
    
        cmd = commands.get()
    # → Ожидание следующей команды.
    
        if cmd is None:
    # → Конец конвейера.
    
            break
    # → Выход.
    
        print(cmd)
    # → Печать учебной команды.
Маркер конца None. В учебных примерах None в очереди часто означает «данных больше не будет». В боевом коде для этого иногда вводят отдельный объект-сигнал, чтобы не путать его с полезными данными. Здесь None достаточно для ясности сценария.
Связь с практикой. Полный текст файла k10_task2_queue_pipeline.py приведён в задаче 10.2.
↑ К оглавлению

6. Знакомство с asyncio

Модуль asyncio реализует кооперативную многозадачность на цикле событий. Функция, объявленная как корутина через async def, может приостанавливаться на await и возвращать управление циклу, чтобы в это время выполнялась другая задача.

Это не то же самое, что два потока операционной системы. Здесь один поток событий по очереди продвигает задачи, когда они ждут. Для сетевых библиотек с async API такой стиль очень естественен. В курсе pymavlink чаще встречается в синхронном виде, поэтому asyncio на M1 нужен как понятная модель «не блокировать всё ожидание навсегда» и как подготовка к возможным async-инструментам позже.

Чистый код

# Пример: две async-задачи в одном цикле событий
import asyncio
from dataclasses import dataclass

@dataclass
class FakeLink:
    ticks: int = 0

async def telemetry_ticker(link: FakeLink, n: int) -> None:
    for i in range(n):
        await asyncio.sleep(0.1)
        link.ticks += 1
        print(f"telemetry tick {i + 1}, total={link.ticks}")

async def status_printer(link: FakeLink, n: int) -> None:
    for i in range(n):
        await asyncio.sleep(0.15)
        print(f"status {i + 1}: ticks_seen={link.ticks}")

async def main() -> None:
    link = FakeLink()
    await asyncio.gather(
        telemetry_ticker(link, 5),
        status_printer(link, 4),
    )
    print("asyncio done, ticks=", link.ticks)

asyncio.run(main())
  • Построчный разбор

    Ниже приведён тот же фрагмент с пояснением после каждой смысловой строки. Строки, которые начинаются с последовательности символов «# →», содержат комментарий преподавателя. Эти строки не являются обязательной частью программы.

    import asyncio
    # → Стандартный модуль асинхронного ввода-вывода и кооперативных задач.
    
    from dataclasses import dataclass
    # → Простое состояние канала.
    
    @dataclass
    # → Учебная структура.
    
    class FakeLink:
    # → Имитация состояния: сколько тиков телеметрии уже «принято».
    
        ticks: int = 0
    # → Счётчик тиков.
    
    async def telemetry_ticker(link: FakeLink, n: int) -> None:
    # → Корутина: async def означает, что внутри можно await.
    
        for i in range(n):
    # → n учебных тиков.
    
            await asyncio.sleep(0.1)
    # → await отдаёт управление циклу событий на время паузы, не блокируя весь процесс так, как это делает обычный time.sleep в однозадачном коде с несколькими задачами.
    
            link.ticks += 1
    # → Увеличиваем счётчик «принятых» тиков.
    
            print(f"telemetry tick {i + 1}, total={link.ticks}")
    # → Печать хода приёма.
    
    async def status_printer(link: FakeLink, n: int) -> None:
    # → Вторая корутина печатает сводку.
    
        for i in range(n):
    # → n сводок.
    
            await asyncio.sleep(0.15)
    # → Другая длительность паузы — задачи чередуются.
    
            print(f"status {i + 1}: ticks_seen={link.ticks}")
    # → Видим, сколько тиков успело накопиться.
    
    async def main() -> None:
    # → Точка сборки асинхронных задач.
    
        link = FakeLink()
    # → Общее состояние для обеих задач.
    
        await asyncio.gather(
    # → gather запускает несколько корутин и ждёт завершения всех.
    
            telemetry_ticker(link, 5),
    # → Первая задача.
    
            status_printer(link, 4),
    # → Вторая задача.
    
        )
    # → Конец gather.
    
        print("asyncio done, ticks=", link.ticks)
    # → Итоговый счётчик.
    
    asyncio.run(main())
    # → Создаёт цикл событий, выполняет main и корректно завершает цикл.
Не смешивайте без нужды time.sleep и asyncio. Обычный time.sleep блокирует весь поток. Внутри async-задач для паузы используют await asyncio.sleep(...). Иначе кооперативность пропадает.
Связь с практикой. Полный текст файла k10_task3_asyncio_tick.py приведён в задаче 10.3.
↑ К оглавлению

7. Ожидание с таймаутом

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

В программах для БПЛА таймаут критичен. Нельзя навсегда «повиснуть» в ожидании сообщения, если канал пропал. На этом уроке таймаут показан на queue.get. Позднее та же идея появится в ожидании сообщений MAVLink.

Чистый код

# Пример: queue.get с timeout
import queue
import threading
import time

def slow_publisher(q: queue.Queue) -> None:
    time.sleep(0.4)
    q.put({"type": "HEARTBEAT", "ok": True})

q: queue.Queue = queue.Queue()
threading.Thread(target=slow_publisher, args=(q,), daemon=True).start()

try:
    msg = q.get(timeout=0.2)
    print("early:", msg)
except queue.Empty:
    print("timeout 0.2s: сообщения ещё нет")

msg = q.get(timeout=0.5)
print("received:", msg)
  • Построчный разбор

    Ниже приведён тот же фрагмент с пояснением после каждой смысловой строки. Строки, которые начинаются с последовательности символов «# →», содержат комментарий преподавателя. Эти строки не являются обязательной частью программы.

    import queue
    # → Очередь и исключение queue.Empty.
    
    import threading
    # → Фоновая публикация сообщения.
    
    import time
    # → Задержка издателя.
    
    def slow_publisher(q: queue.Queue) -> None:
    # → Издатель кладёт сообщение не сразу.
    
        time.sleep(0.4)
    # → Искусственная задержка 0.4 с.
    
        q.put({"type": "HEARTBEAT", "ok": True})
    # → Учебное сообщение, похожее по смыслу на heartbeat.
    
    q: queue.Queue = queue.Queue()
    # → Пустая очередь.
    
    threading.Thread(target=slow_publisher, args=(q,), daemon=True).start()
    # → Запуск медленного издателя.
    
    try:
    # → Первая попытка чтения с коротким таймаутом.
    
        msg = q.get(timeout=0.2)
    # → Ждать элемент не дольше 0.2 с.
    
        print("early:", msg)
    # → Если успели — напечатать (в этом сценарии обычно не успеваем).
    
    except queue.Empty:
    # → Если за 0.2 с ничего не появилось, get поднимает queue.Empty.
    
        print("timeout 0.2s: сообщения ещё нет")
    # → Ожидаемая ветка для короткого таймаута.
    
    msg = q.get(timeout=0.5)
    # → Вторая попытка с большим запасом времени.
    
    print("received:", msg)
    # → К этому моменту сообщение уже должно быть в очереди.
Связь с практикой. Полный текст файла k10_task4_timeout_queue.py приведён в задаче 10.4.
↑ К оглавлению

8. Что выбирать на практике курса

СитуацияПредпочтительный приём на M1
Фоновое чтение и основной цикл решений, синхронные библиотеки threading + Queue или threading + Lock для снимка состояния
Несколько ожиданий ввода-вывода в кооперативном стиле, async API asyncio
Нужно не ждать событие вечно timeout у get / wait; явная ветка «не дождались»
Тяжёлая численная нагрузка на все ядра процессора Не тема этого урока; потоки Python здесь не главный инструмент
Практическое правило курса. Для ближайших модулей с pymavlink чаще всего встретится схема «фоновый приём + очередь или общий снимок состояния + основной цикл». asyncio вы должны узнавать и уметь запустить простой пример, но не обязаны переписывать на него весь курс.
↑ К оглавлению

9. Типичные ошибки

ОшибкаЧто происходитКак следует рассуждать
Общие переменные без Lock Редкие «плавающие» ошибки Либо Queue, либо Lock вокруг общего состояния
Поток без условия остановки Программа не завершается Флаг stop, join с timeout, понятный критерий выхода
get() без timeout в учебном UI-цикле Вечное ожидание timeout + обработка Empty
time.sleep внутри async-задачи Цикл событий блокируется await asyncio.sleep
Смешение «команда напечатана» и «команда ушла в автопилот» Ложное чувство управления Пока это учебные Command; MAVLink-отправка будет отдельно и с подтверждениями
↑ К оглавлению

10. Практика

В каждой карточке задачи ниже в раскрывающемся блоке приведён полный текст соответствующего учебного файла. Этот текст совпадает с файлами в каталоге code/M01_python/. Рекомендуется такой порядок работы. Сначала внимательно прочитайте условие. Затем раскройте аккордеон, скопируйте код в файл проекта под тем же именем, запустите его в виртуальном окружении и сверьте вывод с ожидаемым результатом.

Задача 10.1. Фоновый поток и снимок телеметрии

Файл практики: k10_task1_thread_worker.py.

Запустите фоновый поток, который публикует учебные отсчёты высоты. В основном потоке несколько раз прочитайте snapshot и затем корректно запросите остановку.

Ожидаемый результат. Несколько строк t=… alt=… с растущей высотой, затем done. Точные числа могут чуть отличаться из-за планировщика, но динамика должна быть видна.

  • Полный текст файла k10_task1_thread_worker.py

    Скопируйте приведённое ниже содержимое в файл code/M01_python/k10_task1_thread_worker.py в проекте PyCharm (виртуальное окружение) и запустите программу. Ниже приведён полный учебный текст файла с комментариями.

    # k10_task1_thread_worker.py
    # Цель: фоновый поток имитирует приём телеметрии, основной поток читает последние данные.
    # Это учебная модель, а не реальный MAVLink.
    
    import threading
    import time
    from dataclasses import dataclass
    
    
    @dataclass
    class TelemetrySample:
        t_s: float
        altitude_m: float
    
    
    class TelemetryHub:
        """Хранилище последнего отсчёта. Доступ из двух потоков защищён Lock."""
    
        def __init__(self) -> None:
            self._lock = threading.Lock()
            self._latest: TelemetrySample | None = None
            self._stop = False
    
        def publish(self, sample: TelemetrySample) -> None:
            with self._lock:
                self._latest = sample
    
        def snapshot(self) -> TelemetrySample | None:
            with self._lock:
                return self._latest
    
        def request_stop(self) -> None:
            self._stop = True
    
        @property
        def stop(self) -> bool:
            return self._stop
    
    
    def producer(hub: TelemetryHub) -> None:
        """Фоновый поток: раз в 0.2 с публикует учебную «телеметрию»."""
        t = 0.0
        alt = 0.0
        while not hub.stop:
            hub.publish(TelemetrySample(t_s=t, altitude_m=alt))
            t += 0.2
            alt += 1.5
            time.sleep(0.2)
    
    
    def main() -> None:
        hub = TelemetryHub()
        worker = threading.Thread(target=producer, args=(hub,), daemon=True)
        worker.start()
    
        # Основной поток несколько раз читает снимок
        for _ in range(5):
            time.sleep(0.25)
            sample = hub.snapshot()
            if sample is None:
                print("no data yet")
            else:
                print(f"t={sample.t_s:.1f} alt={sample.altitude_m:.1f}")
    
        hub.request_stop()
        worker.join(timeout=1.0)
        print("done")
    
    
    if __name__ == "__main__":
        main()

Задача 10.2. Конвейер Queue: событие → команда

Файл практики: k10_task2_queue_pipeline.py.

Соберите две очереди. Издатель кладёт TelemetryEvent. Обработчик ставит Command HOLD или CLIMB по учебному порогу 10 м. Основной поток печатает команды до маркера конца.

Ожидаемый результат. Серия строк command: CLIMB/HOLD с деталями высоты, затем pipeline done.

  • Полный текст файла k10_task2_queue_pipeline.py

    Скопируйте приведённое ниже содержимое в файл code/M01_python/k10_task2_queue_pipeline.py в проекте PyCharm (виртуальное окружение) и запустите программу. Ниже приведён полный учебный текст файла с комментариями.

    # k10_task2_queue_pipeline.py
    # Цель: Queue связывает производителя сообщений и обработчика команд.
    # Учебная модель контура «телеметрия пришла → решение → команда в очереди».
    
    import queue
    import threading
    import time
    from dataclasses import dataclass
    
    
    @dataclass(frozen=True)
    class TelemetryEvent:
        altitude_m: float
    
    
    @dataclass(frozen=True)
    class Command:
        name: str
        detail: str
    
    
    def reader(events: queue.Queue, stop_flag: threading.Event) -> None:
        """Имитация потока чтения: кладёт события телеметрии в очередь."""
        altitudes = [0.5, 2.0, 8.0, 15.0, 12.0]
        for alt in altitudes:
            if stop_flag.is_set():
                break
            events.put(TelemetryEvent(altitude_m=alt))
            time.sleep(0.15)
        # Сигнал «данных больше не будет»
        events.put(None)
    
    
    def worker(events: queue.Queue, commands: queue.Queue) -> None:
        """Обработчик: читает события и при необходимости ставит команду."""
        while True:
            item = events.get()
            if item is None:
                commands.put(None)
                events.task_done()
                break
            # Учебное правило: выше 10 м — «удерживать», иначе — «набор»
            if item.altitude_m >= 10.0:
                commands.put(Command("HOLD", f"alt={item.altitude_m}"))
            else:
                commands.put(Command("CLIMB", f"alt={item.altitude_m}"))
            events.task_done()
    
    
    def main() -> None:
        events: queue.Queue = queue.Queue()
        commands: queue.Queue = queue.Queue()
        stop_flag = threading.Event()
    
        t_reader = threading.Thread(target=reader, args=(events, stop_flag), daemon=True)
        t_worker = threading.Thread(target=worker, args=(events, commands), daemon=True)
        t_reader.start()
        t_worker.start()
    
        while True:
            cmd = commands.get()
            if cmd is None:
                commands.task_done()
                break
            print(f"command: {cmd.name} ({cmd.detail})")
            commands.task_done()
    
        t_reader.join(timeout=1.0)
        t_worker.join(timeout=1.0)
        print("pipeline done")
    
    
    if __name__ == "__main__":
        main()

Задача 10.3. Две задачи asyncio

Файл практики: k10_task3_asyncio_tick.py.

Запустите telemetry_ticker и status_printer через asyncio.gather и asyncio.run. Убедитесь, что сообщения чередуются, а не строго блоками одного sleep.

Ожидаемый результат. Чередование строк telemetry tick и status; в конце asyncio done и ненулевой ticks.

  • Полный текст файла k10_task3_asyncio_tick.py

    Скопируйте приведённое ниже содержимое в файл code/M01_python/k10_task3_asyncio_tick.py в проекте PyCharm (виртуальное окружение) и запустите программу. Ниже приведён полный учебный текст файла с комментариями.

    # k10_task3_asyncio_tick.py
    # Цель: две кооперативные async-задачи на одном потоке событий.
    # Это не замена всему concurrent-коду курса, а знакомство с asyncio.
    
    import asyncio
    from dataclasses import dataclass
    
    
    @dataclass
    class FakeLink:
        """Учебное «состояние канала»: счётчик принятых тиков."""
    
        ticks: int = 0
    
    
    async def telemetry_ticker(link: FakeLink, n: int) -> None:
        """Имитация периодического приёма: n раз с паузой."""
        for i in range(n):
            await asyncio.sleep(0.1)
            link.ticks += 1
            print(f"telemetry tick {i + 1}, total={link.ticks}")
    
    
    async def status_printer(link: FakeLink, n: int) -> None:
        """Параллельно (кооперативно) печатает сводку n раз."""
        for i in range(n):
            await asyncio.sleep(0.15)
            print(f"status {i + 1}: ticks_seen={link.ticks}")
    
    
    async def main() -> None:
        link = FakeLink()
        # gather запускает обе корутины в одном цикле событий
        await asyncio.gather(
            telemetry_ticker(link, 5),
            status_printer(link, 4),
        )
        print("asyncio done, ticks=", link.ticks)
    
    
    if __name__ == "__main__":
        asyncio.run(main())

Задача 10.4. get с timeout

Файл практики: k10_task4_timeout_queue.py.

Покажите короткий таймаут, который не успевает получить сообщение, и более длинный таймаут, который сообщение получает.

Ожидаемый результат. Сначала сообщение о timeout 0.2s, затем received с HEARTBEAT.

  • Полный текст файла k10_task4_timeout_queue.py

    Скопируйте приведённое ниже содержимое в файл code/M01_python/k10_task4_timeout_queue.py в проекте PyCharm (виртуальное окружение) и запустите программу. Ниже приведён полный учебный текст файла с комментариями.

    # k10_task4_timeout_queue.py
    # Цель: queue.get с timeout — не ждать событие вечно.
    # Задел под «не дождались сообщения MAVLink за N секунд».
    
    import queue
    import threading
    import time
    
    
    def slow_publisher(q: queue.Queue) -> None:
        """Публикует одно сообщение с задержкой."""
        time.sleep(0.4)
        q.put({"type": "HEARTBEAT", "ok": True})
    
    
    def main() -> None:
        q: queue.Queue = queue.Queue()
        threading.Thread(target=slow_publisher, args=(q,), daemon=True).start()
    
        try:
            # Ждём не дольше 0.2 с — для медленного издателя это мало
            msg = q.get(timeout=0.2)
            print("unexpected early message:", msg)
        except queue.Empty:
            print("timeout 0.2s: сообщения ещё нет (это ожидаемо)")
    
        # Вторая попытка с запасом по времени
        try:
            msg = q.get(timeout=0.5)
            print("received:", msg)
        except queue.Empty:
            print("timeout 0.5s: всё ещё пусто (неожиданно для этого сценария)")
    
    
    if __name__ == "__main__":
        main()

Задача 10.5. Сохранение в Git

Сохраните практику в истории Git.

git add code/M01_python/k10_*.py
git commit -m "K10: threading queue asyncio timeout practice"
↑ К оглавлению

11. Проверьте себя

  1. Зачем в программах для БПЛА отделяют приём телеметрии от принятия решений?
  2. Что такое гонка данных и какую роль играет Lock?
  3. Чем Queue удобнее общей переменной без очереди при связи двух потоков?
  4. Что означает маркер None в учебной очереди событий?
  5. Чем async def и await отличаются по смыслу от обычной функции с time.sleep в многозадачном сценарии?
  6. Зачем queue.get вызывают с timeout?
  7. Какую схему вы выберете для ближайших синхронных скриптов на pymavlink: threading+Queue или сразу только asyncio? Почему?
  8. Почему печать Command в учебном примере ещё не означает отправку команды автопилоту?
↑ К оглавлению

12. Чек-лист самопроверки

  • Я объясняю поток, очередь и asyncio законченными предложениями.
  • Я запускал фоновый Thread и читал общий снимок через Lock.
  • Я собрал конвейер Queue «событие → команда».
  • Я запускал две задачи через asyncio.gather.
  • Я обработал queue.Empty при get с timeout.
  • Я понимаю, что учебные Command не являются MAVLink-отправкой.
  • Я открывал аккордеоны практики и сверял полный код с запуском в виртуальном окружении.
  • Я запускал k10_task1 … k10_task4 в venv.
  • Я сделал git commit по K10.
Результат. Когда пункты закрыты, напишите ментору: K10 чек-лист закрыт. Далее следует урок K11 о тестах pytest и заглушках — то есть о проверке логики без реального автопилота.
↑ К оглавлению

13. Что дальше

  • Предыдущий урок: K09 · Классы, dataclasses и Enum
  • Текущий урок: K10 · Потоки, asyncio и очередь сообщений
  • Следующий урок: K11 · Тесты pytest и моки
  • Горизонт модуля M1: K12 и контрольная точка G1
Стиль изложения. Текст следует голосу внимательного преподавателя: законченные предложения, явные определения, терминология БПЛА, без сленга. Подпись теоретического блока dual-code — «Чистый код». Практика содержит полный исходный текст в аккордеоне. Стандарт: docs/05_EDITOR_STYLE_GUIDE.md.
↑ К оглавлению

 

 

Вторник, 01 сентября 2026
k10. Потоки, asyncio и очередь сообщений