Перейти к содержимому

18.5.5. Потоки данныхAPI на основе корутин​

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

18.5.5.1. Функции потоков данных​

Примечание

Функции верхнего уровня в этом модуле представляют собой лишь обёртки для удобства; в них нет ничего особенного, и если они не выполняют именно то, что требуется, можно скопировать их код.

coroutineasyncio.open_connection(host=None, port=None, *, loop=None, limit=None, **kwds)

Обёртка для create_connection(), возвращающая пару (reader, writer).

Возвращаемый reader является экземпляром StreamReader, а writer – экземпляром StreamWriter.

Все аргументы такие же, как обычно для AbstractEventLoop.create_connection(), за исключением protocol_factory; наиболее часто используются позиционные параметры host и port, за которыми следуют различные необязательные ключевые аргументы.

Дополнительные необязательные ключевые аргументы: loop (для указания экземпляра цикла событий) и limit (для установки лимита буфера, передаваемого в StreamReader).

Эта функция является корутиной.

coroutineasyncio.start_server(client_connected_cb, host=None, port=None, *, loop=None, limit=None, **kwds)

Запускает сокет-сервер с колбэком для каждого подключённого клиента. Возвращаемое значение такое же, как у create_server().

Параметр client_connected_cb вызывается с двумя параметрами: client_reader и client_writer. client_reader – это объект StreamReader, а client_writer – объект StreamWriter. Параметр client_connected_cb может быть как обычной функцией обратного вызова, так и корутинной функцией; если это корутинная функция, она автоматически преобразуется в Task.

Остальные аргументы такие же, как обычно для create_server(), за исключением protocol_factory; наиболее часто используются позиционные host и port, за которыми следуют различные необязательные ключевые аргументы.

Дополнительные необязательные ключевые аргументы: loop (для указания экземпляра цикла событий) и limit (для установки лимита буфера, передаваемого в StreamReader).

Эта функция является корутиной.

coroutineasyncio.open_unix_connection(path=None, *, loop=None, limit=None, **kwds)

Обёртка для create_unix_connection(), возвращающая пару (reader, writer).

См. open_connection() для информации о возвращаемом значении и других деталях.

Эта функция является корутиной.

Доступность: UNIX.

coroutineasyncio.start_unix_server(client_connected_cb, path=None, *, loop=None, limit=None, **kwds)

Запускает сервер UNIX Domain Socket с колбэком для каждого подключённого клиента.

См. start_server() для информации о возвращаемом значении и других деталях.

Эта функция является корутиной.

Доступность: UNIX.

18.5.5.2. StreamReader​

class asyncio.StreamReader(limit=_DEFAULT_LIMIT, loop=None)

Этот класс является не потокобезопасным.

Значение аргумента limit по умолчанию равно _DEFAULT_LIMIT, то есть 2**16 (64 KiB).

exception()

Получить исключение.

feed_eof()

Подтверждает EOF.

feed_data(data)

Помещает байты data во внутренний буфер. Любые операции, ожидающие данные, будут возобновлены.

set_exception(exc)

Устанавливает исключение.

set_transport(транспорт)

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

coroutineread(n=-1)

Читает до n байт. Если n не указан или равен -1, читает до EOF и возвращает все прочитанные байты.

Если был получен EOF и внутренний буфер пуст, возвращает пустой объект bytes.

Этот метод является корутиной.

coroutinereadline()

Читает одну строку, где «строка» – это последовательность байтов, заканчивающаяся на \n.

Если получен EOF, а \n не найден, метод возвращает частично прочитанные байты.

Если был получен EOF и внутренний буфер пуст, возвращает пустой объект bytes.

Этот метод является корутиной.

coroutinereadexactly(n)

Читает ровно n байт. Вызывает IncompleteReadError, если конец потока достигнут до того, как можно прочитать n; атрибут IncompleteReadError.partial исключения содержит частично прочитанные байты.

Этот метод является корутиной.

