جریان‌ها

کد منبع: Lib/asyncio/streams.py


جریان‌ها اولیه‌های سطح بالای آماده برای ناهمگام و 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!'))

همچنین بخش Examples را در زیر ببینید.

توابع جریان

می‌توان از توابع سطح‌بالای 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 روی ۶۴ KiB تنظیم شده است.

بقیه‌ی آرگومان‌ها مستقیماً به 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 روی ۶۴ KiB تنظیم شده است.

بقیه‌ی آرگومان‌ها مستقیماً به loop.create_server() ارسال می‌شوند.

توجه

آرگومان sock مالکیت سوکت را به سرور ایجادشده منتقل می‌کند. برای بستن سوکت، متد close() سرور را فراخوانی کنید.

تغییر یافته در نسخه‌ی 3.7: پارامترهای ssl_handshake_timeout و start_serving اضافه شدند.

تغییر یافته در نسخه‌ی 3.10: پارامتر loop حذف شد.

تغییر یافته در نسخه‌ی 3.11: پارامتر ssl_shutdown_timeout افزوده شد.

تغییر یافته در نسخه‌ی 3.13: پارامتر keep_alive اضافه شد.

سوکت‌های یونیکس

async asyncio.open_unix_connection(path=None, *, limit=65536, ssl=None, sock=None, server_hostname=None, ssl_handshake_timeout=None, ssl_shutdown_timeout=None)

اتصال سوکت یونیکس را برقرار می‌کند و یک جفت (reader, writer) را برمی‌گرداند.

مشابه open_connection() است، اما روی سوکت‌های یونیکس عمل می‌کند.

همچنین مستندات loop.create_unix_connection() را ببینید.

توجه

آرگومان sock مالکیت سوکت را به StreamWriter ایجادشده منتقل می‌کند. برای بستن سوکت، متد close() آن را فراخوانی کنید.

تغییر یافته در نسخه‌ی 3.7: پارامتر ssl_handshake_timeout افزوده شد. پارامتر path اکنون می‌تواند یک path-like object باشد

تغییر یافته در نسخه‌ی 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)

یک سرور سوکت یونیکس را راه‌اندازی کنید.

مشابه start_server() است، اما با سوکت‌های یونیکس کار می‌کند.

اگر cleanup_socket درست باشد، سوکت یونیکس به‌طور خودکار هنگام بسته شدن سرور از سامانه فایل‌بندی حذف می‌شود، مگر اینکه سوکت پس از ایجاد سرور جایگزین شده باشد.

همچنین مستندات loop.create_unix_server() را ببینید.

توجه

آرگومان sock مالکیت سوکت را به سرور ایجادشده منتقل می‌کند. برای بستن سوکت، متد close() سرور را فراخوانی کنید.

تغییر یافته در نسخه‌ی 3.7: پارامترهای ssl_handshake_timeout و start_serving افزوده شدند. پارامتر path اکنون می‌تواند یک path-like object باشد.

تغییر یافته در نسخه‌ی 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()

پایان پرونده را تأیید کنید.

async read(n=-1)

تا n بایت از جریان را بخوانید.

اگر n ارائه‌نشده باشد یا روی -1 تنظیم شده باشد، تا EOF می‌خواند، سپس تمام bytes خوانده‌شده را برمی‌گرداند. اگر EOF دریافت شده باشد و بافر داخلی خالی باشد، یک شیء bytes خالی برمی‌گرداند.

اگر n برابر 0 باشد، بلافاصله یک شیء bytes خالی برمی‌گرداند.

اگر n مثبت باشد، به محض اینکه حداقل ۱ بایت در بافر داخلی در دسترس باشد، حداکثر n بایت در دسترس را به‌صورت bytes برمی‌گرداند. اگر EOF پیش از خوانده‌شدن هیچ بایتی دریافت شود، یک شیء bytes خالی برمی‌گرداند.

async readline()

یک خط را بخوانید، که در آن «سطر» دنباله‌ای از بایت‌ها است که با \n پایان می‌یابد.

اگر EOF دریافت شود و \n یافت نشود، متد داده‌های خوانده‌شده‌ی ناقص را برمی‌گرداند.

در صورتی که EOF دریافت شود و بافر داخلی خالی باشد، یک شیء bytes خالی برمی‌گرداند.

async readexactly(n)

دقیقاً n بایت را بخوانید.

اگر پیش از خواندن n به EOF برسید، یک IncompleteReadError پرتاب می‌شود. برای دریافت داده‌هایی که به‌طور ناقص خوانده شده‌اند، از ویژگی IncompleteReadError.partial استفاده کنید.

