最近、「Python WebSocket」が注目を集めていますね!どうやら、「Python WebSocket」を使うと、REST APIでは難しかったリアルタイムデータ処理が驚くほど簡単に実装できる ようです!

そこで今回は、Pythonのwebsocketsライブラリで「Python WebSocket」の基本から実践的な使い方まで解説を行ってみました!初めてWebSocketを触る方でも読み進めやすい内容 ですので、ぜひ皆さんも記事を読んで試してみてください!

この記事で分かること

  • REST APIとWebSocketの違いと、どちらを選ぶべきかの判断基準
  • websocketsライブラリを使ったBinance仮想通貨オーダーブックのリアルタイム取得コード
  • 自動再接続・Heartbeat・ログ管理など本番運用に必要な安定化テクニック

REST vs WebSocket:どちらを選ぶべきか

「価格データを1秒ごとに取得したい」「チャットメッセージを即座に受け取りたい」――そんなリアルタイム処理で詰まったことはありませんか?従来のREST APIでも実現できますが、WebSocketを使うと はるかに効率的に実装できます。

Python WebSocketでできること

  • サーバーからクライアントへのデータプッシュ配信
  • 単一のTCP接続を維持したまま双方向通信
  • ミリ秒単位の低レイテンシなリアルタイムデータ受信

ここで、「REST APIで定期取得すればよいのでは?」という疑問が出てくると思います。

大きな差としては、「サーバーから能動的にデータをpushできる」が可能ということです。これにより、クライアントが何度もリクエストを送らずに最新データを受け取り続ける ことが出来ます。

今話題の、「リアルタイムストリーミング」というやつですね!

比較すると、以下のようになります。

項目REST APIWebSocket
通信方向クライアント→サーバー(一方向リクエスト)双方向(サーバー・クライアント両側から送信可)
接続の維持リクエストごとに接続・切断一度接続したら持続(コネクション維持)
リアルタイム性ポーリング間隔に依存(遅延あり)イベント発生と同時に受信(低レイテンシ)

「リアルタイムストリーミングって難しそう…」と感じるかもしれませんが、この記事を読めば初心者の方でも大丈夫!順を追って自動再接続や本番運用テクニックまで解説します。

ポーリングの限界とWebSocketが解決すること

つまり、ポーリングとは「一定間隔でサーバーに問い合わせ続ける」手法です。たとえば1秒ごとにREST APIを叩いて価格を取得するイメージですね。しかしこの方法には「更新がなくても毎回リクエストが 発生する」「間隔を短くするほどサーバー負荷が上がる」という限界があります。WebSocketは接続を張りっぱなしにして、変化があったときだけサーバーからデータを送ってくるため、無駄なリクエストをゼロ にできます。

向いているユースケース・向いていないユースケース

WebSocketが向いているのは、仮想通貨の価格配信・オーダーブック更新・チャットアプリ・リアルタイムログ監視など「データが高頻度で変化する」場面です。一方、データ取得頻度が低い(たとえば 1日1回の集計取得)・シンプルなCRUD操作・キャッシュが効くコンテンツ配信には、引き続きREST APIのほうが適しています。用途に合わせて使い分けるのがベストです。

websocketsライブラリの基本セットアップ

手順1:websocketsのインストール

PythonでWebSocketを扱うには、サードパーティライブラリ websockets が最もシンプルでおすすめです。pipで一発インストールできます。Python 3.7以上が必要ですので、バージョンを 事前に確認しておきましょう。

# terminal
# websocketsライブラリをインストール
pip install websockets

# インストール確認
python -c "import websockets; print(websockets.__version__)"

手順2:asyncioとの組み合わせ方

websocketsはPython標準の非同期ライブラリ asyncio と組み合わせて使います。初めて見ると「async defって何?」となりますよね。async defは 「非同期関数の定義」で、awaitは「この処理が終わるまで待つ(ただしその間他の処理を止めない)」という意味です。つまり、待ち時間をムダにせず複数の処理を効率よく並行実行できる仕組みです。

基本的に方法Aがおすすめです。この後の説明も、方法Aを元に行います。

方法A:asyncio.run()で実行(おすすめ)

Python 3.7以降で使える最もシンプルな実行方法です。asyncio.run()に非同期関数を渡すだけでイベントループが自動で管理されます。

# hello_ws.py
import asyncio
import websockets

# async def = 非同期関数の定義
async def hello():
    # wss://echo.websocket.org はエコーサーバー(送ったものが返ってくる)
    async with websockets.connect("wss://echo.websocket.org") as ws:
        await ws.send("Hello WebSocket!")   # await = 完了まで待つ
        reply = await ws.recv()
        print(f"受信: {reply}")

