Транспорты и Протоколы
Предисловие
Транспорты и Протоколы используются 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), где нижний и верхний — положительные числа байтов.Используйте
set_write_buffer_limits()для установки пределов.Новая в версии 3.4.2.
-
WriteTransport.set_write_buffer_limits(high=None, low=None) -
Установить верхний и нижний пороги для управления потоком записи.
Эти два значения (измеряемые в количестве байтов) определяют, когда вызываются методы протокола
protocol.pause_writing()иprotocol.resume_writing(). Если указано, нижний порог должен быть меньше или равен верхнему порогу. Ни верхний, ни нижний пороги не могут быть отрицательными.Метод
pause_writing()вызывается, когда размер буфера становится больше или равен верхнему значению. Если запись была приостановлена, методresume_writing()вызывается, когда размер буфера становится меньше или равен нижнему значению.Значения по умолчанию зависят от реализации. Если указан только верхний порог, нижний порог устанавливается по умолчанию до значения, зависящего от реализации, меньшего или равного верхнему порогу. Установка верхнего порога в ноль устанавливает нижний порог в ноль и заставляет вызывать
pause_writing()всякий раз, когда буфер становится непустым. Установка нижнего порога в ноль заставляет вызывать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 — это непустой объект типа bytes, содержащий входящие данные.
Буферизация, фрагментация или повторная сборка данных зависят от транспорта. В общем случае вы не должны полагаться на определённые семантики, а вместо этого сделайте свой парсинг универсальным и гибким. Однако данные всегда поступают в правильном порядке.
Метод может вызываться любое количество раз в течение открытого соединения.
Однако,
protocol.eof_received()вызывается не более одного раза. После вызоваeof_received(),data_received()больше не вызывается.
-
Protocol.eof_received() -
Вызывается, когда другой конец сигнализирует о том, что больше не будет отправлять данных (например, вызвав
transport.write_eof(), если другой конец также использует asyncio).Этот метод может возвращать ложное значение (включая
None), в этом случае транспорт закроет себя. И наоборот, если этот метод возвращает истинное значение, используемый протокол определяет, нужно ли закрывать транспорт. Поскольку по умолчанию возвращаетсяNone, это подразумевает закрытие соединения.Некоторые транспорты, включая SSL, не поддерживают полузакрытые соединения, поэтому возврат истинного значения из этого метода приведёт к закрытию соединения.
Машина состояний:
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–2022 Python Software Foundation
Licensed under the PSF License.
https://docs.python.org/3.9/library/asyncio-protocol.html