Spec-Zone.ru › Python 3.7

Транспортные средства и протоколы

Предисловие

Транспортные средства и протоколы используются API событийного цикла низкого уровня, такими как loop.create_connection(). Они используют стиль программирования на основе обратных вызовов и позволяют создавать высокопроизводительные реализации сетевых или IPC-протоколов (например, HTTP).

По сути, транспортные средства и протоколы должны использоваться только в библиотеках и фреймворках, а никогда не в приложениях asyncio высокого уровня.

В данном документе рассматриваются как транспортные средства, так и протоколы.

Введение

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

Другой способ сказать то же самое: транспортное средство — это абстракция сокета (или аналогичного конечного пункта ввода-вывода), а протокол — это абстракция приложения с точки зрения транспортного средства.

Еще один взгляд: интерфейсы транспортного средства и протокола вместе определяют абстрактный интерфейс для использования сетевого ввода-вывода и ввода-вывода между процессами.

Всегда существует взаимно однозначное соответствие между объектами транспортного средства и протокола: протокол вызывает методы транспортного средства для отправки данных, а транспортное средство вызывает методы протокола, чтобы передать ему полученные данные.

Большинство методов событийного цикла для ориентированных на подключение соединений (например, loop.create_connection()) обычно принимают аргумент protocol_factory, используемый для создания объекта Protocol для принятого подключения, представленного объектом Transport. Такие методы обычно возвращают кортеж (transport, protocol).

Содержание

Данный документ содержит следующие разделы:

  • В разделе Транспортные средства описываются классы asyncio BaseTransport, ReadTransport, WriteTransport, Transport, DatagramTransport и SubprocessTransport.
  • В разделе Протоколы описываются классы asyncio BaseProtocol, Protocol, BufferedProtocol, DatagramProtocol и SubprocessProtocol.
  • В разделе Примеры демонстрируется работа с транспортными средствами, протоколами и API событийного цикла низкого уровня.

Транспортные средства

Транспортные средства — это классы, предоставляемые asyncio для абстрагирования различных типов каналов связи.

Объекты транспортных средств всегда создаются событийным циклом asyncio.

asyncio реализует транспортные средства для TCP, UDP, SSL и труб подпроцессов. Доступные методы транспортного средства зависят от его типа.

Классы транспортных средств не потокобезопасны.

Иерархия транспортных средств

class asyncio.BaseTransport

Базовый класс для всех транспортных средств. Содержит методы, общие для всех транспортных средств asyncio.

class asyncio.WriteTransport(BaseTransport)

Базовое транспортное средство для односторонних подключений на запись.

Экземпляры класса WriteTransport возвращаются методом loop.connect_write_pipe() событийного цикла, а также используются методами, связанными с подпроцессами, такими как loop.subprocess_exec().

class asyncio.ReadTransport(BaseTransport)

Базовое транспортное средство для односторонних подключений на чтение.

Экземпляры класса ReadTransport возвращаются методом loop.connect_read_pipe() событийного цикла, а также используются методами, связанными с подпроцессами, такими как loop.subprocess_exec().

class asyncio.Transport(WriteTransport, ReadTransport)

Интерфейс, представляющий двустороннее транспортное средство, такое как TCP-соединение.

Пользователь не создаёт транспортное средство напрямую; они вызывают вспомогательную функцию, передавая её фабрику протокола и другую информацию, необходимую для создания транспортного средства и протокола.

Экземпляры класса Transport возвращаются или используются методами событийного цикла, такими как loop.create_connection(), loop.create_unix_connection(), loop.create_server(), loop.sendfile() и т. д.

class asyncio.DatagramTransport(BaseTransport)

Транспортное средство для подключений датаграмм (UDP).

Экземпляры класса DatagramTransport возвращаются методом loop.create_datagram_endpoint() событийного цикла.

class asyncio.SubprocessTransport(BaseTransport)

Абстракция для представления соединения между родительским и дочерним процессами ОС.

Экземпляры класса SubprocessTransport возвращаются методами событийного цикла loop.subprocess_shell() и loop.subprocess_exec().

Базовое транспортное средство

BaseTransport.close()

Закрыть транспортное средство.

Если у транспортного средства есть буфер для исходящих данных, буферизованные данные будут отправлены асинхронно. Больше данных не будет получено. После отправки всех буферизованных данных будет вызван метод протокола protocol.connection_lost() с None в качестве аргумента.

BaseTransport.is_closing()

Возвращает True , если транспортное средство закрывается или закрыто.