async readuntil(separator=b'\n')

داده‌ها را از جریان می‌خواند تا separator یافت شود.

در صورت موفقیت، داده و جداکننده از بافر داخلی حذف می‌شوند (مصرف می‌شوند). داده‌ی برگردانده‌شده شامل جداکننده در انتها خواهد بود.

اگر مقدار داده‌ی خوانده‌شده از محدودیت پیکربندی‌شده‌ی جریان بیشتر شود، استثنای LimitOverrunError پرتاب می‌شود و داده در بافر درونی باقی می‌ماند و می‌تواند دوباره خوانده شود.

اگر پیش از پیدا شدن جداکننده‌ی کامل، به EOF رسیده شود، استثنای IncompleteReadError پرتاب می‌شود و بافر داخلی بازنشانی می‌شود. ویژگی IncompleteReadError.partial ممکن است شامل بخشی از جداکننده باشد.

separator همچنین می‌تواند یک تاپل از جداسازها باشد. در این حالت، مقدار بازگشتی کوتاه‌ترین مقدار ممکن خواهد بود که یکی از جداسازها را به‌عنوان پسوند داشته باشد. از نظر LimitOverrunError، کوتاه‌ترین جداساز ممکن به‌عنوان جداسازی در نظر گرفته می‌شود که مطابقت داشته است.

اضافه شده در نسخه‌ی 3.5.2.

تغییر یافته در نسخه‌ی 3.13: اکنون پارامتر separator می‌تواند یک tuple از جداکننده‌ها باشد.

at_eof()

اگر بافر خالی باشد و feed_eof() فراخوانی شده باشد، True را برمی‌گرداند.

StreamWriter

class asyncio.StreamWriter

نشان‌دهنده‌ی یک شیء نویسنده است که APIهایی را برای نوشتن داده‌ها در جریان IO فراهم می‌کند.

نمونه‌سازی مستقیم اشیای StreamWriter توصیه نمی‌شود؛ در عوض از open_connection() و start_server() استفاده کنید.

write(data)

این متد تلاش می‌کند data را بلافاصله در سوکت زیرین بنویسد. اگر این کار با شکست مواجه شود، داده در یک بافر نوشتن داخلی در صف قرار می‌گیرد تا زمانی که بتوان آن را ارسال کرد.

بافر data باید یک شیء bytes، bytearray یا memoryview یک‌بعدی با چیدمان پیوسته در C (C-contiguous) باشد.

این متد باید همراه با متد 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()

اگر انتقال زیربنایی از متد write_eof() پشتیبانی کند، True را برمی‌گرداند، در غیر این صورت False را برمی‌گرداند.

write_eof()

پس از تخلیه‌شدن داده‌های نوشتاری بافرشده، پایانه نوشتاری جریان را ببندید.

transport

انتقال زیربنایی asyncio را برمی‌گرداند.

get_extra_info(name, default=None)

به اطلاعات اختیاری انتقال دسترسی پیدا کنید؛ برای جزئیات، BaseTransport.get_extra_info() را ببینید.

async drain()

صبر کنید تا زمان مناسب برای ازسرگیری نوشتن در جریان برسد. مثال:

writer.write(data)
await writer.drain()

این یک متد کنترل جریان است که با بافر نوشتاری زیربنایی IO تعامل دارد. هنگامی که اندازه‌ی بافر به آستانه‌ی بالا (high watermark) برسد، drain() مسدود می‌شود تا زمانی که بافر تخلیه شود و اندازه‌ی آن به آستانه‌ی پایین (low watermark) برسد و بتوان نوشتن را از سر گرفت. هنگامی که چیزی برای انتظار وجود نداشته باشد، drain() بلافاصله بازمی‌گردد.

توجه

هنگامی که بافر نوشتن کمتر از آستانه بالا باشد، drain() بلافاصله و بدون واگذاری کنترل به حلقه رویداد بازمی‌گردد. در نتیجه، کدی که به‌طور مکرر write() و سپس await drain() را فراخوانی می‌کند، ممکن است مانع اجرای سایر taskها شود. برای جلوگیری از رفتار مسدودکننده، به‌صراحت با 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 پیش از قطع اتصال است. اگر None باشد، 60.0 ثانیه (پیش‌فرض).

  • ssl_shutdown_timeout زمان بر حسب ثانیه است که باید برای تکمیل خاموش‌سازی SSL پیش از قطع اتصال منتظر بمانید. اگر None باشد، 30.0 ثانیه است (پیش‌فرض).

اضافه شده در نسخه‌ی 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() برای پایش یک توصیف‌گر پرونده استفاده می‌کند.