# asyncio.run() でイベントループを起動
asyncio.run(hello())

方法B:Jupyter Notebookで試す(お試し向け)

Jupyter環境ではすでにイベントループが動いているため、asyncio.run()の代わりにawaitを直接使います。動作確認にサクッと使えます。

# jupyter_ws.ipynb (セル内)
import websockets

# Jupyter では await を直接呼び出せる
async with websockets.connect("wss://echo.websocket.org") as ws:
    await ws.send("test")
    print(await ws.recv())

手順3:メッセージ送受信の基本パターン

実際のアプリでは、接続後にループでメッセージを受信し続けるパターンが基本です。以下のコードがそのひな形になります。async forを使うと、接続が切れるまで自動でメッセージを 受け取り続けてくれます。

# ws_basic.py
import asyncio
import websockets

async def receive_loop(uri: str):
    async with websockets.connect(uri) as ws:
        # async for = 接続が続く限りメッセージを受信し続けるループ
        async for message in ws:
            print(f"[受信] {message}")

asyncio.run(receive_loop("wss://echo.websocket.org"))

このひな形を覚えておくだけで、あとは受信した message を自由に加工するだけです。次のセクションでBinanceの実例を見てみましょう。

実例①:Binanceオーダーブックをリアルタイム表示

Binance WebSocket APIのエンドポイント仕様

BinanceはAPIキー不要で使えるWebSocket Market Streams(公開ストリーム)を提供しています。ただ、エンドポイントのURLが独特で最初は戸惑うかもしれません。そこも含めて、実際のコードで 確認してみましょう!

今回はBTC/USDTのオーダーブック差分ストリーム(depth stream)を使うことにしました。

Binanceのdepth streamは、注文板(売り・買いの価格と数量の一覧)が更新されるたびに差分データをpushしてくれます。つまり、常に最新の板情報をリアルタイムで受け取れます。

今回は、その受信データを整形してターミナルに表示するアプリを作成しました!最終的に、以下のようにリアルタイム更新される表示まで実装できます。

今回実装するオーダーブック表示ツールの機能

  • APIキー不要
    BinanceのPublic Streamsを使うため、アカウント登録なしで動作します。
  • リアルタイム差分受信
    価格・数量の更新があるたびに即座にデータが届きます。
  • ターミナル上書き表示
    ANSIエスケープコードで画面をリフレッシュし、常に最新板を表示します。

実際に動かしてみると、ターミナルがビットコインの板情報でリアルタイムに更新されていく様子は圧巻です。REST APIのポーリングとは比べ物にならないスムーズさを実感できますよ。

受信データのパースと整形(実装コード)

Binanceから届くデータはJSON形式です。bids(買い注文)と asks(売り注文)がそれぞれ [[価格, 数量], ...] の配列で入っています。 json.loads()でパースして、見やすく整形しましょう。

ただ、これもコードを一から書く必要はなく、以下のひな形をベースに改造するだけです。

# binance_orderbook.py
import asyncio
import json
import websockets

# Binance公開ストリーム(BTCUSDT板情報・差分・APIキー不要)
BINANCE_WS_URL = "wss://stream.binance.com:9443/ws/btcusdt@depth"

def format_orders(orders: list, label: str, n: int = 5) -> str:
    """上位n件を整形して文字列で返す"""
    lines = [f"  {'価格':>12}  {'数量':>12}"]
    for price, qty in orders[:n]:
        lines.append(f"  {float(price):>12,.2f}  {float(qty):>12.5f}")
    return f"【{label}】\n" + "\n".join(lines)

async def watch_orderbook():
    async with websockets.connect(BINANCE_WS_URL) as ws:
        async for raw in ws:
            data = json.loads(raw)          # JSON文字列 → dict
            bids = data.get("b", [])        # 買い注文リスト
            asks = data.get("a", [])        # 売り注文リスト
            yield bids, asks                # 呼び出し元に渡す(非同期ジェネレータ)

受信データの中身は更新差分なので、本格的な板管理をする場合はローカルでオーダーブックを保持・マージする処理が必要になります。今回はシンプルに「受信差分をそのまま表示」する形で進めます。

ターミナルにリアルタイム表示する