BaseTransport.get_extra_info(name, default=None)

Возвращает информацию о транспорте или используемых им базовых ресурсах.

name — строка, представляющая запрашиваемую информацию, специфичную для транспорта.

default — значение, возвращаемое, если информация недоступна или транспорт не поддерживает её запрос с использованием данного стороннего цикла событий или на текущей платформе.

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

sock = transport.get_extra_info('socket')
if sock is not None:
    print(sock.getsockopt(...))

Категории информации, которые могут быть запрошены для некоторых транспортов:

  • сокет:

    • 'peername': удалённый адрес, к которому подключён сокет, результат socket.socket.getpeername() (None при ошибке)
    • 'socket': экземпляр socket.socket
    • 'sockname': собственный адрес сокета, результат socket.socket.getsockname()
  • сокет SSL:

    • 'compression': алгоритм сжатия в виде строки или None если подключение не сжато; результат ssl.SSLSocket.compression()
    • 'cipher': тройка значений, содержащая название используемого шифра, версию протокола SSL, определяющего его использование, и количество используемых секретных битов; результат ssl.SSLSocket.cipher()
    • 'peercert': сертификат удалённого узла; результат ssl.SSLSocket.getpeercert()
    • 'sslcontext': экземпляр ssl.SSLContext
    • 'ssl_object': экземпляр ssl.SSLObject или ssl.SSLSocket
  • пайп:

    • 'pipe': объект пайпа
  • подпроцесс:

    • 'subprocess': экземпляр subprocess.Popen
BaseTransport.set_protocol(protocol)

Устанавливает новый протокол.

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

BaseTransport.get_protocol()

Возвращает текущий протокол.

Только чтение Транспорты

ReadTransport.is_reading()

Возвращает True, если транспорт получает новые данные.

Добавлен в версии 3.7.

ReadTransport.pause_reading()

Приостанавливает приём данных по транспорту. Данные не будут передаваться методу протокола protocol.data_received(), пока не будет вызван resume_reading().

Изменено в версии 3.7: Метод идемпотентен, т.е. может вызываться, когда транспорт уже приостановлен или закрыт.

ReadTransport.resume_reading()

Возобновляет приём. Метод протокола protocol.data_received() будет вызван снова, если доступны данные для чтения.

Изменено в версии 3.7: Метод идемпотентен, т.е. может вызываться, когда транспорт уже читает.

Только запись Транспорты

WriteTransport.abort()

Немедленно закрывает транспорт, не дожидаясь завершения текущих операций. Буферизованные данные будут потеряны. Больше данных получено не будет. Метод протокола protocol.connection_lost() в конечном итоге будет вызван с аргументом None.

WriteTransport.can_write_eof()

Возвращает True, если транспорт поддерживает write_eof(), и False в противном случае.

WriteTransport.get_write_buffer_size()

Возвращает текущий размер буфера вывода, используемого транспортом.

WriteTransport.get_write_buffer_limits()

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

Используйте set_write_buffer_limits() для установки ограничений.

Добавлен в версии 3.4.2.

WriteTransport.set_write_buffer_limits(high=None, low=None)

Устанавливает верхнюю и нижнюю отметки для управления потоком записи.

Эти два значения (измеряемые в количестве байтов) управляют тем, когда вызываются методы протокола protocol.pause_writing() и protocol.resume_writing(). Если указаны, нижняя отметка должна быть меньше или равна верхней отметке. Ни high, ни low не могут быть отрицательными.

pause_writing() вызывается, когда размер буфера становится больше или равен значению high. Если запись приостановлена, resume_writing() вызывается, когда размер буфера становится меньше или равен значению low.

Значения по умолчанию зависят от реализации. Если задана только верхняя отметка, нижняя отметка устанавливается по умолчанию до реализации-специфического значения, меньшего или равного верхней отметке. Установка high в ноль устанавливает low в ноль и вызывает pause_writing() всякий раз, когда буфер становится непустым. Установка low в ноль вызывает resume_writing() только когда буфер пустой. Использование нуля для любого из ограничений обычно неэффективно, так как это уменьшает возможности одновременного выполнения ввода-вывода и вычислений.

Используйте get_write_buffer_limits() для получения ограничений.

WriteTransport.write(data)

Записывает некоторые байты data в транспорт.

Этот метод не блокируется; он буферизует данные и организует их асинхронную отправку.

WriteTransport.writelines(list_of_data)