coroutinereaduntil(separator=b'\n')

Читает данные из потока до тех пор, пока не будет найден separator.

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

Настроенное ограничение потока используется для проверки результата. Ограничение задаёт максимальную длину данных, которые могут быть возвращены, не считая разделителя.

Если происходит EOF и полный разделитель всё ещё не найден, будет вызвано исключение IncompleteReadError, а внутренний буфер будет сброшен. Атрибут IncompleteReadError.partial может содержать разделитель частично.

Если данные не удаётся прочитать из-за превышения ограничения, будет вызвано исключение LimitOverrunError, а данные останутся во внутреннем буфере, чтобы их можно было прочитать снова.

Новое в версии 3.5.2.

at_eof()

Возвращает True, если буфер пуст и был вызван feed_eof().

18.5.5.3. StreamWriter​

classasyncio.StreamWriter(транспорт, протокол, читатель, цикл событий)

Оборачивает транспорт.

Это предоставляет write(), writelines(), can_write_eof(), write_eof(), get_extra_info() и close(). Он добавляет drain(), который возвращает необязательный Future, на котором можно ожидать управления потоком. Он также добавляет атрибут транспорта, который ссылается непосредственно на Transport.

Этот класс является не потокобезопасным.

transport

Транспорт.

can_write_eof()

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

close()

Закрывает транспорт: см. BaseTransport.close().

coroutinedrain()

Даёт возможность сбросить буфер записи нижележащего транспорта.

Предполагаемое использование: записать:

python
w.write(data)
yield from w.drain()

Когда размер буфера транспорта достигает верхнего предела (протокол приостановлен), блокирует выполнение до тех пор, пока размер буфера не уменьшится до нижнего предела и протокол не будет возобновлён. Если ожидать нечего, yield-from продолжается немедленно.

Yield from из drain() даёт циклу событий возможность запланировать операцию записи и сбросить буфер. Особенно это следует использовать, когда в транспорт записывается возможно большой объём данных, и корутина не делает yield-from между вызовами write().

Этот метод является корутиной.

get_extra_info(имя, по умолчанию=None)

Возвращает дополнительную информацию о транспорте: см. BaseTransport.get_extra_info().

write(данные)

Записывает некоторое количество байт данных в транспорт: см. WriteTransport.write().

writelines(данные)

Записывает список (или любую итерируемую последовательность) байтов данных в транспорт: см. WriteTransport.writelines().

write_eof()

Закрывает конец записи транспорта после сброса буферизованных данных: см. WriteTransport.write_eof().

18.5.5.4. StreamReaderProtocol​

classasyncio.StreamReaderProtocol(читатель потока, колбэк подключения клиента=None, цикл событий=None)

Вспомогательный класс для адаптации между Protocol и StreamReader. Подкласс Protocol.

stream_reader – это экземпляр StreamReader, client_connected_cb – необязательная функция, вызываемая с (stream_reader, stream_writer) при установлении соединения, а loop – используемый экземпляр цикла событий.

(Это вспомогательный класс, а не подкласс Protocol самого StreamReader, поскольку у StreamReader есть и другие потенциальные применения, и чтобы пользователь StreamReader случайно не вызвал неподходящие методы протокола.)

18.5.5.5. IncompleteReadError​

exceptionasyncio.IncompleteReadError

Ошибка неполного чтения, подкласс EOFError.

expected

Общее количество ожидаемых байтов (int).

partial

Строка прочитанных байтов до достижения конца потока (bytes).

18.5.5.6. LimitOverrunError​

exceptionasyncio.LimitOverrunError

Достигнут лимит буфера при поиске разделителя.

consumed

Общее количество байтов, подлежащих обработке.

18.5.5.7. Примеры использования потоков​

18.5.5.7.1. TCP-эхо-клиент с использованием потоков​

TCP-эхо-клиент с использованием функции asyncio.open_connection():

python
import asyncio

