コンテンツにスキップ

新しいブロックをリッスンする⚓︎

初級

blocks WS WebSocket チャネル は、新しい ブロック がチェーンに追加されるたびにリアルタイム通知を送信します。 /chain/height GET エンドポイントをポーリングする場合と比べ、WebSocket は API を繰り返し呼び出すオーバーヘッドなしに、発生した更新をプッシュします。

このチュートリアルでは、チャネルをサブスクライブして、受け取った更新を表示する方法を説明します。

ポーリングによる方法

ポーリングを使う方法については、チェーン高と不可逆高を照会する チュートリアルを参照してください。

前提条件⚓︎

NEM は SockJS 上で STOMP メッセージングプロトコルを使って WebSocket を提供するため、STOMP クライアントと WebSocket トランスポートが必要です。

stomperwebsockets ライブラリをインストールします。

pip install stomper websockets

@stomp/stompjssockjs-client ライブラリをインストールします。

npm install @stomp/stompjs sockjs-client

接続プロトコルの詳細については、WebSocket リファレンス を参照してください。

完全なコード⚓︎

import asyncio
import json
import os
import random
import uuid

import stomper
from websockets import connect

NODE_URL = os.getenv('NODE_URL', 'http://libertalia.nemtest.net:7890')
WS_URL = NODE_URL.replace(':7890', ':7778')
print(f'Using node {NODE_URL}')


# SockJS has no Python client library.
# These helpers wrap the raw WebSocket transport.
def sockjs_url(endpoint_url):
    # SockJS raw WebSocket transport adds a random server and session id
    server = random.randint(100, 999)
    session = uuid.uuid4().hex
    ws_base = endpoint_url.replace('http', 'ws', 1)
    return f'{ws_base}/{server}/{session}/websocket'


async def send_frame(websocket, frame):
    # SockJS wraps each client payload as a JSON array of frame strings
    await websocket.send(json.dumps([frame]))


async def stomp_connect(websocket):
    await websocket.recv()  # consume the SockJS open frame
    await send_frame(
        websocket, stomper.connect('', '', WS_URL, heartbeats=(0, 0)))


async def stomp_subscribe(websocket, destination, sub_id):
    await send_frame(websocket, stomper.subscribe(destination, sub_id))


async def stomp_unsubscribe(websocket, sub_id):
    await send_frame(websocket, stomper.unsubscribe(sub_id))


async def stomp_disconnect(websocket):
    await send_frame(websocket, stomper.disconnect())


def stomp_messages(raw_frame):
    # Yield the JSON body of each STOMP MESSAGE in a SockJS data frame
    if 'a' != raw_frame[0]:  # skip 'o' open, 'h' heartbeat, 'c' close
        return
    for payload in json.loads(raw_frame[1:]):
        frame = stomper.unpack_frame(payload)
        if 'MESSAGE' == frame['cmd']:
            yield json.loads(frame['body'])


async def main():
    # Open connection
    async with connect(sockjs_url(f'{WS_URL}/w/messages')) as websocket:
        await stomp_connect(websocket)
        print(f'Connected to {WS_URL}')

        # Subscribe to the new block channel
        destination = '/blocks'
        await stomp_subscribe(websocket, destination, 'id-0')
        print(f'Subscribed to {destination} channel')

        # Read and format each new block
        try:
            async for raw_frame in websocket:
                for block in stomp_messages(raw_frame):
                    print(
                        f'New block: height={block["height"]:,}'
                        f' harvester={block["signer"][:16].upper()}...'
                    )

        # Unsubscribe on exit
        finally:
            await stomp_unsubscribe(websocket, 'id-0')
            await stomp_disconnect(websocket)
            print('Unsubscribed and disconnected')


try:
    asyncio.run(main())
except KeyboardInterrupt:
    pass
except Exception as error:
    print(error)

Download source

import { Client } from '@stomp/stompjs';
import SockJS from 'sockjs-client';

const NODE_URL = process.env.NODE_URL ||
    'http://libertalia.nemtest.net:7890';
const WS_URL = NODE_URL.replace(':7890', ':7778');
console.log(`Using node ${NODE_URL}`);

// Open connection
const client = new Client({
    webSocketFactory: () => new SockJS(`${WS_URL}/w/messages`)
});
await new Promise(resolve => {
    client.onConnect = resolve;
    client.activate();
});
console.log(`Connected to ${WS_URL}`);

// Read and format each new block
function formatBlock(message) {
    const block = JSON.parse(message.body);
    console.log(
        `New block: height=${block.height.toLocaleString()}` +
        ` harvester=${block.signer.substring(0, 16).toUpperCase()}...`
    );
}

// Subscribe to the new block channel
const destination = '/blocks';
const subscription = client.subscribe(destination, formatBlock, {
    id: 'id-0'
});
console.log(`Subscribed to ${destination} channel`);

// Unsubscribe on exit
process.on('SIGINT', () => {
    subscription.unsubscribe();
    client.deactivate();
    console.log('Unsubscribed and disconnected');
    process.exit(0);
});

Download source

スニペットでは、NODE_URL 環境変数を使って NEM ノード を指定します。 値が指定されていない場合は、デフォルト値を使用します。

WS_URL は同じノードの WebSocket エンドポイントを定義します。 NODE_URL のポート 7890(デフォルトの HTTP API ポート)を 7778(デフォルトの NIS WebSocket ポート)に置き換えて導出します。

