Asyncio なしでシングルスレッドのノンブロッキング非同期サーバーを作る(Feat. Event Loop を理解する)

Dev Team2025. 01. 22
リンクをコピーしました。
Asyncio なしでシングルスレッドのノンブロッキング非同期サーバーを作る(Feat. Event Loop を理解する)

こんにちは。バックエンド開発者のパク・ジョンインです.
Webサーバーを開発するにあたり、同じリソースでより多くのリクエストを処理するために、ますます多くの場面で非同期方式による開発が行われています。
私たちは社内で、このような非同期処理、とりわけ現在ほとんどの場合に使われているイベントループベースの非同期処理方式についての理解を深めるため、これに関連するセミナーを独自に準備して実施しました。
そして、このセミナーの内容が私たちの会社以外の方々にとっても、イベントループの動作方式を理解する助けになると思い、このように関連内容を共有することにしました。
この投稿では、Pythonの asyncio ライブラリと await & async 構文を使わずに、直接 socket を利用して非同期的にリクエストを受け取る簡単なサーバーを作ってみて、この過程で asyncio ライブラリの中核要素である event loop の原理について見ていきます。

1. ただのTelnetサーバー

最初は簡単に、Telnetクライアントからリクエストを受けたら、そのリクエスト内容を単純に echo するサーバーを作ってみましょう。

1.1. まずはリクエストをちゃんと受け取ってみよう

次のコードは、12345番ポートでTCP接続リクエストを受けるサーバーです。Telnetクライアントで12345番ポートにリクエストしてみると、runserver 関数の中にあるコードによってリクエストが処理されます。

# basic_server.py
import socket

def runserver():
    # Create a listening socket
    server_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
    server_socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
    server_socket.bind(('0.0.0.0', 12345))
    server_socket.listen()

    connection_socket, client_address = server_socket.accept()
    print(f"Connection established with {client_address}")

    data = connection_socket.recv(1024)
    print(data)
    connection_socket.close()

if __name__ == "__main__":
    runserver()

1.1.1. 動作方式

上のコードを1行ずつ見ていきましょう。

  1. サーバーで使用するソケットを生成します。
server_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)

socket.AF_INET は IPv4 アドレス体系を意味します。
socket.SOCK_STREAM は TCP を意味します。
つまり上のコードは、IPv4 アドレス体系を使用する TCP ソケットを生成して server_socket に割り当てます。

  1. 実習の便宜のため、ソケットに簡単なオプションを追加します。
server_socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)

socket.SOL_SOCKET はソケットレベルで設定を定義するという意味です。この引数の位置には socket.IPPROTO_IP を入れて IP レベルの設定をしたり、socket.IPPROTO_TCP を入れて TCP レベルの設定をしたりできますが、これ以上の説明は本投稿の範囲を外れるテーマなので省略し、ひとまず今はソケットに設定を行うという意図を持つ引数だと考えていただければ大丈夫です。
socket.SO_REUSEADDR は、このソケットがアドレスを再利用するという意味です。ここでいうアドレスとは IP、PORT の組を意味します。このオプションを使うことで、私たちがサーバーを起動したり停止したりするときに Address already in use エラーを防ぎます。
1 は上の socket.SO_REUSEADDR オプションを enable するという flag です。

  1. ソケットを 0.0.0.0 アドレスの 12345 番ポートにバインドします。
server_socket.bind(('0.0.0.0', 12345))

0.0.0.0 は、サーバーコンピュータにあるすべてのネットワークインターフェースを指す特別な IP です。つまり (‘0.0.0.0’, 12345) は、イーサネット、Wi-Fi などすべてのネットワークインターフェースを通じて 12345 番ポートに入ってくる通信を、このソケットが担当することを意味します。
ローカルでのみテストしたい場合は、ループバックインターフェースであり localhost を指す 127.0.0.1 でも構いません。

  1. ソケットで LISTEN を開始します。
server_socket.listen()

これで、このコンピュータの 12345 番ポートに入ってくる TCP 接続を受ける準備が整いました。

  1. リクエストが入ってくるまで待ちます。
connection_socket, client_address = server_socket.accept()

このサーバーに 12345 番ポートで TCP 接続リクエストが入ると、接続を成立させて connection_socket オブジェクトとクライアントのアドレス(IP, PORT)を返します。
サーバーに接続リクエストが入ってくるまでは、次のコード行へ進みません。

  1. リクエストを処理し、コネクションを閉じます。
data = connection_socket.recv(1024)
print(data)
connection.close()
  1. server_sockerconnection_socketは下の図でそれぞれ welcoming socket と connection socket を表します。
