Потоки
Исходный код: Lib/asyncio/streams.py
Потоки — это высокоуровневые примитивы, готовые к использованию с async/await, для работы с сетевыми соединениями. Потоки позволяют отправлять и получать данные без использования обратных вызовов, низкоуровневых протоколов и транспортов.
Пример TCP-клиента эхо-сервера, написанного с использованием потоков asyncio:
import asyncio
async def tcp_echo_client(message):
reader, writer = await asyncio.open_connection(
'127.0.0.1', 8888)
print(f'Send: {message!r}')
writer.write(message.encode())
await writer.drain()
data = await reader.read(100)
print(f'Received: {data.decode()!r}')
print('Close the connection')
writer.close()
await writer.wait_closed()
asyncio.run(tcp_echo_client('Hello World!'))
См. также раздел Примеры ниже.
Функции потоков
Для создания потоков и работы с ними можно использовать следующие функции asyncio верхнего уровня:
-
async asyncio.open_connection(host=None, port=None, *, limit=65536, ssl=None, family=0, proto=0, flags=0, sock=None, local_addr=None, server_hostname=None, ssl_handshake_timeout=None, ssl_shutdown_timeout=None, happy_eyeballs_delay=None, interleave=None) -
Устанавливает сетевое соединение и возвращает пару объектов
(reader, writer).Возвращаемые объекты reader и writer являются экземплярами классов
StreamReaderиStreamWriter.limit определяет ограничение размера буфера для возвращаемого экземпляра
StreamReader. По умолчанию limit установлен в 64 КиБ.Остальные аргументы передаются непосредственно в
loop.create_connection().Примечание
Аргумент sock передаёт владение сокетом созданному объекту
StreamWriter. Чтобы закрыть сокет, вызовите его методclose().Изменено в версии 3.7: Добавлен параметр ssl_handshake_timeout.
Изменено в версии 3.8: Добавлены параметры happy_eyeballs_delay и interleave.
Изменено в версии 3.10: Удалён параметр loop.
Изменено в версии 3.11: Добавлен параметр ssl_shutdown_timeout.
-
async asyncio.start_server(client_connected_cb, host=None, port=None, *, limit=65536, family=socket.AF_UNSPEC, flags=socket.AI_PASSIVE, sock=None, backlog=100, ssl=None, reuse_address=None, reuse_port=None, keep_alive=None, ssl_handshake_timeout=None, ssl_shutdown_timeout=None, start_serving=True) -
Запускает сервер сокетов.
Функция обратного вызова client_connected_cb вызывается при установлении нового клиентского соединения. Она получает пару
(reader, writer)в качестве двух аргументов — экземпляры классовStreamReaderиStreamWriter.client_connected_cb может быть обычным вызываемым объектом или функцией-корутиной; если это функция-корутина, она будет автоматически запланирована как
Task.limit определяет ограничение размера буфера для возвращаемого экземпляра
StreamReader. По умолчанию limit установлен в 64 КиБ.Остальные аргументы передаются непосредственно в
loop.create_server().Примечание
Аргумент sock передаёт владение сокетом созданному серверу. Чтобы закрыть сокет, вызовите метод
close()сервера.Изменено в версии 3.7: Добавлены параметры ssl_handshake_timeout и start_serving.
Изменено в версии 3.10: Удалён параметр loop.
Изменено в версии 3.11: Добавлен параметр ssl_shutdown_timeout.
Изменено в версии 3.13: Добавлен параметр keep_alive.
Сокеты Unix
-
async asyncio.open_unix_connection(path=None, *, limit=65536, ssl=None, sock=None, server_hostname=None, ssl_handshake_timeout=None, ssl_shutdown_timeout=None) -
Устанавливает соединение через сокет Unix и возвращает пару
(reader, writer).Аналогична
open_connection(), но работает с сокетами Unix.См. также документацию по
loop.create_unix_connection().Примечание
Аргумент sock передаёт владение сокетом созданному объекту
StreamWriter. Чтобы закрыть сокет, вызовите его методclose().Доступность: Unix.
Изменено в версии 3.7: Добавлен параметр ssl_handshake_timeout. Теперь параметр path может быть объектом, подобным пути.
Изменено в версии 3.10: Удалён параметр loop.
Изменено в версии 3.11: Добавлен параметр ssl_shutdown_timeout.
-
async asyncio.start_unix_server(client_connected_cb, path=None, *, limit=65536, sock=None, backlog=100, ssl=None, ssl_handshake_timeout=None, ssl_shutdown_timeout=None, start_serving=True, cleanup_socket=True) -
Запускает сервер сокетов Unix.
Аналогична
start_server(), но работает с сокетами Unix.Если cleanup_socket имеет значение true, сокет Unix будет автоматически удалён из файловой системы при закрытии сервера, если только сокет не был заменён после создания сервера.
См. также документацию по
loop.create_unix_server().Примечание
Аргумент sock передаёт владение сокетом созданному серверу. Чтобы закрыть сокет, вызовите метод
close()сервера.Доступность: Unix.
Изменено в версии 3.7: Добавлены параметры ssl_handshake_timeout и start_serving. Теперь параметр path может быть объектом, подобным пути.
Изменено в версии 3.10: Удалён параметр loop.
Изменено в версии 3.11: Добавлен параметр ssl_shutdown_timeout.
Изменено в версии 3.13: Добавлен параметр cleanup_socket.
StreamReader
-
class asyncio.StreamReader -
Представляет объект для чтения, предоставляющий API для чтения данных из потока ввода-вывода. Будучи асинхронным итерируемым объектом, он поддерживает инструкцию
async for.Не рекомендуется создавать объекты StreamReader напрямую; вместо этого используйте
open_connection()иstart_server().-
feed_eof() -
Сообщает о достижении EOF.
-
async read(n=-1) -
Читает из потока до n байт.
Если n не задан или равен
-1, читает до EOF, а затем возвращает все прочитанные данные типаbytes. Если получен EOF и внутренний буфер пуст, возвращает пустой объектbytes.Если n равен
0, немедленно возвращает пустой объектbytes.Если n положительно, возвращает не более n доступных значений
bytes, как только во внутреннем буфере появляется хотя бы 1 байт. Если EOF получен до чтения хотя бы одного байта, возвращает пустой объектbytes.
-
async readline() -
Читает одну строку, где «строка» — это последовательность байтов, заканчивающаяся на
\n.Если получен EOF, а
\nне найден, метод возвращает частично прочитанные данные.Если получен EOF и внутренний буфер пуст, возвращает пустой объект
bytes.
-
async readexactly(n) -
Читает ровно n байт.
Вызывает исключение
IncompleteReadError, если достигнут EOF до чтения n байт. Используйте атрибутIncompleteReadError.partial, чтобы получить частично прочитанные данные.
-
async readuntil(separator=b'\n') -
Читает данные из потока, пока не будет найден separator.
При успешном выполнении данные и разделитель удаляются из внутреннего буфера (потребляются). Возвращаемые данные включают разделитель в конце.
Если объём прочитанных данных превышает настроенное ограничение потока, вызывается исключение
LimitOverrunError, а данные остаются во внутреннем буфере и могут быть прочитаны повторно.Если EOF достигнут до обнаружения полного разделителя, вызывается исключение
IncompleteReadError, а внутренний буфер сбрасывается. АтрибутIncompleteReadError.partialможет содержать часть разделителя.separator также может быть кортежем разделителей. В этом случае возвращаемое значение будет самым коротким из возможных и будет оканчиваться на любой из разделителей. Для целей исключения
LimitOverrunErrorсовпавшим считается самый короткий из возможных разделителей.Добавлено в версии 3.5.2.
Изменено в версии 3.13: Теперь параметр separator может быть
tupleразделителей.
-
at_eof() -
Возвращает
True, если буфер пуст и был вызванfeed_eof().
-
StreamWriter
-
class asyncio.StreamWriter -
Представляет объект для записи, предоставляющий API для записи данных в поток ввода-вывода.
Не рекомендуется создавать объекты StreamWriter напрямую; вместо этого используйте
open_connection()иstart_server().-
write(data) -
Метод пытается немедленно записать data в базовый сокет. Если это не удаётся, данные помещаются во внутренний буфер записи до тех пор, пока их не удастся отправить.
Буфер data должен быть объектом bytes, bytearray или одномерным объектом memoryview с непрерывной областью памяти в стиле C.
Метод следует использовать вместе с методом
drain():stream.write(data) await stream.drain()
-
writelines(data) -
Метод немедленно записывает список (или любой итерируемый объект) байтов в базовый сокет. Если это не удаётся, данные помещаются во внутренний буфер записи до тех пор, пока их не удастся отправить.
Метод следует использовать вместе с методом
drain():stream.writelines(lines) await stream.drain()
-
close() -
Метод закрывает поток и базовый сокет.
Метод следует использовать вместе с методом
wait_closed(), хотя это и не обязательно:stream.close() await stream.wait_closed()
-
can_write_eof() -
Возвращает
True, если базовый транспорт поддерживает методwrite_eof(), иFalseв противном случае.
-
write_eof() -
Закрывает конец потока для записи после сброса буферизованных данных записи.
-
transport -
Возвращает базовый транспорт asyncio.
-
get_extra_info(name, default=None) -
Предоставляет доступ к необязательной информации о транспорте; подробности см. в
BaseTransport.get_extra_info().
-
async drain() -
Ожидает момента, когда можно будет возобновить запись в поток. Пример:
writer.write(data) await writer.drain()
Это метод управления потоком, взаимодействующий с базовым буфером записи ввода-вывода. Когда размер буфера достигает верхнего порогового значения, drain() блокируется, пока размер буфера не уменьшится до нижнего порогового значения и запись не станет возможной. Если ожидать нечего,
drain()немедленно возвращает управление.Примечание
Если размер буфера записи меньше верхнего порогового значения,
drain()возвращает управление немедленно, не передавая управление циклу событий. В результате код, который многократно вызываетwrite(), а затемawait drain(), может препятствовать выполнению других задач. Чтобы избежать блокировки, явно передайте управление циклу событий с помощьюawait asyncio.sleep(0)(см.asyncio.sleep()).
-
async start_tls(sslcontext, *, server_hostname=None, ssl_handshake_timeout=None, ssl_shutdown_timeout=None) -
Обновляет существующее соединение на основе потока до TLS.
Параметры:
-
sslcontext: настроенный экземпляр
SSLContext. - server_hostname: задаёт или переопределяет имя хоста, с которым будет сравниваться сертификат целевого сервера.
-
ssl_handshake_timeout — время в секундах ожидания завершения TLS-рукопожатия, после которого соединение прерывается.
60.0секунд, если значение равноNone(по умолчанию). -
ssl_shutdown_timeout — время в секундах ожидания завершения отключения SSL, после которого соединение прерывается.
30.0секунд, если значение равноNone(по умолчанию).
Добавлено в версии 3.11.
Изменено в версии 3.12: Добавлен параметр ssl_shutdown_timeout.
-
sslcontext: настроенный экземпляр
-
is_closing() -
Возвращает
True, если поток закрыт или находится в процессе закрытия.Добавлено в версии 3.7.
-
async wait_closed() -
Ожидает закрытия потока.
Этот метод следует вызывать после
close(), чтобы дождаться закрытия базового соединения и убедиться, что все данные отправлены, прежде чем, например, завершить программу.Добавлено в версии 3.7.
-
Примеры
TCP-клиент эхо-сервера с использованием потоков
TCP-клиент эхо-сервера с использованием функции asyncio.open_connection():
import asyncio
async def tcp_echo_client(message):
reader, writer = await asyncio.open_connection(
'127.0.0.1', 8888)
print(f'Send: {message!r}')
writer.write(message.encode())
await writer.drain()
data = await reader.read(100)
print(f'Received: {data.decode()!r}')
print('Close the connection')
writer.close()
await writer.wait_closed()
asyncio.run(tcp_echo_client('Hello World!'))
См. также
В примере протокол TCP-клиента эхо-сервера используется низкоуровневый метод loop.create_connection().
TCP-сервер эхо-сервера с использованием потоков
TCP-сервер эхо-сервера с использованием функции asyncio.start_server():
import asyncio
async def handle_echo(reader, writer):
data = await reader.read(100)
message = data.decode()
addr = writer.get_extra_info('peername')
print(f"Received {message!r} from {addr!r}")
print(f"Send: {message!r}")
writer.write(data)
await writer.drain()
print("Close the connection")
writer.close()
await writer.wait_closed()
async def main():
server = await asyncio.start_server(
handle_echo, '127.0.0.1', 8888)
addrs = ', '.join(str(sock.getsockname()) for sock in server.sockets)
print(f'Serving on {addrs}')
async with server:
await server.serve_forever()
asyncio.run(main())
См. также
В примере протокол TCP-сервера эхо-сервера используется метод loop.create_server().
Получение заголовков HTTP
Простой пример получения заголовков HTTP для URL, переданного в командной строке:
import asyncio
import urllib.parse
import sys
async def print_http_headers(url):
url = urllib.parse.urlsplit(url)
if url.scheme == 'https':
reader, writer = await asyncio.open_connection(
url.hostname, 443, ssl=True)
else:
reader, writer = await asyncio.open_connection(
url.hostname, 80)
query = (
f"HEAD {url.path or '/'} HTTP/1.0\r\n"
f"Host: {url.hostname}\r\n"
f"\r\n"
)
writer.write(query.encode('latin-1'))
while True:
line = await reader.readline()
if not line:
break
line = line.decode('latin1').rstrip()
if line:
print(f'HTTP header> {line}')
# Ignore the body, close the socket
writer.close()
await writer.wait_closed()
url = sys.argv[1]
asyncio.run(print_http_headers(url))
Использование:
python example.py http://example.com/path/page.html
или с HTTPS:
python example.py https://example.com/path/page.html
Ожидание данных на открытом сокете с использованием потоков
Корутина, ожидающая получения данных сокетом с использованием функции open_connection():
import asyncio
import socket
async def wait_for_data():
# Get a reference to the current event loop because
# we want to access low-level APIs.
loop = asyncio.get_running_loop()
# Create a pair of connected sockets.
rsock, wsock = socket.socketpair()
# Register the open socket to wait for data.
reader, writer = await asyncio.open_connection(sock=rsock)
# Simulate the reception of data from the network
loop.call_soon(wsock.send, 'abc'.encode())
# Wait for data
data = await reader.read(100)
# Got data, we are done: close the socket
print("Received:", data.decode())
writer.close()
await writer.wait_closed()
# Close the second socket
wsock.close()
asyncio.run(wait_for_data())
См. также
В примере ожидание данных на открытом сокете с использованием протокола используется низкоуровневый протокол и метод loop.create_connection().
В примере отслеживание событий чтения файлового дескриптора для отслеживания файлового дескриптора используется низкоуровневый метод loop.add_reader().
© 2001 Python Software Foundation
Licensed under the PSF License.
https://docs.python.org/3.14/library/asyncio-stream.html