【実務・中級編】VB.NETでのTPL Dataflowを用いたパイプライン処理の構築:非同期メッセージングによる疎結合なデータ処理フロー – Visual Basic (VB / VB.NET)解析バイブル

スポンサーリンク

VB.NETでのTPL Dataflowを用いたパイプライン処理の構築:非同期メッセージングによる疎結合なデータ処理フロー

業務システムやデータ処理ツールを開発していると、必ず直面する壁がある。
「数万件のCSVファイルを読み込み、バリデーションを行い、外部APIを叩いてDBに書き込む」――この一連の処理を、従来のシングルスレッドや単純な `For Each` で実装していないだろうか?

処理が重くなるとUIがフリーズし、例外処理でどこが壊れたか分からなくなったり、メモリが枯渇したりする。かといって、生粋の `Task` を直接ハンドリングして並列制御を書こうとすれば、スレッドセーフティの担保や排他制御(Lock)の嵐で、コードは保守不可能な「スパゲッティ」へと成り下がる。

この地獄から抜け出すための特効薬が、.NETが誇る最強の非同期メッセージング基盤「TPL Dataflow(Task Parallel Library Dataflow)」だ。

今回は、VB.NETを用いて、データ加工・変換・出力の各工程を独立した「ブロック」に分割し、バッファリングと並列度を完璧に制御しながら流し込むモダンなパイプラインアーキテクチャを伝授する。

なぜ従来の書き方は非効率なのか?

多くの開発者は、データを処理する際に以下のようなアプローチをとる。

1. すべてのデータをメモリに読み込む
2. ループを回して1件ずつ処理・DB保存を行う

このアプローチの致命的な欠点は「結合度(Coupling)の高さ」と「スロットリング(流量制御)の欠如」にある。
「読む」「変換する」「書く」の処理が密結合しているため、データベースの書き込みがボトルネックになると、読み込み側や変換側の処理まで引きずられて全体がスローダウンする。また、数百万件のデータを一気にメモリに載せれば、容赦なく `OutOfMemoryException` が発生する。

TPL Dataflowによる解決

TPL Dataflowでは、処理を「ブロック(Block)」という独立した部品に分割し、それらを「パイプ(LinkTo)」で繋ぐ。
各ブロックは独自のバッファ(キュー)を持ち、自身のリソース状況に応じてデータを非同期で引き受けて処理する。つまり、「プロデューサー・コンシューマー問題」を自前で実装することなく、宣言的に安全な並列パイプラインを構築できるのだ。

アーキテクチャの全体像

今回構築するパイプラインは以下の3段構えとする。

[CSVファイル読み込み (Producer)]
↓ (BufferBlock)
[データ変換・バリデーション (TransformBlock)]
↓ (BoundedCapacity付きBuffer)
[DB一括書き込み (ActionBlock – 並列度制御)]

  • BufferBlock: 読み込んだ生データを一旦保持するバッファ。
  • TransformBlock: 生データをドメインモデルに変換し、バリデーションを行う。
  • ActionBlock: 変換済みデータをデータベースへ非同期で書き込む(ここで並列度 `MaxDegreeOfParallelism` を絞ることで、DBのコネクション枯渇を防ぐ)。

プロダクションコード例:堅牢なパイプライン実装

実務でそのまま使える、堅牢なVB.NETのコードを示す。
必要十分なエラーハンドリング、キャンセル処理(`CancellationToken`)、そしてブロック間の適切なバッファ制限(BoundedCapacity)を網羅している。

Imports System.IO
Imports System.Threading
Imports System.Threading.Tasks
Imports System.Threading.Tasks.Dataflow

Namespace DataflowPipeline

‘ 処理対象のデータモデル
Public Class SalesRecord
Public Property OrderId As String
Public Property Amount As Decimal
Public Property CustomerName As String
Public Property IsValid As Boolean = True
Public Property ErrorMessage As String = String.Empty
End Class

Public Class PipelineRunner

Public Async Function ExecutePipelineAsync(csvFilePath As String, cancellationToken As CancellationToken) As Task

‘ —————————————————————–
‘ 1. ブロックの設定(実行オプション)
‘ —————————————————————–
‘ メモリ枯渇を防ぐため、各ブロックのバッファ上限(BoundedCapacity)を必ず設定する
Dim executionOptions = New ExecutionDataflowBlockOptions With {
.MaxDegreeOfParallelism = Environment.ProcessorCount, ‘ CPUコア数に応じた並列度
.CancellationToken = cancellationToken,
.BoundedCapacity = 1000 ‘ バッファに溜める最大件数を制限し、バックプレッシャーを効かせます
}

‘ —————————————————————–
‘ 2. 各パイプラインブロックの定義
‘ —————————————————————–

