はじめに
前回の第2回では、Orleans のステート管理と永続化を取り上げ、グレインが自身の状態を安全に保持する仕組みを見てきました。本稿では、複数のグレインが協調して動くための「グレイン間通信」と、時系列に発生するイベントを扱う「Orleans Streams」を解説します。単体で完結するグレインではなく、注文処理やリアルタイム集計のように、多数のアクターが役割を分担しながら連携するアプリケーションを設計するうえで欠かせない要素です。
対象とするのは .NET 10 世代の Orleans です。メソッド呼び出しによる同期的な連携と、ストリームを介した疎結合な pub/sub の両方を扱い、それぞれがどのような場面に向くのかを整理します。インメモリから Azure Event Hubs までのプロバイダの選び方、明示的・暗黙的サブスクリプションの違い、そして信頼性とバックプレッシャの考え方まで、実務で判断に迷いやすい点を中心に順を追って見ていきます。

グレイン間のメソッド呼び出し
グレイン間の最も基本的な連携は、インターフェース越しのメソッド呼び出しです。あるグレインが GrainFactory.GetGrain<T>() で相手の参照を取得し、その非同期メソッドを await するだけで、呼び出し先がクラスタ内のどのサイロで動いているかを意識せずに通信できます。物理的な配置を隠蔽するこの性質は「位置透過性」と呼ばれ、Orleans がスケールアウトを容易にしている中核の仕組みです。
public interface IOrderGrain : IGrainWithGuidKey
{
Task<OrderStatus> PlaceOrderAsync(OrderRequest request);
}
public interface IInventoryGrain : IGrainWithStringKey
{
Task<bool> TryReserveAsync(IReadOnlyList<OrderItem> items, Guid orderId);
Task ReleaseAsync(Guid orderId);
}
public sealed class OrderGrain : Grain, IOrderGrain
{
private readonly IPersistentState<OrderState> _state;
public OrderGrain(
[PersistentState("order", "orderStore")] IPersistentState<OrderState> state)
{
_state = state;
}
public async Task<OrderStatus> PlaceOrderAsync(OrderRequest request)
{
_state.State.OrderId = this.GetPrimaryKey();
_state.State.Items = request.Items;
// 別グレインへのメソッド呼び出し。呼び出し先がどのサイロにあるかは意識しない
var inventory = GrainFactory.GetGrain<IInventoryGrain>("main");
var reserved = await inventory.TryReserveAsync(request.Items, _state.State.OrderId);
if (!reserved)
{
_state.State.Status = OrderStatus.Rejected;
await _state.WriteStateAsync();
return OrderStatus.Rejected;
}
_state.State.Status = OrderStatus.Confirmed;
await _state.WriteStateAsync();
return OrderStatus.Confirmed;
}
}
この呼び出しは要求と応答が対になった同期的なやり取りで、戻り値を待ってから次の処理へ進みます。注文グレインが在庫グレインへ予約を依頼し、その成否を受けて確定するといった、結果に依存するワークフローに向いた形です。ただし Orleans のグレインは既定で単一スレッド・非再入実行のため、相互に結果を待ち合う循環した呼び出しはデッドロックを招くことがあります。依存の向きを一方向に保つか、どうしても再入が必要な箇所にのみ [Reentrant] を検討するのが安全な設計です。
Orleans Streams の基本
メソッド呼び出しが一対一の同期的な連携だとすれば、Orleans Streams はイベントを介した一対多の疎結合な連携を担います。中心にあるのは「仮想ストリーム」という考え方です。ストリームは名前空間とキーからなる StreamId で識別され、あらかじめ実体を作成しておく必要はありません。プロデューサが発行を始めた時点で論理的に存在し、コンシューマはいつでも同じ StreamId を指して購読できます。
プロデューサはストリームへ OnNextAsync でイベントを流し込み、コンシューマは購読して受け取ります。両者は互いを直接参照せず、間にストリームプロバイダが立つことで完全に切り離されています。この pub/sub 構造により、送り手を変えずに受け手を増減でき、クラスタ全体へイベントを波及させる処理を素直に表現できます。
public sealed class SensorGrain : Grain, ISensorGrain
{
private IAsyncStream<Reading> _stream = null!;
public override Task OnActivateAsync(CancellationToken cancellationToken)
{
var provider = this.GetStreamProvider("StreamProvider");
// StreamId は 名前空間 + キー で仮想ストリームを一意に指す
var streamId = StreamId.Create("readings", this.GetPrimaryKeyString());
_stream = provider.GetStream<Reading>(streamId);
return base.OnActivateAsync(cancellationToken);
}
// プロデューサはイベントを発行するだけ。購読者が何個いるかは意識しない
public Task PublishAsync(Reading reading) => _stream.OnNextAsync(reading);
}
ストリームプロバイダの選択
ストリームの配送を実際に担うのがストリームプロバイダです。用途に応じて複数を使い分けられ、サイロ構成時に登録します。開発やテストにはインメモリのプロバイダが手軽で、外部依存なしにストリームの挙動を確認できます。購読状態を保持する PubSubStore 用のストレージを併せて登録しておく点に注意します。
var builder = Host.CreateApplicationBuilder(args);
builder.UseOrleans(silo =>
{
silo.UseLocalhostClustering();
// サブスクリプション状態(PubSub)を保持するストレージ
silo.AddMemoryGrainStorage("PubSubStore");
// 開発・テスト向けのインメモリストリーム
silo.AddMemoryStreams("StreamProvider");
});
await builder.Build().RunAsync();
本番では永続的なメッセージ基盤に接続します。毎秒大量のイベントを取り込み、あとから再生する要件には Azure Event Hubs が向きます。より単純で安価なキュー配送で足りる場合は Azure Queue を選べます。いずれもアプリケーション側のプロデューサ・コンシューマのコードは変えず、登録するプロバイダを差し替えるだけで配送基盤を切り替えられます。
// Azure Event Hubs: 高スループットな取り込みと再生に向く
silo.AddEventHubStreams("StreamProvider", configurator =>
{
configurator.ConfigureEventHub(ob => ob.Configure(options =>
{
options.ConfigureEventHubConnection(connectionString, hubName, consumerGroup);
}));
configurator.UseAzureTableCheckpointer(ob => ob.Configure(options =>
{
options.TableServiceClient = new TableServiceClient(storageConnectionString);
}));
});
// Azure Queue: シンプルで安価なキュー配送
silo.AddAzureQueueStreams("StreamProvider", configurator =>
{
configurator.ConfigureAzureQueue(ob => ob.Configure(options =>
{
options.QueueServiceClient = new QueueServiceClient(storageConnectionString);
}));
});
明示的・暗黙的サブスクリプション
購読の張り方には二つの方式があります。明示的サブスクリプションでは、コンシューマ側が SubscribeAsync を呼んで自ら購読を登録し、返される StreamSubscriptionHandle を通じて後から解除もできます。購読の開始と終了をコードで制御したい場合に適しています。
public sealed class DashboardGrain : Grain, IDashboardGrain
{
private StreamSubscriptionHandle<Reading>? _handle;
public override async Task OnActivateAsync(CancellationToken cancellationToken)
{
var provider = this.GetStreamProvider("StreamProvider");
var streamId = StreamId.Create("readings", this.GetPrimaryKeyString());
var stream = provider.GetStream<Reading>(streamId);
// 明示的サブスクリプション: 購読するグレインが自分で購読を張る
_handle = await stream.SubscribeAsync(OnNextAsync);
await base.OnActivateAsync(cancellationToken);
}
private Task OnNextAsync(Reading reading, StreamSequenceToken? token)
{
// 受信したイベントを処理する
return Task.CompletedTask;
}
}
暗黙的サブスクリプションは、グレインのクラスに [ImplicitStreamSubscription] を付け、対象の名前空間を宣言する方式です。その名前空間に、グレインのキーと一致する StreamId でイベントが流れると、対応するグレインが自動的に起動して受け取ります。プロデューサは購読者の存在を一切意識せず、発行するだけで処理が始まります。デバイスごと・ユーザーごとに処理を割り当てるような、キーで宛先が定まるファンアウトに向いた形です。
// 名前空間 "readings" に流れたイベントを、キーが一致するこのグレインが自動で受け取る
[ImplicitStreamSubscription("readings")]
public sealed class ReadingProcessorGrain
: Grain, IReadingProcessorGrain, IStreamSubscriptionObserver
{
public Task OnSubscribed(IStreamSubscriptionHandleFactory handleFactory)
{
var handle = handleFactory.Create<Reading>();
return handle.ResumeAsync(OnNextAsync);
}
private Task OnNextAsync(Reading reading, StreamSequenceToken? token)
{
// プロデューサ側は購読者を知らない。イベント発行だけで処理が起動する
return Task.CompletedTask;
}
}
信頼性とバックプレッシャ
ストリームの配送保証はプロバイダによって異なります。インメモリプロバイダは軽量な反面、サイロ障害時にイベントを失う可能性があります。Event Hubs や Azure Queue のような永続プロバイダは、配送されるまでイベントを保持し、コンシューマ側の障害から回復したあとに処理を再開できます。要件が可用性やデータ欠損の許容度に依存するため、プロバイダの選択はそのまま信頼性の設計になります。
永続プロバイダでは各イベントに StreamSequenceToken が付与され、どこまで処理したかを表す位置として使えます。コンシューマはこのトークンを指定して購読を再開し、中断した箇所からイベントを読み直せます。加えて、コンシューマの処理が追いつかないときには、プロバイダがキューにイベントを滞留させることで自然にバックプレッシャがかかります。プロデューサが一方的に押し込み続けて処理側が溢れる事態を、キュー基盤が緩衝することで防ぎます。
ユースケース: リアルタイム集計と通知
これらを組み合わせると、リアルタイム集計が素直に書けます。注文イベントが流れる名前空間を暗黙的サブスクリプションで受け取る集計グレインを用意し、受信のたびに合計や件数を積み上げます。そのうえでグレインタイマーを使い、一定間隔で集計結果を下流のストリームへ発行すれば、絶えず更新される売上サマリを配信できます。
[ImplicitStreamSubscription("orders")]
public sealed class SalesAggregatorGrain
: Grain, ISalesAggregatorGrain, IStreamSubscriptionObserver
{
private decimal _total;
private int _count;
public override Task OnActivateAsync(CancellationToken cancellationToken)
{
// 一定間隔で集計結果を下流へ通知する
this.RegisterGrainTimer(
PublishSummaryAsync,
new GrainTimerCreationOptions(
dueTime: TimeSpan.FromSeconds(10),
period: TimeSpan.FromSeconds(10)));
return base.OnActivateAsync(cancellationToken);
}
public Task OnSubscribed(IStreamSubscriptionHandleFactory factory)
=> factory.Create<OrderPlaced>().ResumeAsync(OnOrderAsync);
private Task OnOrderAsync(OrderPlaced order, StreamSequenceToken? token)
{
_total += order.Amount;
_count++;
return Task.CompletedTask;
}
private async Task PublishSummaryAsync(CancellationToken cancellationToken)
{
var provider = this.GetStreamProvider("StreamProvider");
var summary = provider.GetStream<SalesSummary>(
StreamId.Create("summary", "sales"));
await summary.OnNextAsync(new SalesSummary(_count, _total, DateTime.UtcNow));
}
}
通知の配信も同じ pub/sub で表現できます。あるグレインがイベントを発行すると、暗黙的サブスクリプションで結びついた宛先グレインが起動し、ユーザーやデバイス単位で処理を分散します。SignalR などのリアルタイム通信と組み合わせれば、サーバー内部のイベントをそのままクライアントへ届ける経路を、送り手と受け手を疎に保ったまま構築できます。
まとめ
本稿では、グレイン間のメソッド呼び出しと Orleans Streams による pub/sub を通じて、Orleans での連携の設計を見てきました。結果を待ち合わせる同期的な処理にはメソッド呼び出しを、送り手と受け手を切り離したイベント配送にはストリームを選ぶという使い分けが基本になります。プロバイダの選択が信頼性とバックプレッシャの性質を決めること、そして明示的・暗黙的サブスクリプションで購読の制御粒度を選べることを押さえておけば、リアルタイム処理の多くを無理なく組み立てられます。
次回の第4回では、クラスタリングと高可用性を取り上げ、マルチノード構成やフェイルオーバーを含む本番運用の設計を解説します。
エンハンスド株式会社では、Orleans を用いた分散システムやリアルタイム処理基盤の設計・実装を支援しています。既存システムのスケール課題の解消から、ストリーム処理を含むアーキテクチャの新規設計、Azure 上での本番運用まで、実務に即した形でご相談を承ります。.NET でのスケーラブルなシステム構築を検討されている方は、お気軽にお問い合わせください。
このテーマの全体像は .NET モダナイゼーション完全ガイド にまとめています。
この記事をシェア
関連記事

【.NET Orleans入門】第2回 ステート管理と永続化 - 信頼性の高いステートフルサービスの構築
第1回では、.NET Orleans の全体像として、バーチャルアクターモデル、グレインとサイロという中心概念、そして高スケールなワークロードに向く理由を整理しました。…

【.NET Orleans入門】第6回 実践ケーススタディ - 注文と在庫のリアルタイム処理システム
「.NET Orleans入門」も今回で第6回、シリーズの最終回を迎えます。これまで、バーチャルアクターモデルの基礎(第1回)、状態管理と永続化(第2回)、グレイン間通信とストリーミング(第3回)、…

Orleansで金融取引基盤を建て直した話 ある証券系プロジェクトの現場から
金融のリアルタイム基盤は、平常時の姿だけを見ていると意外なほど静かに動いている。負荷が本当の意味で問われるのは、相場が動いた一瞬だ。ここで紹介するのは、…
