Spec-Zone.ru › Python 3.10

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

Предисловие

Транспортные средства и Протоколы используются API событийного цикла нижележащего уровня, такими как loop.create_connection(). Они используют стиль программирования на основе обратного вызова и позволяют создавать высокопроизводительные реализации сетевых или межпроцессных протоколов (например, 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 событийного цикла низкого уровня.

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

Исходный код: Lib/asyncio/transports.py

Транспортные средства — это классы, предоставляемые 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.

Протоколы

Исходный код: Lib/asyncio/protocols.py

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).

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

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

Автомат состояний:

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

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

Введено в версии 3.7.

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

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()

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

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

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

Spec-Zone.ru

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