受信データを整形したら、あとはターミナルに表示するだけです。\033[H\033[JというANSIエスケープコードを使うと、前の出力を消して上書き表示できます。つまり、データが届くたびに 画面がスムーズに更新されるリアルタイム表示の完成です。

# binance_display.py
import asyncio
import json
import os
import websockets

BINANCE_WS_URL = "wss://stream.binance.com:9443/ws/btcusdt@depth"

def clear():
    """ターミナル画面をクリアして上書き表示(ANSI)"""
    print("\033[H\033[J", end="")

def fmt(orders, label, n=5):
    lines = [f"  {'価格(USDT)':>14}  {'数量(BTC)':>12}"]
    for p, q in orders[:n]:
        lines.append(f"  {float(p):>14,.2f}  {float(q):>12.5f}")
    return f"─── {label} ───\n" + "\n".join(lines)

async def main():
    print("Binance BTC/USDT オーダーブック監視中... (Ctrl+Cで終了)")
    async with websockets.connect(BINANCE_WS_URL) as ws:
        async for raw in ws:
            data = json.loads(raw)
            bids = data.get("b", [])   # 買い板
            asks = data.get("a", [])   # 売り板
            if not bids and not asks:
                continue
            clear()
            print("=== BTC/USDT リアルタイムオーダーブック ===\n")
            print(fmt(asks[::-1], "売り (Asks)"))   # 高い順に表示
            print()
            print(fmt(bids, "買い (Bids)"))
            print()

asyncio.run(main())

このコードをそのまま実行するだけで、APIキー不要でビットコインのリアルタイム板情報がターミナルに流れ始めます。実際に動かしてみると、WebSocketの低遅延をリアルに体験できます。

実例②:WebSocket経由で暗号化コードを配信する

サーバー側:AES暗号化してコードをpush

「ツールの機能をサーバーから動的に配信したい」「クライアントにソースを見せずにコードを実行させたい」――そんなユースケースへの応用例です。サーバー側ではAES(Advanced Encryption Standard) で実行コードを暗号化し、WebSocketでクライアントにpushします。つまり、配信内容を通信傍受されても暗号化済みなので中身が読めません。概念的なフローは以下のとおりです。

# server_concept.py(概念コード・実運用時はセキュリティ設計を別途行うこと)
import asyncio
import websockets
from Crypto.Cipher import AES  # pip install pycryptodome
import base64, os, json

SECRET_KEY = os.environ["AES_KEY"].encode()  # 32バイトのAESキー(環境変数から取得)

def encrypt_code(code: str) -> str:
    """Pythonコード文字列をAES-GCMで暗号化しBase64返却"""
    cipher = AES.new(SECRET_KEY, AES.MODE_GCM)
    ct, tag = cipher.encrypt_and_digest(code.encode())
    return base64.b64encode(cipher.nonce + tag + ct).decode()

async def handler(websocket):
    code = 'print("Hello from server!")'   # 配信したいコード
    payload = json.dumps({"type": "code", "data": encrypt_code(code)})
    await websocket.send(payload)

asyncio.run(websockets.serve(handler, "localhost", 8765))

AESキーはハードコードせず、必ず環境変数や安全なキー管理サービスから取得するようにしてください。

クライアント側:復号してeval実行

クライアント側では受信したBase64データを復号し、Pythonの exec() で実行します。たとえばSaaS型ツールで「機能ライセンスが有効な間だけコードを配信する」仕組みに応用できます。 注意点として、exec()は受け取ったコードをそのまま実行するため、信頼できるサーバーからのデータのみに使う設計にすることが大前提です。

# client_concept.py(概念コード)
import asyncio
import websockets
from Crypto.Cipher import AES
import base64, os, json

SECRET_KEY = os.environ["AES_KEY"].encode()

def decrypt_code(b64data: str) -> str:
    """Base64データをAES-GCMで復号してコード文字列を返す"""
    raw = base64.b64decode(b64data)
    nonce, tag, ct = raw[:16], raw[16:32], raw[32:]
    cipher = AES.new(SECRET_KEY, AES.MODE_GCM, nonce=nonce)
    return cipher.decrypt_and_verify(ct, tag).decode()

async def main():
    async with websockets.connect("ws://localhost:8765") as ws:
        msg = json.loads(await ws.recv())
        if msg.get("type") == "code":
            code = decrypt_code(msg["data"])
            exec(code)   # 復号したコードを実行

asyncio.run(main())

ライセンス管理・不正利用防止への応用

この仕組みを応用すると、サーバー側でライセンス認証を行い「有効なユーザーにのみコードをpushする」ライセンス管理システムが構築できます。たとえば、接続時にユーザーIDとライセンスキーを送信させ、 サーバーがDBで検証してからコードを配信するフローです。不正コピーを防ぐためには、ライセンスキーをデバイスIDと紐づけたり、配信コードに有効期限トークンを含める設計が有効です。

本番運用のための安定化テクニック

WebSocketは便利な反面、ネットワーク瞬断・サーバー再起動・タイムアウトなどで接続が突然切れることがあります。本番環境で安定稼働させるには、以下の3点を必ず実装しましょう。

Heartbeat(ping/pong)の実装

WebSocket接続は無通信が続くとルーターやロードバランサーに「死んだ接続」と判断されて切断されることがあります。これを防ぐのがHeartbeatです。つまり、定期的に小さなping信号を送り続けることで 「接続は生きています」と通知し続けます。websocketsライブラリは ping_interval と ping_timeout パラメータで自動Heartbeatを設定できます。

# ws_heartbeat.py
import asyncio
import websockets

async def main():
    # ping_interval=20: 20秒ごとにpingを送信
    # ping_timeout=10 : 10秒以内にpongが返らなければ切断と判定
    async with websockets.connect(
        "wss://stream.binance.com:9443/ws/btcusdt@trade",
        ping_interval=20,
        ping_timeout=10,
    ) as ws:
        async for msg in ws:
            print(msg[:80], "...")  # 先頭80文字だけ表示

asyncio.run(main())

例外ハンドリングと自動再接続ロジック

接続が切れたとき、そのまま落ちるのではなく自動で再接続する仕組みが本番では必須です。指数バックオフ(再接続間隔を徐々に延ばす手法)を組み込むと、サーバー障害時にリクエストが集中して さらに負荷をかける「再接続ストーム」を防げます。

# ws_reconnect.py
import asyncio
import websockets

URI = "wss://stream.binance.com:9443/ws/btcusdt@trade"

async def connect_with_retry():
    wait = 1  # 初回再接続待機秒数
    while True:
        try:
            async with websockets.connect(URI, ping_interval=20, ping_timeout=10) as ws:
                wait = 1  # 接続成功したらバックオフをリセット
                print("[INFO] 接続成功")
                async for msg in ws:
                    print(msg[:60])
        except (websockets.ConnectionClosed, OSError) as e:
            print(f"[WARN] 切断: {e}. {wait}秒後に再接続...")
            await asyncio.sleep(wait)
            wait = min(wait * 2, 60)  # 最大60秒まで待機を延ばす(指数バックオフ)

asyncio.run(connect_with_retry())

接続ステータスのログ管理

障害発生時の原因調査のために、接続・切断・再接続のイベントをログファイルに残す習慣をつけましょう。Pythonの標準 logging モジュールを使えば、ログレベル・タイムスタンプ付きで 管理できます。

# ws_logging.py
import asyncio
import logging
import websockets

# ログ設定:INFOレベル以上をファイルとコンソールに出力
logging.basicConfig(
    level=logging.INFO,
    format="%(asctime)s [%(levelname)s] %(message)s",
    handlers=[
        logging.FileHandler("ws_status.log", encoding="utf-8"),
        logging.StreamHandler(),
    ],
)
logger = logging.getLogger(__name__)

URI = "wss://stream.binance.com:9443/ws/btcusdt@trade"

async def main():
    wait = 1
    while True:
        try:
            async with websockets.connect(URI) as ws:
                logger.info("WebSocket接続確立: %s", URI)
                wait = 1
                async for _ in ws:
                    pass  # 実際の処理をここに書く
        except Exception as e:
            logger.warning("切断検知: %s | %d秒後に再接続", e, wait)
            await asyncio.sleep(wait)
            wait = min(wait * 2, 60)

asyncio.run(main())

ログファイルを見れば「いつ・どれくらいの頻度で切断が起きているか」が一目でわかります。本番運用の品質を大きく上げる一手です。

まとめ:Python WebSocketで広がるリアルタイム開発

今回は、Python WebSocketを使ったリアルタイムデータ処理ツールの構築に挑戦してみました。

REST APIのポーリングも便利ですが、「Python WebSocket(websocketsライブラリ)」はさらに便利で、低レイテンシ・双方向通信・サーバーpushが一度の接続で実現できる ため、これは使わないのはもったいない!と感じました。

導入も簡単ですので、みなさんも今回の記事を参考に、ぜひ「Python WebSocket」を活用してみてください!

次のステップ

  • WebSocketで取得したデータをSQLiteやCSVに保存する方法は「Python asyncioで非同期DB書き込みを実装する」をご覧ください。
  • Binanceの高度なAPIをもっと活用したい方は「Binance APIでボット開発入門【注文・残高取得】」もおすすめです。
  • 暗号化・ライセンス管理をさらに深掘りしたい方は「PythonでAES暗号化ツールを作る【pycryptodome実践】」もチェックしてみてください。