本文へスキップ
【.NET Orleans入門】第6回 実践ケーススタディ - 注文と在庫のリアルタイム処理システムのアイキャッチ画像
Architecture

【.NET Orleans入門】第6回 実践ケーススタディ - 注文と在庫のリアルタイム処理システム

公開: 更新: 約11分で読めます

はじめに

「.NET Orleans入門」も今回で第6回、シリーズの最終回を迎えます。これまで、バーチャルアクターモデルの基礎(第1回)、状態管理と永続化(第2回)、グレイン間通信とストリーミング(第3回)、クラスタリングと高可用性(第4回)、そしてパフォーマンスチューニング(第5回)と、Orleans を構成する要素を一つずつ見てきました。最終回では、それらをどう組み合わせて一つのシステムに仕立てるのかを、具体的な題材に沿って解説します。

取り上げるのは、EC サイトの注文と在庫をリアルタイムに処理するシステムです。セール開始直後のように、限られた在庫へ注文が集中する場面を想定します。個々の要素技術は前回までに触れていますので、本記事では設計判断の理由と、それらを一つのアーキテクチャへ結線していく過程に焦点を当てます。ランタイムは .NET 10 を前提とします。

注文と在庫をリアルタイム処理する Orleans システムの構成図。クライアントと API 層、複数サイロからなるクラスター上の注文グレインと SKU 単位の在庫グレイン、ホットな SKU のシャード分割、状態の永続化、確定注文を流すイベントストリームと集計用の読み取りモデルを示す。
注文グレインと在庫グレインをクラスターへ分散配置し、確定注文をストリームで読み取りモデルへ流す構成を表しています。

題材と要件

題材は、商品ごとの在庫を管理しながら注文を受け付けるサービスです。難しさの中心は、在庫という共有資源に対して多数の注文が同時に到達する点にあります。素朴に実装すると、複数のリクエストが同じ在庫数を読み込み、それぞれが引き当て可能と判断し、結果として在庫を超える注文を受け付けてしまう、いわゆる売り越しが起こります。

この課題を踏まえ、システムに求める性質を整理します。

  • 整合性 — 在庫を超える引き当てを起こさない。売り越しは金銭的な補償や信用の低下に直結するため、最優先で守る対象とする
  • 低レイテンシ — 注文者には短時間で確定または品切れを返す。待たせるほど離脱が増える
  • スケール — 商品点数と注文数の増加に対して、サイロを追加することで線形に近い形で対応できる
  • 可観測性 — 在庫の推移や注文の成否をリアルタイムに把握し、異常を早期に検知できる

従来のアーキテクチャでは、これらを両立させるために、データベースの行ロックや分散ロックへ頼ることが一般的でした。しかしロックは競合が激しいほど待ち行列を生み、レイテンシとスループットの双方を圧迫します。ここに Orleans のバーチャルアクターモデルがうまく噛み合います。グレインは同一キーにつき一度に一つのリクエストしか処理しないため、在庫という単位ごとに整合性の境界を自然に引けるからです。

アーキテクチャの全体像

システムは大きく三つの層で構成します。前段に注文を受け付ける ASP.NET Core の API 層、中核に Orleans のサイロクラスター、背後に永続化ストアとメッセージングを置きます。API 層はグレインへの入口に徹し、業務ロジックはグレインへ寄せます。

グレインの割り当て方は、そのままシステムの整合性設計になります。ここでは在庫を SKU(商品識別子)ごとの IInventoryGrain に対応させ、注文を注文 ID ごとの IOrderGrain に対応させます。SKU をキーにすることで、ある商品の在庫に関する判断はすべて単一のグレインへ集約され、そのグレインが一度に一つの引き当てしか処理しないという性質が、ロックなしの整合性を生みます。注文 ID をキーにすることで、同じ注文を二重に送っても同じグレインへ届き、冪等な処理として扱えます。

サイロの起動構成は次のようになります。クラスタリング、グレインストレージ、ストリーミングを一つのホストにまとめて設定します。

var builder = WebApplication.CreateBuilder(args);