プログラムは Ctrl+C で割り込みを受けるまで実行され、接続を閉じる前にサブスクライブ解除の手順が実行されます。

Python の SockJS ヘルパー

Python 用の SockJS クライアントライブラリはないため、便宜上、ファイルの先頭に小さなヘルパーメソッドをいくつか定義しています。

コードの説明⚓︎

WebSocket に接続する⚓︎

    # Open connection
    async with connect(sockjs_url(f'{WS_URL}/w/messages')) as websocket:
        await stomp_connect(websocket)
        print(f'Connected to {WS_URL}')
// Open connection
const client = new Client({
    webSocketFactory: () => new SockJS(`${WS_URL}/w/messages`)
});
await new Promise(resolve => {
    client.onConnect = resolve;
    client.activate();
});
console.log(`Connected to ${WS_URL}`);

最初の手順では、ノードの /w/messages エンドポイントへの接続を開き、その上で STOMP セッション を開始します。

チャネルをサブスクライブする⚓︎

        # Subscribe to the new block channel
        destination = '/blocks'
        await stomp_subscribe(websocket, destination, 'id-0')
        print(f'Subscribed to {destination} channel')
// Subscribe to the new block channel
const destination = '/blocks';
const subscription = client.subscribe(destination, formatBlock, {
    id: 'id-0'
});
console.log(`Subscribed to ${destination} channel`);

コードは blocks WS チャネルをサブスクライブします。 ノードは新しいブロックがチェーンに追加されるたびに、サブスクライバーへ通知します(およそ 1 分ごとです)。

通知は常に均等な間隔で届くとは限りません。 例えば、ノードがピアに追いつくために一度に複数のブロックを追加すると、チャネルはブロックごとに 1 通知を、短い間隔でまとめて配信します。

サブスクリプションには idid-0)が付与されます。終了時のサブスクライブ解除に使います。

受信した各メッセージは、以下の整形処理に渡されます。

メッセージを整形する⚓︎

        # Read and format each new block
        try:
            async for raw_frame in websocket:
                for block in stomp_messages(raw_frame):
                    print(
                        f'New block: height={block["height"]:,}'
                        f' harvester={block["signer"][:16].upper()}...'
                    )
// Read and format each new block
function formatBlock(message) {
    const block = JSON.parse(message.body);
    console.log(
        `New block: height=${block.height.toLocaleString()}` +
        ` harvester=${block.signer.substring(0, 16).toUpperCase()}...`
    );
}

受信する各メッセージの本文は、ブロック情報です。Block スキーマに従います。

スニペットでは、各メッセージから次の 2 フィールドを表示します。

NEM のブロックは、このペイロードに自身のハッシュを含まないため、このチュートリアルでは各ブロックを height で識別し、ハーベスターの signer も表示します。

新しいブロックはまだ確定していません

新しいブロックはすでにチェーンの一部ですが、まだ不可逆ではありません。 後続ブロックが 書き換え制限 を超えるのに十分な数だけ追加されるまでは、ロールバック が発生する可能性があります。 ロールバック後は、受信済みのブロックを置き換えるブロックが届くため、チャネルがすでに受け取ったブロックより低い高さのブロックを報告することもあります。

チェーン高と不可逆高を照会する チュートリアルでは、不可逆高の計算方法を説明しています。

終了時にサブスクライブ解除する⚓︎

        # Unsubscribe on exit
        finally:
            await stomp_unsubscribe(websocket, 'id-0')
            await stomp_disconnect(websocket)
            print('Unsubscribed and disconnected')
// Unsubscribe on exit
process.on('SIGINT', () => {
    subscription.unsubscribe();
    client.deactivate();
    console.log('Unsubscribed and disconnected');
    process.exit(0);
});

プログラムが割り込み(Ctrl+C)を受けると、コードはチャネルのサブスクライブを解除し、接続を閉じる前に STOMP セッションを終了します。 これにより、ノードからきれいに切断できます。

出力⚓︎

以下の出力は、新しいブロックをリッスンした場合の実行例です。

1
2
3
4
5
6
7
8
9
Using node http://libertalia.nemtest.net:7890
Connected to http://libertalia.nemtest.net:7778
Subscribed to /blocks channel
New block: height=682,352 harvester=95BA0E864CD5EDB3...
New block: height=682,353 harvester=95BA0E864CD5EDB3...
New block: height=682,354 harvester=7C814B3B44B1A98B...
New block: height=682,355 harvester=95BA0E864CD5EDB3...
New block: height=682,356 harvester=95BA0E864CD5EDB3...
Unsubscribed and disconnected

出力には次の内容が表示されます。

  • 接続(2 行目): ノードのポート 7778 の WebSocket エンドポイント上で STOMP セッションが確立されます。
  • サブスクリプション(3 行目): /blocks チャネルをサブスクライブします。
  • 新しいブロック(4~8 行目): およそ 1 分ごとに新しいブロック通知が届きます。
  • サブスクライブ解除(9 行目): Ctrl+C でコードがサブスクライブを解除して切断します。

まとめ⚓︎

このチュートリアルでは、次の方法を説明しました。

手順 関連ドキュメント
ブロックチャネルをサブスクライブする blocks WS
ブロックメッセージを整形する Block