Записывает список (или любой итерируемый объект) байтов данных в транспорт. Это функционально эквивалентно вызову write() для каждого элемента, полученного из итерируемого объекта, но может быть реализовано более эффективно.

WriteTransport.write_eof()

Закрывает конец записи транспорта после сброса всех буферизованных данных. Данные могут всё ещё получаться.

Этот метод может вызвать NotImplementedError, если транспорт (например, SSL) не поддерживает полузакрытые подключения.

Транспорты датаграмм

DatagramTransport.sendto(data, addr=None)

Отправляет байты data удалённому узлу, указанному параметром addr (адрес назначения, зависящий от транспорта). Если addr — None, данные отправляются по адресу назначения, заданному при создании транспорта.

Этот метод не блокируется; он буферизует данные и организует их асинхронную отправку.

DatagramTransport.abort()

Немедленно закрывает транспорт, не дожидаясь завершения текущих операций. Буферизованные данные будут потеряны. Больше данных получено не будет. Метод протокола protocol.connection_lost() в конечном итоге будет вызван с аргументом None.

Транспорты подпроцессов

SubprocessTransport.get_pid()

Возвращает идентификатор процесса подпроцесса в виде целого числа.

SubprocessTransport.get_pipe_transport(fd)

Возвращает транспорт для канала связи, соответствующий целочисленному дескриптору файла fd:

  • 0: транспорт потокового чтения стандартного ввода (stdin), или None, если подпроцесс не был создан с stdin=PIPE
  • 1: транспорт потокового записи стандартного вывода (stdout), или None, если подпроцесс не был создан с stdout=PIPE
  • 2: транспорт потокового записи стандартной ошибки (stderr), или None, если подпроцесс не был создан с stderr=PIPE
  • другие fd: None
SubprocessTransport.get_returncode()

Возвращает код завершения подпроцесса как целое число или None, если он ещё не завершился, что аналогично атрибуту subprocess.Popen.returncode.

SubprocessTransport.kill()

Завершает подпроцесс.

В системах POSIX функция отправляет SIGKILL подпроцессу. В Windows этот метод является алиасом для terminate().

См. также subprocess.Popen.kill().

SubprocessTransport.send_signal(signal)

Отправляет номер сигнала подпроцессу, как в subprocess.Popen.send_signal().

SubprocessTransport.terminate()

Останавливает подпроцесс.

В системах POSIX этот метод отправляет SIGTERM подпроцессу. В Windows вызывается функция Windows API TerminateProcess() для остановки подпроцесса.

См. также subprocess.Popen.terminate().

SubprocessTransport.close()

Завершает подпроцесс, вызвав метод kill().

Если подпроцесс ещё не завершился, и закрывает транспортные средства stdin, stdout и stderr каналов.

Протоколы

asyncio предоставляет набор абстрактных базовых классов, которые следует использовать для реализации сетевых протоколов. Эти классы предназначены для совместного использования с транспортами.

Подклассы абстрактных базовых классов протокола могут реализовывать некоторые или все методы. Все эти методы являются обратными вызовами: они вызываются транспортом при определённых событиях, например, при получении данных. Метод базового протокола должен вызываться соответствующим транспортом.

Базовые Протоколы

class asyncio.BaseProtocol

Базовый протокол со методами, общими для всех протоколов.

class asyncio.Protocol(BaseProtocol)

Базовый класс для реализации потоковых протоколов (TCP, сокеты Unix и т. д.).

class asyncio.BufferedProtocol(BaseProtocol)

Базовый класс для реализации потоковых протоколов с ручным управлением буфером приёма.

class asyncio.DatagramProtocol(BaseProtocol)

Базовый класс для реализации протоколов датаграмм (UDP).

class asyncio.SubprocessProtocol(BaseProtocol)

Базовый класс для реализации протоколов, взаимодействующих с дочерними процессами (однонаправленные каналы).

Базовый Протокол

Все протоколы asyncio могут реализовывать обратные вызовы Базового Протокола.

Обратные вызовы подключения

Обратные вызовы подключения вызываются на всех протоколах ровно один раз при успешном подключении. Все остальные обратные вызовы протокола могут быть вызваны только между этими двумя методами.

BaseProtocol.connection_made(transport)

Вызывается при установлении подключения.

Аргумент transport — это транспорт, представляющий подключение. Протокол отвечает за хранение ссылки на свой транспорт.

BaseProtocol.connection_lost(exc)

Вызывается при потере или закрытии подключения.