@asyncio.coroutine
def tcp_echo_client(message, loop):
    reader, writer = yield from asyncio.open_connection('127.0.0.1', 8888,
                                                        loop=loop)

    print('Send: %r' % message)
    writer.write(message.encode())

    data = yield from reader.read(100)
    print('Received: %r' % data.decode())

    print('Close the socket')
    writer.close()

message = 'Hello World!'
loop = asyncio.get_event_loop()
loop.run_until_complete(tcp_echo_client(message, loop))
loop.close()

Смотрите также

В примере протокол TCP-эхо-клиента используется метод AbstractEventLoop.create_connection().

18.5.5.7.2. TCP-эхо-сервер с использованием потоков​

TCP-эхо-сервер с использованием функции asyncio.start_server():

python
import asyncio

@asyncio.coroutine
def handle_echo(reader, writer):
    data = yield from reader.read(100)
    message = data.decode()
    addr = writer.get_extra_info('peername')
    print("Received %r from %r" % (message, addr))

    print("Send: %r" % message)
    writer.write(data)
    yield from writer.drain()

    print("Close the client socket")
    writer.close()

loop = asyncio.get_event_loop()
coro = asyncio.start_server(handle_echo, '127.0.0.1', 8888, loop=loop)
server = loop.run_until_complete(coro)

# Обслуживать запросы до нажатия Ctrl+C
print('Serving on {}'.format(server.sockets[0].getsockname()))
try:
    loop.run_forever()
except KeyboardInterrupt:
    pass

# Закрыть сервер
server.close()
loop.run_until_complete(server.wait_closed())
loop.close()

Смотрите также

Пример протокола TCP-эхо-сервера использует метод AbstractEventLoop.create_server().

18.5.5.7.3. Получение HTTP-заголовков​

Простой пример запроса HTTP-заголовков URL, переданного в командной строке:

python
import asyncio
import urllib.parse
import sys

@asyncio.coroutine
def print_http_headers(url):
    url = urllib.parse.urlsplit(url)
    if url.scheme == 'https':
        connect = asyncio.open_connection(url.hostname, 443, ssl=True)
    else:
        connect = asyncio.open_connection(url.hostname, 80)
    reader, writer = yield from connect
    query = ('HEAD {path} HTTP/1.0\r\n'
             'Host: {hostname}\r\n'
             '\r\n').format(path=url.path or '/', hostname=url.hostname)
    writer.write(query.encode('latin-1'))
    while True:
        line = yield from reader.readline()
        if not line:
            break
        line = line.decode('latin1').rstrip()
        if line:
            print('HTTP header> %s' % line)

    # Игнорировать тело, закрыть сокет
    writer.close()

url = sys.argv[1]
loop = asyncio.get_event_loop()
task = asyncio.ensure_future(print_http_headers(url))
loop.run_until_complete(task)
loop.close()

Использование:

python
python example.py http://example.com/path/page.html

или с HTTPS:

python
python example.py https://example.com/path/page.html

18.5.5.7.4. Регистрация открытого сокета для ожидания данных с использованием потоков​

Корутина, ожидающая получения данных сокетом с помощью функции open_connection():

python
import asyncio
try:
    from socket import socketpair
except ImportError:
    from asyncio.windows_utils import socketpair

@asyncio.coroutine
def wait_for_data(loop):
    # Создать пару соединённых сокетов
    rsock, wsock = socketpair()

    # Зарегистрировать открытый сокет для ожидания данных
    reader, writer = yield from asyncio.open_connection(sock=rsock, loop=loop)

    # Симулировать приём данных из сети
    loop.call_soon(wsock.send, 'abc'.encode())

    # Ожидать данные
    data = yield from reader.read(100)

    # Получены данные, готово: закрыть сокет
    print("Received:", data.decode())
    writer.close()

    # Закрыть второй сокет
    wsock.close()

loop = asyncio.get_event_loop()
loop.run_until_complete(wait_for_data(loop))
loop.close()

Смотрите также

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

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