‘ 【Step 1: 変換・バリデーションブロック】
‘ CSVの1行(String)を受け取り、SalesRecordオブジェクトに変換する
Dim transformBlock = New TransformBlock(Of String, SalesRecord)(
Function(line)
Dim cols = line.Split(“,”c)
Dim record As New SalesRecord()

Try
‘ 簡易パースとバリデーション
record.OrderId = cols(0).Trim()
record.Amount = Decimal.Parse(cols(1).Trim())
record.CustomerName = cols(2).Trim()

If record.Amount < 0 Then record.IsValid = False record.ErrorMessage = "金額が負の値です。" End If Catch ex As Exception record.IsValid = False record.ErrorMessage = $"パースエラー: {ex.Message}" End Try Return record end Function, executionOptions ) ' 【Step 2: 出力(DB書き込み)ブロック】 ' 変換されたデータを受け取り、外部ストレージやDBへ非同期で書き込む ' DBの過負荷を防ぐため、ここだけ並列度を「2」に絞るなどのチューニングが可能 Dim dbWriteOptions = New ExecutionDataflowBlockOptions With { .MaxDegreeOfParallelism = 2, ' DBのコネクションプールを保護 .CancellationToken = cancellationToken, .BoundedCapacity = 500 } Dim dbActionBlock = New ActionBlock(Of SalesRecord)( Async Function(record) If Not record.IsValid Then Console.WriteLine($"[スキップ] OrderId: {record.OrderId}, 理由: {record.ErrorMessage}") Return End Function ' 非同期DB書き込みのシミュレーション(実際はEntity FrameworkやDapper等を使用) Await Task.Delay(100, cancellationToken) Console.WriteLine($"[DB保存完了] OrderId: {record.OrderId}, 金額: {record.Amount}") End Function, dbWriteOptions ) ' ----------------------------------------------------------------- ' 3. ブロック同士の結合(LinkTo) ' ----------------------------------------------------------------- ' 伝播オプション:データ転送時にソースを自動完了させ、例外を伝播させる Dim linkOptions = New DataflowLinkOptions With {.PropagateCompletion = True} ' transformBlock の出力結果を dbActionBlock の入力へ繋ぐ transformBlock.LinkTo(dbActionBlock, linkOptions) ' ----------------------------------------------------------------- ' 4. データの流し込み(プロデューサー) ' ----------------------------------------------------------------- Try Using reader As New StreamReader(csvFilePath) ' ヘッダーをスキップ Await reader.ReadLineAsync() While Not reader.EndOfStream Dim line = Await reader.ReadLineAsync() If Not String.IsNullOrWhiteSpace(line) then ' transformBlockへデータを非同期でプッシュ(バッファがいっぱいの場合はここで待機する) Await transformBlock.SendAsync(line, cancellationToken) End If End While End Using ' データの供給源が尽きたことを通知 transformBlock.Complete() Catch ex As Exception Console.WriteLine($"パイプライン異常終了: {ex.Message}") CType(transformBlock, IDataflowBlock).Fault(ex) End Try ' ----------------------------------------------------------------- ' 5. パイプライン全体の完了待ち ' ----------------------------------------------------------------- ' 最後のブロック(dbActionBlock)の処理がすべて終わるまで非同期で待機 Await dbActionBlock.Completion Console.WriteLine("すべてのデータ処理パイプラインが正常に完了しました。") End Function End Class End Namespace ---

現場で絶対に外せない設計の勘所

プロのアーキテクトとして、このコードを運用する上で死守すべき鉄則を授ける。

1. `BoundedCapacity`(バッファ上限)の強制

TPL Dataflowで最も多いバグは、`BoundedCapacity` を指定せずに無限にデータを流し込み、メモリが爆発することだ。
プロデューサー(読み込み側)がいくら速くても、コンシューマー(DB書き込み)が遅ければ、バッファが無限に膨れ上がる。`BoundedCapacity` を設定することで、コンシューマーの処理速度に合わせてプロデューサー側が自動的にスロットリング(一時停止)する「バックプレッシャー(背圧)」の仕組みが機能する。これがプロの設計だ。

2. 例外伝播と `PropagateCompletion`

パイプラインの途中で例外が発生した場合、何もしないと処理が途中でフリーズしたままになりかねない。
`LinkTo` のオプションに `.PropagateCompletion = True` を指定し、さらにソース側でエラーが発生した場合は `IDataflowBlock.Fault(ex)` を明示的に呼び出すことで、パイプライン全体に異常を伝播させ、デッドロックを防ぐことができる。

3. スレッドセーフティと外部リソース

各ブロックのラムダ式内で外部の可変オブジェクト(コレクションなど)にアクセスする場合、マルチスレッド環境下での競合に注意しなければならない。原則として、ブロック内では状態を持たない(ステートレスな)コードを心がけ、集約処理が必要な場合は `BatchBlock` や専用の集約用ブロックを挟むこと。

まとめ

VB.NETにおけるTPL Dataflowの導入は、単なる「処理の高速化」にとどまらない。
「データの読み込み」「変換」「永続化」という関心事を完全に分離(疎結合化)し、スレッド管理の複雑さをフレームワークに丸投げすることで、「保守性が高く、破綻しない堅牢なバッチ・ETLアーキテクチャ」を手に入れることができる。

「なぜこの書き方は非効率なのか」を理解したあなたなら、もう愚直なシングルスレッドのループに戻ることはできないはずだ。
次の業務システムのバッチ処理やファイル連携モジュールから、ぜひこのモダンな非同期パイプラインパターンを適用してほしい。圧倒的なパフォーマンスとコードの美しさに、チームメンバーも驚くことだろう。

タイトルとURLをコピーしました