builder.Host.UseOrleans(silo =>
{
    // クラスターメンバーシップ(サイロ同士の発見と障害検知)
    silo.UseAzureStorageClustering(options =>
        options.TableServiceClient = new TableServiceClient(clusteringConnection));

    // グレインの状態を永続化するストア
    silo.AddAzureTableGrainStorage("inventoryStore", options =>
        options.TableServiceClient = new TableServiceClient(storageConnection));
    silo.AddAzureTableGrainStorage("orderStore", options =>
        options.TableServiceClient = new TableServiceClient(storageConnection));

    // 読み取りモデルへ注文イベントを流すストリーム
    silo.AddMemoryGrainStorage("PubSubStore");
    silo.AddMemoryStreams("orders");
});

var app = builder.Build();

API 層は注文リクエストを受け取り、注文 ID から IOrderGrain を引いて処理を委譲します。結果に応じて HTTP ステータスを返し分けるだけで、同時実行の制御そのものはグレイン側が引き受けます。

app.MapPost("/orders", async (OrderRequest request, IGrainFactory grains) =>
{
    var order = grains.GetGrain<IOrderGrain>(request.OrderId);
    var result = await order.PlaceAsync(request);

    return result.Status switch
    {
        OrderStatus.Confirmed => Results.Ok(result),
        OrderStatus.OutOfStock => Results.Conflict(result),
        _ => Results.UnprocessableEntity(result)
    };
});

app.Run();

この段階で、クライアントから見た入口は単純な HTTP エンドポイントに収まっています。負荷が増えたときにサイロを足すという運用は、この API のコードには影響しません。

主要グレインの実装

まず在庫グレインです。状態は現在の手持ち数と、注文 ID ごとの引き当て(予約)を保持します。Orleans のシリアライザーに載せるため、状態やイベントの型には [GenerateSerializer] と [Id] を付けます。

public interface IInventoryGrain : IGrainWithStringKey
{
    Task<ReservationResult> ReserveAsync(string orderId, int quantity);
    Task ConfirmAsync(string orderId);
    Task ReleaseAsync(string orderId);
    Task<int> GetAvailableAsync();
}

[GenerateSerializer]
public sealed class InventoryState
{
    [Id(0)] public int OnHand { get; set; }
    [Id(1)] public Dictionary<string, int> Reservations { get; set; } = new();
}

public sealed class InventoryGrain : Grain, IInventoryGrain
{
    private readonly IPersistentState<InventoryState> _state;

    public InventoryGrain(
        [PersistentState("inventory", "inventoryStore")] IPersistentState<InventoryState> state)
        => _state = state;

    public async Task<ReservationResult> ReserveAsync(string orderId, int quantity)
    {
        // 同じ注文からの再送は同じ結果を返す(冪等性)
        if (_state.State.Reservations.ContainsKey(orderId))
            return ReservationResult.Reserved;

        var reserved = _state.State.Reservations.Values.Sum();
        var available = _state.State.OnHand - reserved;
        if (quantity > available)
            return ReservationResult.OutOfStock;

        _state.State.Reservations[orderId] = quantity;
        await _state.WriteStateAsync();
        return ReservationResult.Reserved;
    }

    public async Task ConfirmAsync(string orderId)
    {
        if (_state.State.Reservations.Remove(orderId, out var quantity))
        {
            _state.State.OnHand -= quantity;
            await _state.WriteStateAsync();
        }
    }

    public async Task ReleaseAsync(string orderId)
    {
        if (_state.State.Reservations.Remove(orderId))
            await _state.WriteStateAsync();
    }

    public Task<int> GetAvailableAsync()
        => Task.FromResult(_state.State.OnHand - _state.State.Reservations.Values.Sum());
}

ここで注目したいのは、在庫数を読んで判定し書き戻すという一連の操作に、ロックが一切ないことです。グレインは同一キーに対して一度に一つのリクエストしか処理しないため、この確認と更新は他のリクエストに割り込まれません。分散ロックや行ロックで守っていた不変条件を、グレインの実行モデルそのものが保証してくれます。引き当てを状態として永続化しているのも意図的で、サイロが落ちても別のサイロでグレインが復元され、処理中だった予約が失われません。

次に注文グレインです。在庫の予約、決済などの外部処理、確定という流れを一つのグレインが調整します。途中で失敗した場合は予約を解放し、整合性を保ちます。

public interface IOrderGrain : IGrainWithStringKey
{
    Task<OrderResult> PlaceAsync(OrderRequest request);
}