Аргументом является либо объект исключения, либо None. Последнее означает, что получен стандартный EOF, или подключение было прервано или закрыто этой стороной соединения.

Обратные вызовы управления потоком

Обратные вызовы управления потоком могут вызываться транспортом для приостановки или возобновления записи, выполняемой протоколом.

См. документацию метода set_write_buffer_limits() для получения более подробной информации.

BaseProtocol.pause_writing()

Вызывается, когда буфер транспорта превышает верхний порог.

BaseProtocol.resume_writing()

Вызывается, когда буфер транспорта падает ниже нижнего порога.

Если размер буфера равен верхнему порогу, pause_writing() не вызывается: размер буфера должен строго превышать значение.

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

Потоковые протоколы

Методы событий, такие как loop.create_server(), loop.create_unix_server(), loop.create_connection(), loop.create_unix_connection(), loop.connect_accepted_socket(), loop.connect_read_pipe() и loop.connect_write_pipe() принимают фабрики, возвращающие потоковые протоколы.

Protocol.data_received(data)

Вызывается при получении данных. data — это непустой байтовый объект, содержащий входящие данные.

Буферизация, фрагментация или повторная сборка данных зависят от транспорта. В целом, вы не должны полагаться на конкретную семантику, а вместо этого сделайте свой парсинг универсальным и гибким. Однако данные всегда поступают в правильном порядке.

Метод может вызываться произвольное число раз во время открытого соединения.

Однако protocol.eof_received() вызывается не более одного раза. После вызова eof_received() вызов data_received() больше не вызывается.

Protocol.eof_received()

Вызывается, когда другая сторона сигнализирует о том, что больше не будет отправлять данных (например, вызвав transport.write_eof(), если другая сторона также использует asyncio).

Этот метод может возвращать ложное значение (включая None), в этом случае транспорт закроет себя. Напротив, если этот метод возвращает истинное значение, протокол определяет, нужно ли закрывать транспорт. Поскольку по умолчанию реализация возвращает None, она неявно закрывает соединение.

Некоторые транспорты, включая SSL, не поддерживают полузакрытые соединения, поэтому возвращение значения true из этого метода приведет к закрытию соединения.

Схему состояния:

start -> connection_made
    [-> data_received]*
    [-> eof_received]?
-> connection_lost -> end

Буферизованные потоковые протоколы

New in version 3.7: Важно: это было добавлено в asyncio в Python 3.7 на временной основе! Это экспериментальный API, который может быть изменён или полностью удалён в Python 3.8.

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

BufferedProtocol реализации позволяют явно выделять и контролировать буфер приёма. Циклы событий могут затем использовать буфер, предоставляемый протоколом, чтобы избежать ненужных копий данных. Это может привести к заметному повышению производительности для протоколов, которые принимают большие объёмы данных. Сложные реализации протоколов могут значительно уменьшить количество выделений буферов.

Следующие обратные вызовы вызываются на объектах BufferedProtocol:

BufferedProtocol.get_buffer(sizehint)

Вызывается для выделения нового буфера приема.

sizehint — рекомендуемый минимальный размер возвращаемого буфера. Допустимо возвращать буферы меньшего или большего размера, чем предполагает sizehint. При значении -1 размер буфера может быть произвольным. Возвращение буфера нулевого размера является ошибкой.

get_buffer() должен вернуть объект, реализующий протокол буфера.

BufferedProtocol.buffer_updated(nbytes)

Вызывается, когда буфер был обновлен полученными данными.

nbytes — общее количество байтов, записанных в буфер.

BufferedProtocol.eof_received()

См. документацию метода protocol.eof_received().

get_buffer() может быть вызван произвольное количество раз во время подключения. Однако protocol.eof_received() вызывается не более одного раза, и, если он был вызван, get_buffer() и buffer_updated() больше не будут вызываться после него.

Схема работы:

start -> connection_made
    [-> get_buffer
        [-> buffer_updated]?
    ]*
    [-> eof_received]?
-> connection_lost -> end

Протоколы дейтаграмм

Экземпляры протокола дейтаграмм должны создаваться фабриками протоколов, передаваемыми методу loop.create_datagram_endpoint().

DatagramProtocol.datagram_received(data, addr)

Вызывается при получении дейтаграммы. data — объект bytes, содержащий входящие данные. addr — адрес отправителя данных; точный формат зависит от транспорта.

DatagramProtocol.error_received(exc)

Вызывается, когда предыдущая операция отправки или получения вызывает исключение OSError. exc — экземпляр OSError.

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

