Spec-Zone.ru › Python 3.14

Потоки

Исходный код: 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.

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

Spec-Zone.ru

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