[GenerateSerializer]
public sealed class OrderState
{
    [Id(0)] public string Sku { get; set; } = "";
    [Id(1)] public int Quantity { get; set; }
    [Id(2)] public OrderStatus Status { get; set; }
}

public sealed class OrderGrain : Grain, IOrderGrain
{
    private readonly IPersistentState<OrderState> _state;
    private IAsyncStream<OrderEvent> _events = default!;

    public OrderGrain(
        [PersistentState("order", "orderStore")] IPersistentState<OrderState> state)
        => _state = state;

    public override Task OnActivateAsync(CancellationToken cancellationToken)
    {
        var streams = this.GetStreamProvider("orders");
        _events = streams.GetStream<OrderEvent>(StreamId.Create("order-events", "all"));
        return Task.CompletedTask;
    }

    public async Task<OrderResult> PlaceAsync(OrderRequest request)
    {
        // 二重送信への防御。確定済みなら同じ結果を返す
        if (_state.State.Status == OrderStatus.Confirmed)
            return new OrderResult(request.OrderId, OrderStatus.Confirmed);

        var inventory = GrainFactory.GetGrain<IInventoryGrain>(request.Sku);
        var reservation = await inventory.ReserveAsync(request.OrderId, request.Quantity);
        if (reservation != ReservationResult.Reserved)
            return new OrderResult(request.OrderId, OrderStatus.OutOfStock);

        _state.State.Sku = request.Sku;
        _state.State.Quantity = request.Quantity;
        _state.State.Status = OrderStatus.Reserved;
        await _state.WriteStateAsync();

        try
        {
            await ChargePaymentAsync(request);          // 外部決済など
            await inventory.ConfirmAsync(request.OrderId);
            _state.State.Status = OrderStatus.Confirmed;
            await _state.WriteStateAsync();
        }
        catch
        {
            // 失敗したら予約を解放し、在庫を取り戻す
            await inventory.ReleaseAsync(request.OrderId);
            _state.State.Status = OrderStatus.Failed;
            await _state.WriteStateAsync();
            return new OrderResult(request.OrderId, OrderStatus.Failed);
        }

        await _events.OnNextAsync(new OrderEvent(request.OrderId, request.Sku, request.Quantity));
        return new OrderResult(request.OrderId, OrderStatus.Confirmed);
    }
}

予約してから確定するという二段構えは、決済のような失敗しうる外部処理を挟むための工夫です。予約の時点では在庫を確保するだけで実際には減らさず、確定で初めて手持ち数を差し引きます。決済に失敗すれば予約を解放するので、在庫が宙に浮くことはありません。これは分散トランザクションを避け、補償によって整合性を保つ考え方をグレイン単位で実装したものです。

確定した注文はストリームへ発行し、読み取りモデルの構築を書き込み経路から切り離します。たとえば商品ごとの販売数を集計するダッシュボードは、このストリームを購読するだけで済みます。

public sealed class SalesDashboardGrain : Grain, ISalesDashboardGrain
{
    private readonly Dictionary<string, int> _soldBySku = new();

    public override async Task OnActivateAsync(CancellationToken cancellationToken)
    {
        var stream = this.GetStreamProvider("orders")
            .GetStream<OrderEvent>(StreamId.Create("order-events", "all"));

        await stream.SubscribeAsync((evt, token) =>
        {
            _soldBySku[evt.Sku] = _soldBySku.GetValueOrDefault(evt.Sku) + evt.Quantity;
            return Task.CompletedTask;
        });
    }

    public Task<int> GetSoldAsync(string sku)
        => Task.FromResult(_soldBySku.GetValueOrDefault(sku));
}

集計処理が注文の確定を遅らせることはありません。読み取り側の負荷や障害が書き込み側へ波及しない構造は、リアルタイム系のシステムでは安定運用の土台になります。

スケールと運用

スケールの基本は、グレインがサイロ群へ分散配置されることにあります。SKU をキーにした在庫グレインは、商品点数が多いほど自然に多数のサイロへ散らばり、負荷が広がります。サイロを増やせばクラスターは新しいメンバーを取り込み、活性化するグレインの配置が広がっていきます。アプリケーション側のコードを変えずに水平方向へ伸ばせるのが、この構造の利点です。