Примечание

В системах BSD (macOS, FreeBSD и т. д.) управление потоком для протоколов дейтаграмм не поддерживается, потому что нет надежного способа обнаружить ошибки отправки, вызванные отправкой слишком большого количества пакетов.

Сокет всегда выглядит «готовым», и лишние пакеты отбрасываются. Исключение OSError с errno, установленным в errno.ENOBUFS, может быть или не быть возбуждено; если оно возбуждается, оно будет передано методу DatagramProtocol.error_received(), в противном случае будет проигнорировано.

Протоколы подпроцессов

Экземпляры протокола дейтаграмм должны создаваться фабриками протоколов, передаваемыми методам loop.subprocess_exec() и loop.subprocess_shell().

SubprocessProtocol.pipe_data_received(fd, data)

Вызывается, когда дочерний процесс записывает данные в свой канал stdout или stderr.

fd — целочисленный дескриптор файла канала.

data — непустой объект bytes, содержащий полученные данные.

SubprocessProtocol.pipe_connection_lost(fd, exc)

Вызывается, когда один из каналов, общающихся с дочерним процессом, закрыт.

fd — целочисленный дескриптор файла, который был закрыт.

SubprocessProtocol.process_exited()

Вызывается, когда дочерний процесс завершил работу.

Примеры

Сервер TCP-эха

Создайте сервер TCP-эха, используя метод loop.create_server(), отправьте обратно полученные данные и закройте подключение:

import asyncio


class EchoServerProtocol(asyncio.Protocol):
    def connection_made(self, transport):
        peername = transport.get_extra_info('peername')
        print('Connection from {}'.format(peername))
        self.transport = transport

    def data_received(self, data):
        message = data.decode()
        print('Data received: {!r}'.format(message))

        print('Send: {!r}'.format(message))
        self.transport.write(data)

        print('Close the client socket')
        self.transport.close()


async def main():
    # Get a reference to the event loop as we plan to use
    # low-level APIs.
    loop = asyncio.get_running_loop()

    server = await loop.create_server(
        lambda: EchoServerProtocol(),
        '127.0.0.1', 8888)

    async with server:
        await server.serve_forever()


asyncio.run(main())

См. также

Пример сервера TCP-эха, использующего потоки использует высокоуровневую функцию asyncio.start_server().

Клиент TCP-эха

Клиент TCP-эха, использующий метод loop.create_connection(), отправляет данные и ждет закрытия соединения:

import asyncio


class EchoClientProtocol(asyncio.Protocol):
    def __init__(self, message, on_con_lost):
        self.message = message
        self.on_con_lost = on_con_lost

    def connection_made(self, transport):
        transport.write(self.message.encode())
        print('Data sent: {!r}'.format(self.message))

    def data_received(self, data):
        print('Data received: {!r}'.format(data.decode()))

    def connection_lost(self, exc):
        print('The server closed the connection')
        self.on_con_lost.set_result(True)


async def main():
    # Get a reference to the event loop as we plan to use
    # low-level APIs.
    loop = asyncio.get_running_loop()

    on_con_lost = loop.create_future()
    message = 'Hello World!'

    transport, protocol = await loop.create_connection(
        lambda: EchoClientProtocol(message, on_con_lost),
        '127.0.0.1', 8888)

    # Wait until the protocol signals that the connection
    # is lost and close the transport.
    try:
        await on_con_lost
    finally:
        transport.close()


asyncio.run(main())

См. также

Пример клиента TCP-эха, использующего потоки использует высокоуровневую функцию asyncio.open_connection().

Сервер UDP-эха

Сервер UDP-эха, использующий метод loop.create_datagram_endpoint(), отправляет обратно полученные данные:

import asyncio


class EchoServerProtocol:
    def connection_made(self, transport):
        self.transport = transport

    def datagram_received(self, data, addr):
        message = data.decode()
        print('Received %r from %s' % (message, addr))
        print('Send %r to %s' % (message, addr))
        self.transport.sendto(data, addr)


async def main():
    print("Starting UDP server")

    # Get a reference to the event loop as we plan to use
    # low-level APIs.
    loop = asyncio.get_running_loop()

    # One protocol instance will be created to serve all
    # client requests.
    transport, protocol = await loop.create_datagram_endpoint(
        lambda: EchoServerProtocol(),
        local_addr=('127.0.0.1', 9999))

    try:
        await asyncio.sleep(3600)  # Serve for 1 hour.
    finally:
        transport.close()


asyncio.run(main())