James F. Kurose & Keith W. Ross, 2022
James F. Kurose & Keith W. Ross, 2022

サーバーが私たちの設定した 12345 番ポートにバインドされたソケットでリクエストを受けると、別のソケットをもう1つ生成し、そのソケットでクライアントとデータをやり取りするようになります。この別のソケットがまさに connection_socket です。

1.1.2. 実行してみる

  1. 上のコードを basic_server.py という名前のファイルとして保存し、ターミナルで実行してみましょう。
python basic_server.py
  1. もう1つターミナルを開いて、Telnetクライアントを localhost 12345 に接続します。
telnet localhost 12345

# 以下の stdout を確認できる
Trying 127.0.0.1...
Connected to localhost.
Escape character is '^]'.
  1. Telnetクライアント上でメッセージを1つ送ってみます。
hi
Connection closed by foreign host.

hi というメッセージを送ると、connection_socket, client_address = server_socket.accept() のコードが進み、data = connection.recv(1024) で hi を読み取って data に割り当て、print(data) を通じて hi を basic_server.py が実行されているターミナルに出力します。
最終的に connection_socket.close() 行でコネクションを終了すると、ターミナルに Connection closed by foreign host. というメッセージとともに接続が終了します。

1.1.3. 改善すべき部分

  1. 現時点では1つのコネクションだけを処理してすぐに connection_socket.close() を呼び出しているため、Telnetクライアントから複数のメッセージを送ってもコネクションを閉じないようにするとよさそうです。
  1. 現時点ではTelnetクライアントのターミナルで見ると、自分が入力した hi だけが見えますが、echo サーバーを作るにはサーバー側でクライアントから受け取った hi を再びクライアントへ送るコードが必要です。

1.2. echo サーバーを作ってみよう

上のコードで、1つのコネクションだけを処理してコネクションを閉じる問題と、コネクションで送られたメッセージを echo する機能を追加したサーバーのコードは次のとおりです。

# echo_server.py
import socket

def runserver():
    server_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
    server_socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
    server_socket.bind(('0.0.0.0', 12345))
    server_socket.listen()

    while True:
        connection_socket, client_address = server_socket.accept()
        print(f"Connection established with {client_address}")

        while True:
            data = connection_socket.recv(1024)
            print(data)
            msg = b'echo: '
            msg += data
            connection_socket.send(msg)

            if data == b'bye\r\n':
                connection_socket.close()
                break

if __name__ == "__main__":
    runserver()

1.2.1. 変わった部分

変わった部分だけを見てみましょう。

data = connection_socket.recv(1024)
print(data)
connection_socket.close()

従来の、一度だけデータを受け取ってコネクションを受けるロジックから、

while True:
    data = connection_socket.recv(1024)
    print(data)
    msg = b'echo: '
    msg += data
    connection_socket.send(msg)

    if data == b'bye\r\n':
        connection_socket.close()
        break

bye というデータが来るまでループを回しながらコネクションからデータを受け取り、そのメッセージを再び send するロジックに変わりました。
このアップデートにより、次のように機能が改善されました.

1.
while ループを使用することで、複数のメッセージを送ってもコネクションが終了しない(1.1.2の1を解決)
2.
telnet クライアントでメッセージを送ると、応答として送ったメッセージをそのまま返してくれる(1.1.2の2を解決)
3.
telnet で bye を送信して、クライアントがコネクションを切断できる。

1.2.2. 実行してみる

  1. コードを実行します。
python echo_server.py
  1. ターミナルをもう1つ開いて、次のコマンドを実行します。
telnet localhost 12345

# 다음과 같은 stdout을 확인할 수 있음
Trying 127.0.0.1...
Connected to localhost.
Escape character is '^]'.
  1. 1.1.2.の3.とは異なり、これで hi を何度も送れるようになります。
hi
echo: hi
hi
echo: hi
hi
echo: hi
hi
echo: hi
hi
echo: hi
  1. bye でコネクションを終了します。
bye
echo: bye
Connection closed by foreign host.

2. 複数コネクションを処理するノンブロッキング非同期 echo サーバー

python の wsgi 実装体は、コネクションを生成するとき、結局は上のような方式で動作します

benoitc, 2014
benoitc, 2014

ここにマルチスレッドあるいはマルチプロセシングを利用して、同時に複数のコネクションを受け付けます。しかし私たちの目的は、1つのプロセスで複数のリクエストを非同期に受けられるサーバーを作ることです。そしてこの方式は asgi の実装体が動作する方式でもあります。これからその方向に向けて、サーバーを少しずつ改善していきましょう。