一方で、注意すべき偏りもあります。ある一つの SKU に注文が集中する場合、その在庫グレインは一つの活性化として一つのサイロに存在し、リクエストを直列に捌きます。これはグレインの単一スレッド実行の裏返しで、キー単位のスループットには上限があります。話題の商品が一点だけ突出して売れるような場面では、この上限が効いてきます。

対策の一つは、ホットな SKU に限って在庫を複数のシャードへ分割する方法です。総在庫をいくつかの子グレインへ配り、注文をそれらへ振り分けます。各シャードは別々のサイロに載りうるため、負荷が分散します。

public sealed class ShardedInventoryGrain : Grain, IShardedInventoryGrain
{
    private const int ShardCount = 16;

    public async Task<ReservationResult> ReserveAsync(string orderId, int quantity)
    {
        var sku = this.GetPrimaryKeyString();
        var start = (int)(Hash(orderId) % ShardCount);

        // ハッシュで始点を散らし、在庫が見つかるまでシャードを順に探る
        for (var i = 0; i < ShardCount; i++)
        {
            var index = (start + i) % ShardCount;
            var shard = GrainFactory.GetGrain<IInventoryGrain>($"{sku}:{index}");
            if (await shard.ReserveAsync(orderId, quantity) == ReservationResult.Reserved)
                return ReservationResult.Reserved;
        }
        return ReservationResult.OutOfStock;
    }
}

この手法には代償があります。一つのリクエストが在庫を見つけるまで複数のシャードを問い合わせることがあり、在庫全体の残数を知るにはシャードへのファンアウトが要ります。したがって、あらかじめホットになると分かっている一部の SKU にのみ適用し、大多数の商品は単一グレインのままにするのが現実的です。整合性を犠牲にせず、必要な箇所だけスループットを稼ぐという判断です。

障害からの回復は、状態の永続化とクラスターメンバーシップが支えます。サイロが失われても、そこにいたグレインは呼び出しを受けた時点で別のサイロに再活性化し、ストアから状態を読み戻します。引き当てを永続化しているため、処理中だった予約も引き継がれます。デプロイの際は、サイロを段階的に入れ替えるローリング更新とし、離脱するサイロのグレインが穏やかに移るようにします。

運用では可観測性を最初から組み込みます。第5回で扱ったように、グレイン呼び出しのレイテンシや活性化数、ストリームの滞留を OpenTelemetry のメトリクスとして継続的に収集します。加えて、このシステム固有の指標として、引き当ての拒否率と品切れの発生状況を見ておくと、在庫の枯渇やホットスポットの兆候を早期に捉えられます。ストリームを本番で使う場合は、メモリではなく Azure Event Hubs や各種キューを裏に持つ永続的なプロバイダーへ差し替え、再起動を挟んでも読み取りモデルを再構築できるようにします。

まとめ

本シリーズの締めくくりとして、注文と在庫のリアルタイム処理を題材に、前回までの要素を一つのシステムへ組み上げました。設計の中心にあったのは、在庫という不変条件を単一のグレインの内側へ置き、グレインの単一スレッド実行によってロックなしで整合性を得るという判断です。そのうえで、状態の永続化が障害復旧を、ストリーミングが読み取りモデルの分離を、クラスタリングが水平スケールを担い、ホットスポットにはシャーディングという代償つきの手段で対処しました。

強調しておきたいのは、万能の解はないという点です。グレインの単一スレッド実行はキー単位のスループットを制約し、キー設計そのものがシステムの整合性とスケールの両方を決めます。逆に言えば、ドメインの不変条件をどのキーへ寄せるかを見極められれば、Orleans は分散システムの難所を大きく肩代わりしてくれます。第1回から追ってきたバーチャルアクターモデルの考え方は、現実のシステム設計においても一貫して有効でした。全6回にわたりお付き合いいただき、ありがとうございました。

エンハンスド株式会社では、Orleans を用いた分散システムの設計・開発をご支援しています。リアルタイム性と整合性、そして本番運用に耐えるスケーラビリティを同時に求められるサービスについて、アーキテクチャの妥当性検証から、グレイン設計、既存システムのモダナイゼーション、運用設計までを、実務に即した形でご一緒します。Orleans の採用を具体的に検討されている方は、お気軽にご相談ください。

この記事をシェア

コピーしました

関連記事