Клиент UDP-эха

Клиент UDP-эха, использующий метод loop.create_datagram_endpoint(), отправляет данные и закрывает транспорт при получении ответа:

import asyncio


class EchoClientProtocol:
    def __init__(self, message, on_con_lost):
        self.message = message
        self.on_con_lost = on_con_lost
        self.transport = None

    def connection_made(self, transport):
        self.transport = transport
        print('Send:', self.message)
        self.transport.sendto(self.message.encode())

    def datagram_received(self, data, addr):
        print("Received:", data.decode())

        print("Close the socket")
        self.transport.close()

    def error_received(self, exc):
        print('Error received:', exc)

    def connection_lost(self, exc):
        print("Connection closed")
        self.on_con_lost.set_result(True)


async def main():
    # Get a reference to the event loop as we plan to use
    # low-level APIs.
    loop = asyncio.get_running_loop()

    on_con_lost = loop.create_future()
    message = "Hello World!"

    transport, protocol = await loop.create_datagram_endpoint(
        lambda: EchoClientProtocol(message, on_con_lost),
        remote_addr=('127.0.0.1', 9999))

    try:
        await on_con_lost
    finally:
        transport.close()


asyncio.run(main())

Подключение существующих сокетов

Ожидание получения данных сокетом, используя метод loop.create_connection() с протоколом:

import asyncio
import socket


class MyProtocol(asyncio.Protocol):

    def __init__(self, on_con_lost):
        self.transport = None
        self.on_con_lost = on_con_lost

    def connection_made(self, transport):
        self.transport = transport

    def data_received(self, data):
        print("Received:", data.decode())

        # We are done: close the transport;
        # connection_lost() will be called automatically.
        self.transport.close()

    def connection_lost(self, exc):
        # The socket has been closed
        self.on_con_lost.set_result(True)


async def main():
    # Get a reference to the event loop as we plan to use
    # low-level APIs.
    loop = asyncio.get_running_loop()
    on_con_lost = loop.create_future()

    # Create a pair of connected sockets
    rsock, wsock = socket.socketpair()

    # Register the socket to wait for data.
    transport, protocol = await loop.create_connection(
        lambda: MyProtocol(on_con_lost), sock=rsock)

    # Simulate the reception of data from the network.
    loop.call_soon(wsock.send, 'abc'.encode())

    try:
        await protocol.on_con_lost
    finally:
        transport.close()
        wsock.close()

asyncio.run(main())

См. также

Пример отслеживание дескриптора файла для событий чтения использует низкоуровневый метод loop.add_reader() для регистрации FD.

Пример регистрации открытого сокета для ожидания данных, используя потоки использует высокоуровневые потоки, созданные функцией open_connection() в корутине.

loop.subprocess_exec() и SubprocessProtocol

Пример использования протокола подпроцесса для получения вывода подпроцесса и ожидания завершения подпроцесса.

Подпроцесс создается методом loop.subprocess_exec():

import asyncio
import sys

class DateProtocol(asyncio.SubprocessProtocol):
    def __init__(self, exit_future):
        self.exit_future = exit_future
        self.output = bytearray()

    def pipe_data_received(self, fd, data):
        self.output.extend(data)

    def process_exited(self):
        self.exit_future.set_result(True)

async def get_date():
    # Get a reference to the event loop as we plan to use
    # low-level APIs.
    loop = asyncio.get_running_loop()

    code = 'import datetime; print(datetime.datetime.now())'
    exit_future = asyncio.Future(loop=loop)

    # Create the subprocess controlled by DateProtocol;
    # redirect the standard output into a pipe.
    transport, protocol = await loop.subprocess_exec(
        lambda: DateProtocol(exit_future),
        sys.executable, '-c', code,
        stdin=None, stderr=None)

    # Wait for the subprocess exit using the process_exited()
    # method of the protocol.
    await exit_future

    # Close the stdout pipe.
    transport.close()

    # Read the output which was collected by the
    # pipe_data_received() method of the protocol.
    data = bytes(protocol.output)
    return data.decode('ascii').rstrip()

if sys.platform == "win32":
    asyncio.set_event_loop_policy(
        asyncio.WindowsProactorEventLoopPolicy())

date = asyncio.run(get_date())
print(f"Current date: {date}")

См. также аналогичный пример, написанный с использованием высокоуровневых API.

© 2001–2020 Python Software Foundation
Licensed under the PSF License.
https://docs.python.org/3.7/library/asyncio-protocol.html

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API