2.1. コードをノンブロッキングに変える

既存のサーバーが1つのコネクションしか受けられない理由は、次のコードで Telnet リクエストが来るまでプログラムが止まっており、その外側にある while ループが回らないからです。

while True:
    connection_socket, client_address = server_socket.accept()

このように、あるコードが自分の仕事を終えるまで次のコードに進まないことをブロッキングと言います。
今の場合、上のコードでブロッキングされている間は、単にクライアントからのコネクションを待っている状態なので、CPU は何の処理もしていません。これは非効率なので、コネクションが来なくてもコードのループが回り続けられるようにすることが、複数のコネクションを非同期に受けるための第一歩になるでしょう。
そこで、次のようにコードを書くと、コネクションが来なくてもひとまず while ループが回るようにできます。

import socket

def runserver():
    server_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
    server_socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
    server_socket.bind(('0.0.0.0', 12345))
    server_socket.listen()
    server_socket.setblocking(False)

    while True:
        try:
            connection_socket, client_address = server_socket.accept()
            print(f"Connection established with {client_address}")
        except BlockingIOError:
            continue

        while True:
            try:
                data = connection_socket.recv(1024)
            except BlockingIOError:
                continue

            print(data)
            msg = b'echo: '
            msg += data
            connection_socket.send(msg)

            if data == b'bye\r\n':
                connection_socket.close()
                print(f"connection with {client_address} closed")
                break

if __name__ == "__main__":
    runserver()

2.1.1. 変わった部分

  1. サーバーソケットを生成する部分にコードが1行追加されました。
server_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
server_socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
server_socket.bind(('0.0.0.0', 12345))
server_socket.listen()
server_socket.setblocking(False) # 추가된 부분

ソケットに .setblocking(False) を設定すると、そのソケットで .accept().recv() のようにクライアント側の動作を期待する関数が呼ばれたとき、クライアントから何の動作もなければ BlockingIOError Exception を発生させるようになります。

  1. クライアントからコネクション生成の試みがまったくなく、.accept() の呼び出しが BlockingIOError Exception を投げた場合には、これを pass してコードがそのまま進行できるようにします。
while True:
    try:                                # 추가된 부분
        connection_socket, client_address = server_socket.accept()
        print(f"Connection established with {client_address}")
        connections.append(connection)  # 추가된 부분
    except BlockingIOError:             # 추가된 부분
        pass                            # 추가된 부분
  1. コネクションからクライアントのデータを受け取るとき、クライアントがまだ何のデータも送っておらず、.recv(1024) の呼び出しが BlockingIOError Exception を投げた場合には、これを pass してコードが続けて進行できるようにします。
while True:
    try:                     # 추가된 부분
        data = connection_socket.recv(1024)
    except BlockingIOError:  # 추가된 부분
        continue             # 추가된 부분

    print(data)

これにより、サーバーはクライアントから何の動作がなくても while ループを回し続けられるようになりました。しかし、まだ1つのコネクションしか受けられません。ここであと1段階進めば、このサーバーは複数のコネクションを受けられるようになります。

2.2. 複数のコネクションを処理する

次のコードは、2.2.1.のコードから複数のコネクションを受けるために改善されたコードです。

import socket

def runserver():
    server_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
    server_socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
    server_socket.bind(('0.0.0.0', 12345))
    server_socket.listen()
    server_socket.setblocking(False)
    connections = []

    while True:
        try:
            connection_socket, client_address = server_socket.accept()
            print(f"Connection established with {client_address}")
            connections.append(connection_socket)
        except BlockingIOError:
            pass

        for connection_socket in connections:
            try:
                data = connection_socket.recv(1024)
            except BlockingIOError:
                continue

            print(f"send to {client_address}: {data}")
            msg = b'echo: '
            msg += data
            connection_socket.send(msg)

            if data == b'bye\r\n':
                connection_socket.close()
                print(f"connection with {client_address} closed")
                connections.remove(connection_socket)
                break

if __name__ == "__main__":
    runserver()

2.2.1. 変わった部分

  1. 次のように、複数のコネクションを入れるリスト宣言を追加しました。
server_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
server_socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
server_socket.bind(('0.0.0.0', 12345))
server_socket.listen()
server_socket.setblocking(False)
connections = []  # 추가된 부분
  1. 1つのコネクションの recv だけを待っていた部分を、connections リストを回りながら複数のコネクションから recv を行うように改善しました。
for conn in connections:  # 변경된 부분
    try:
        data = connection_socket.recv(1024)
    except BlockingIOError:
        continue
  1. クライアントがコネクションを切断したら、connections リストから該当コネクションを削除するコードを追加しました。
