注: 以下の翻訳の正確性は検証されていません。AIPを利用して英語版の原文から機械的に翻訳されたものです。

WebSocket 経由でオブジェクトセットの変更を購読する

TypeScript オントロジー SDK の .subscribe メソッドは、WebSocket プロトコルを使用してオブジェクトの更新をクライアントにストリーミングします。また、Object Set Watcher エンドポイント(/api/v2/ontologySubscriptions/ontologies/{ontology}/streamSubscriptions)に直接接続して、オブジェクトセットを購読することもできます。

エンドポイント

wss://{foundryUrl}/api/v2/ontologySubscriptions/ontologies/{ontology}/streamSubscriptions

{foundryUrl} を Foundry インスタンスの URL(https:// の接頭辞を除く)に、{ontology} をオントロジーのAPI名に置き換えます。

認証

接続の認証は、WebSocket のサブプロトコルとして "Bearer-{token}" を渡すことで行われます。トークンの前には、REST API で使用する空白の代わりに - を使用します。

メッセージ形式

通信は JSON にシリアライズされたメッセージを通じて行われます。クライアントは購読リクエストを送信します。サーバーは、購読の確認またはエラーに加え、購読中のオブジェクトセットに対するオブジェクトの更新を返します。

オブジェクトセットを購読する

購読を開始するには、メッセージを送信します。

Copied!
1 2 3 4 5 6 7 8 9 10 11 12 13 { "id": "550e8400-e29b-41d4-a716-446655440000", "requests": [ { "objectSet": { "type": "base", "objectType": "Country" }, "propertySet": ["population", "countryName"], "referenceSet": [] } ] }
フィールドデータ型必須説明
idUUIDはいレスポンスと対応付けるための一意のリクエストID
requestsarrayはい購読するオブジェクトセットのリスト
requests[].objectSetobjectはいオブジェクトセットの定義
requests[].propertySetstring[]いいえ更新に含めるプロパティのAPI名。省略すると、すべてのプロパティが返されます。
requests[].referenceSetstring[]いいえ購読する参照プロパティのAPI名(たとえば、地理時系列)
requests[].objectLoadingResponseOptionsobjectいいえ任意のレスポンス設定
requests[].objectLoadingResponseOptions.shouldLoadObjectRidsbooleanいいえレスポンスにオブジェクトの RID を含めるかどうか(デフォルト: false)。RID を含めると遅延が増えるため、必要な場合にのみ有効にしてください。

標準のオブジェクトセット構文を使用して、オブジェクトセットを絞り込めます。

Copied!
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 { "id": "550e8400-e29b-41d4-a716-446655440001", "requests":[ { "objectSet": { "type": "filter", "objectSet": { "type": "base", "objectType": "Country" }, "where": { "type": "gte", "field": "population", "value": 1000000 } }, "propertySet": ["population"] } ] }

購読を管理する

新しい購読メッセージを送信すると、アクティブな購読が更新されます。サーバーは新しいリクエストリストを既存の購読と比較します。一致する購読は開いたまま維持され、リストにない購読は閉じられ、新しい購読は開かれます。

すべてのオブジェクトセットの購読を解除するには、空の requests 配列を送信します。

Copied!
1 2 3 4 { "id": "550e8400-e29b-41d4-a716-446655440002", "requests": [] }

オブジェクトセットの更新

クライアントは、購読のライフサイクル全体を通じてさまざまな種類のメッセージを受信します。

subscribeResponses

購読リクエストを受信した後に送信され、リクエストされた購読ごとのレスポンスが含まれます。

Copied!
1 2 3 4 5 6 7 8 9 10 { "type": "subscribeResponses", "id": "550e8400-e29b-41d4-a716-446655440000", "responses": [ { "type": "success", "id": "a1b2c3d4-e5f6-7890-abcd-ef1234567890" } ] }

成功レスポンスには、その購読の更新を識別する購読IDが含まれます。エラーレスポンスには診断情報が含まれます。

Copied!
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 { "type": "subscribeResponses", "id": "550e8400-e29b-41d4-a716-446655440000", "responses": [ { "type": "error", "errors": [ { "error": "INVALID_OBJECT_TYPE", "args": [{ "name": "objectType", "value": "InvalidType" }] } ] } ] }

type が "qos" の購読レスポンスを受信した場合、サーバーに高い負荷がかかっています。ジッター付きの指数バックオフを使用してリクエストを再試行することを推奨します。

objectSetChanged

購読中のオブジェクトセット内のオブジェクトが追加、更新、または削除されたときに送信されます。

Copied!
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 { "type": "objectSetChanged", "id": "a1b2c3d4-e5f6-7890-abcd-ef1234567890", "updates": [ { "type": "object", "object": { "__apiName": "Country", "__primaryKey": "US", "countryName": "United States", "population": 331000000 }, "state": "ADDED_OR_UPDATED" } ] }
フィールド説明
objectリクエストされたプロパティを含む、更新されたオブジェクト
object.__apiNameオブジェクトタイプのAPI名
object.__primaryKeyオブジェクトの主キーの値
state"ADDED_OR_UPDATED" または "REMOVED"

state が "REMOVED" の場合、オブジェクトは削除されたか、オブジェクトセットのフィルターに一致しなくなりました。

参照プロパティ(地理時系列など)を含む購読では、参照プロパティの更新を受信することがあります。

Copied!
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 { "type": "objectSetChanged", "id": "a1b2c3d4-e5f6-7890-abcd-ef1234567890", "updates": [ { "type": "reference", "objectType": "Aircraft", "primaryKey": { "aircraftId": "AC-123" }, "property": "currentLocation", "value": { "type": "geotimeSeriesValue", "position": { "type": "Point", "coordinates": [-122.4194, 37.7749] }, "timestamp": "2024-01-15T10:30:00Z" } } ] }

refreshObjectSet

サーバーがインクリメンタル(差分)更新を提供できないことを示します。標準のオブジェクトセットを読み込むエンドポイントへの通常の HTTP リクエストで、オブジェクトセットを再読み込みしてください。

Copied!
1 2 3 4 5 { "type": "refreshObjectSet", "id": "a1b2c3d4-e5f6-7890-abcd-ef1234567890", "objectType": "Country" }

subscriptionClosed

購読が閉じられたことを示します。これは、エラーなど、さまざまな理由で発生することがあります。

Copied!
1 2 3 4 5 6 7 8 9 { "type": "subscriptionClosed", "id": "a1b2c3d4-e5f6-7890-abcd-ef1234567890", "cause": { "type": "error", "error": "SUBSCRIPTION_MEMORY_LIMIT_EXCEEDED", "args": [] } }

リクエストリストから購読を削除した場合にも、購読は閉じられます。

Copied!
1 2 3 4 5 6 7 8 { "type": "subscriptionClosed", "id": "a1b2c3d4-e5f6-7890-abcd-ef1234567890", "cause": { "type": "reason", "reason": "USER_CLOSED" } }

制限事項

  • 結合を使用して構築されたオブジェクトセットには、完全には対応していません。サーバーは選択された部分を監視し、変更されたときに通知を送信しますが、完全性は保証できません。
  • 購読はサーバーの再起動後には維持されません。クライアントに再接続のロジックを実装してください。
  • メモリーの上限は購読ごとに適用されます。多数のオブジェクトを追跡する購読は、上限を超えると閉じられることがあります。
  • 購読の絞り込みでは、一部の高度な全文検索フィルターに対応していません。
  • インターフェースのオブジェクトセットを購読する場合、オブジェクトの更新は基となるオブジェクトタイプとして返されます。更新を処理する際は、オブジェクトのプロパティをインターフェースのプロパティに再マッピングする必要があります。

インターフェース

インターフェースのオブジェクトセットの購読を処理するには、オブジェクトのプロパティタイプから、対応するインターフェースのプロパティタイプへのマッピングを取得する必要があります。これは、オブジェクトタイプの完全なメタデータを読み込むエンドポイントで行えます。このベータ版エンドポイントは、実装されているすべてのインターフェースのプロパティマッピングを含む、リクエストされたオブジェクトタイプの完全なメタデータを返します。

オブジェクトタイプのAPI名には、更新されたオブジェクトの "__apiName" キーでアクセスできます。API名が interfaceType のインターフェースと、オブジェクトタイプの完全なメタデータを読み込むのレスポンス response がある場合、返されたインターフェースのマッピングには response["implementsInterfaces2"][interfaceType]["properties"] でアクセスできます。このマッピングは、オブジェクトのプロパティタイプのAPI名から、対応するインターフェースのプロパティタイプのAPI名へのマッピングです。

購読が存続する間、これらのプロパティマッピングをキャッシュすることを推奨します。

Python の例

以下は、絞り込まれたオブジェクトセットを購読し、更新を処理する方法を示す Python スクリプトのサンプルです。この例では、インターフェースのプロパティを再マッピングする方法は示していません。

Copied!
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 # /// script # requires-python = ">=3.10" # dependencies = [ # "websockets", # ] # /// import asyncio import json import os import uuid from websockets.asyncio.client import connect TOKEN = ... FOUNDRY_URL = ... ONTOLOGY_RID = ... async def main() -> None: filtered_object_set = { "type": "filter", "objectSet": { "type": "base", "objectType": "Country", }, "where": { "type": "gte", "field": "population", "value": 10000000, }, } subscribe_msg = { "id": str(uuid.uuid4()), "requests": [ { "objectSet": filtered_object_set, "propertySet": ["timeUntilNextFlight", "aircraftRegistration"], }, ], } open_subscriptions: set[str] = set() async with connect( f"wss://{FOUNDRY_URL}/api/v2/ontologySubscriptions/ontologies/{ONTOLOGY_RID}/streamSubscriptions", subprotocols=[f"Bearer-{TOKEN}"], # Bearer トークン認証 ) as websocket: await websocket.send(json.dumps(subscribe_msg)) async for message in websocket: match json.loads(message): case {"type": "subscribeResponses", "responses": [*responses]}: for response in responses: match response: case {"type": "success", "id": identifier}: print(f"Subscribed: ID {identifier}") open_subscriptions.add(identifier) case {"type": "error", "errors": [*errors]}: print(f"Subscription errors: {errors}") if not any(response["type"] == "success" for response in responses): print("All subscriptions failed. Exiting.") return case {"type": "objectSetChanged", "updates": [*updates]}: for update in updates: match update: case {"type": "object", "state": state, "object": {**obj}}: print(f"{state}: {obj}") case { "type": "reference", "property": prop, "primaryKey": primary_key, }: print(f"Reference update: {primary_key} <- {prop}") case {"type": "refreshObjectSet", "objectType": object_type}: print(f"Refresh required for object type: {object_type}") case { "type": "subscriptionClosed", "id": identifier, "cause": {**cause}, }: print(f"Subscription {identifier} was closed: {cause}") open_subscriptions.remove(identifier) if not open_subscriptions: print("All subscriptions closed. Exiting.") return if __name__ == "__main__": asyncio.run(main())