if data == b'bye\r\n':
    connection_socket.close()
    print(f"connection with {client_address} closed")
    connections.remove(connection_socket)    # 추가된 부분
    break                                    # 추가된 부분

これにより、このサーバープログラムは1つのプロセスだけでも複数のコネクションを受けられる非同期サーバーになりました。しかし、現在このプログラムには while ループを休みなく繰り返し、CPU リソースをすべて消費してしまう致命的な問題があります。だからといって1秒ずつ休みながらループを回すと、クライアントはサーバーが休んでいる時間分の遅延を経験することになります。次のパートではこの問題を解決してみましょう。

3. I/O event notification を利用した非同期 echo サーバー

3.1. I/O event を使う

私たちが使う OS には、あるファイルで読み取りや書き込みなどの I/O event が発生したとき、プロセスがその event について通知を受け取れる API が存在します。

  • Linux: epoll
  • MacOS: kqueue
  • Windows: IOCP

Python では selectors というライブラリを使うことで、上記 API を利用できます。ソケットも本質的にはファイルなので、上記 API をソケットに対しても使えます。そのため、ソケットに対する I/O event を購読し、そのソケットに変更が発生したら通知を受ける形で実装すれば、クライアントの動作を待つために While ループを休みなく回し続ける必要がなくなります。
まず selectors を利用したサーバー実装コードを見て、どのように改善したのかを一つずつ見ていきましょう。

import selectors
import socket

def runserver():
    server_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
    server_socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
    server_socket.bind(('0.0.0.0', 12345))
    server_socket.listen()
    server_socket.setblocking(False)

    selector = selectors.DefaultSelector()
    selector.register(server_socket, selectors.EVENT_READ)

    while True:
        events = selector.select(timeout=1)

        if len(events) == 0:
            continue

        for event, _ in events:
            event_socket = event.fileobj

            if event_socket == server_socket:
                connection_socket, client_address = server_socket.accept()
                print(f"Connection established with {client_address}")
                selector.register(connection_socket, selectors.EVENT_READ)
            else:
                connection_socket = event_socket
                data = connection_socket.recv(1024)
                msg = b'echo: '
                msg += data
                connection_socket.send(msg)

                if data == b'bye\r\n':
                    connection_socket.close()
                    print(f"connection with {client_address} closed")
                   break

if __name__ == "__main__":
    runserver()

3.2. 変わった部分

  1. ソケットに I/O event を購読する部分が追加されました。
selector = selectors.DefaultSelector()
selector.register(server_socket, selectors.EVENT_READ)

上のコードでは seletor オブジェクトを生成し、そのオブジェクトが server_socket の読み取りイベント、つまり connection 作成リクエストに関する通知を受け取るように設定します。

  1. selector に登録された socket たちから I/O event を待ちます。
while True:
    events = selector.select(timeout=1)

    if len(events) == 0:
        continue

上のコードでは1秒間、selector オブジェクトに登録されたソケットに対してクライアントからのリクエストを待ちます。1秒が経つと次のコードへ進み、if len(events) == 0: の部分に到達し、1秒の間に何のイベントもなければ continue によって while ループが回ります。
このコードは重要な部分です。先ほどのコードでは while ループを休みなく回しながら、ソケットにクライアントのリクエストがあるかないかを確認するために CPU リソースをすべて消費していましたが、このコードでは selector を通じてソケットにクライアントのリクエストがあるかを確認する役割を、epoll や kqueue API を通じて OS に委任できます。OS はその役割を実行するとき、CPU を使わず効率的に処理します。

  1. もしクライアントからリクエストがあったなら、該当ソケットを取り出します。
for event, _ in events:
    event_socket = event.fileobj

events を for 文に渡すと、SelectorKey オブジェクトの event と int 型の file descriptor 番号である _ にアンパックされます。file descriptor 番号は使わないので _ にアンパックしました。
SelectorKey オブジェクトの fileobj プロパティを参照すると、先ほどの .accept.recv の呼び出しが可能なソケットオブジェクトを受け取れます。
ここで受け取る event_socket は、コネクションリクエストを受けた server_socket である場合もあれば、サーバーの応答を待っている connection_socket である場合もあります。
実は、現時点までを見ると connection_socket にはまだ何の通知も購読されていないため、event_socket がどうして connection_socket にもなり得るのか理解しにくいかもしれません。しかし、このあとすぐ続く 4 で説明します。

  1. もし event_socketserver_socket であれば、これはコネクション作成リクエストを意味するので、コネクションを受け取って connections リストに追加しておきます。
if event_socket == server_socket:
    connection_socket, client_address = server_socket.accept()
    print(f"Connection established with {client_address}")
    selector.register(connection_socket, selectors.EVENT_READ)

上のコードの最後の行を見ると、上の 1 で server_socket に対して行ったのと同じように connection_socket に対しても I/O event を購読しています。これにより、event_socketconnection_socket でもあり得るようになりました。

  1. もし event_socketconnection_socket であれば、これはクライアントがサーバーへコネクションを通じてリクエストを送ったことを意味するため、echo 応答を返します。
else:
    connection_socket = event_socket
    data = connection_socket.recv(1024)
    msg = b'echo: '
    msg += data
    connection_socket.send(msg)
  1. コネクションを閉じるリクエストが来たら、コネクションを整理することも忘れません。
if data == b'bye\r\n':
    connection_socket.close()
    print(f"connection with {client_address} closed")
   break

この改善を通じて、CPU リソースを過度に占有することなく、1つのサーバープロセスだけで複数のリクエストを非同期に受けられるサーバーを実装してみました。

4. 実際の asyncio の Event Loop 実装と比較する

今回のパートでは、上で実装したサーバーをもとに asyncio の event loop を探ってみましょう。
実は、上で selectors 非同期サーバーを実装しながら、私たちはすでに簡単な event loop を作っていました。3.2 の 1 の部分のコードが、まさに簡略化された event loop です。

while True:
    events = selector.select(timeout=1)

    if len(events) == 0:
        continue
    
    ...

このアイデアをもとに、asyncio の event loop が実際にはどのように実装されているのかを見てみましょう。
BaseEventLoop.run_forever を見ると、while: true 構文を見つけることができます。(benoitc et al, 2023)

def run_forever(self):
    """Run until stop() is called."""

  ...

  events._set_running_loop(self)
  while True:
      self._run_once()
      if self._stopping:
          break
  ...

while: true の中にある self._run_once の定義を見てみましょう。

def _run_once(self):
    """Run one full iteration of the event loop.

    This calls all currently ready callbacks, polls for I/O,
    schedules the resulting callbacks, and finally schedules
    'call_later' callbacks.
    """
  
  ...

  event_list = self._selector.select(timeout)
  self._process_events(event_list)

  ...

上に抜粋した部分を見ると、asyncio でも私たちが使っていた selectors を使用している様子を確認できます。
この姿は、結局のところ私たちが 3.2 の 2 のコードで見たものとよく似ています。

while True:
    events = selector.select(timeout=1)

    if len(events) == 0:
        continue

その次に呼び出される _process_events の実装部分を見ると、次のとおりです。

def _process_events(self, event_list):
   for key, mask in event_list:
       fileobj, (reader, writer) = key.fileobj, key.data
       if mask & selectors.EVENT_READ and reader is not None:
           if reader._cancelled:
               self._remove_reader(fileobj)
           else:
               self._add_callback(reader)
       if mask & selectors.EVENT_WRITE and writer is not None:
           if writer._cancelled:
               self._remove_writer(fileobj)
            else:
               self._add_callback(writer)

反復文の前半部分を見ると、3.2 の 3 の部分と似ていることがわかります。

for event, _ in events:
    event_socket = event.fileobj

asyncio の event loop では、単に read I/O event だけでなく、より一般的なケースをカバーできるように実装されています。
このように asyncio の event loop もまた while: true 構文の中で selectors を使って I/O event を待ちながら非同期的な動作を実装していることが確認できます。

5. 結論

これまで asyncio なしで非同期ノンブロッキングサーバーを実装する方法を見てきて、これを通じて asyncio が提供する event loop の原理を見てきました。
event loop とは、結局 while: true の中で I/O event を制御するループであることがわかりました。
それでは今後さらに考えてみるべき部分は、asyncio の Coroutine とは具体的に何なのかということでしょう。近い将来、Coroutine に関する投稿も扱ってみます。

引用資料:
benoitc. (2014, October 25). optimize the sync worker. https://github.com/benoitc/gunicorn/commit/4c601ce447fafbeed27f0f0a238e0e48c928b6f9
benoitc, et al. (2023, December 7). n.d. https://github.com/benoitc/gunicorn/blob/9802e21f779d9f1f208a1a3288218bd5b843ad46/gunicorn/workers/sync.py
James F. Kurose & Keith W. Ross. (2022). Computer Networking A Top-Down Approach (8th edition). Pearson.
mindmajix, (2023, April 4), Express JS Interview Questions https://mindmajix.com/express-js-interview-questions

ストーリー一覧

最新ストーリー