Dagsterで会話の分岐と評価データを管理する
Go・NATS・Dagsterで会話の再実行、評価、学習データ生成を設計。AIシステム開発のデータ基盤と、レビューで見つかった4つの問題を記録します。
追加した会話データ基盤
稼働中の Go + NATS + Dagster イベントパイプラインに、会話の分岐、評価、学習データ生成を扱う3つのアセットグループを設計・実装しました。LLM 会話を、任意の地点から再実行できるグラフとして記録するためです。
PostgreSQL は再実行に必要なリネージュ、DuckDB は評価・クラスタリング、Iceberg / Parquet は凍結済みデータセットを保存します。用途を分け、各層を独立に変更できる構成にしました。
実装レビューでは、グルーピング、メタデータ、turn_index、未実装の asset に4つの問題が見つかりました。修正方針と残課題も記録します。
前提
既存パイプラインの状態
agent-gateway は Go + Gin の OpenAI 互換プロキシです。NATS JetStream、Dagster、PostgreSQL + pgvector、DuckDB を使い、イベント処理、検索、分析を行います。
既存の pipeline_runtime では11個の Dagster アセットが稼働していました。NATS センサーが拾ったイベントを pipeline_event_inbox から llm_chat_turns、flow_nodes、flow_edges へ書き込みます。
不足していた機能
既存処理はリクエスト単位の正規化と蓄積までで、次の機能がありませんでした。
- 会話全体を 分岐可能なグラフ として扱う仕組み
- トラジェクトリの評価
- トレーニング用データセットの自動生成
ローカル OSS LLM の fine-tuning、モデル評価、system prompt の最適化に使うデータ基盤として、これらを Dagster に追加します。
3つのアセットグループの設計
Conversation Lineage Assets
会話ツリーを記録するグループです。
| アセット | 役割 |
|---|---|
seed_themes | 実験テーマの参照テーブル |
conversation_sessions | セッションの管理 |
conversation_branch_nodes | 会話ツリーの全ノードを保持 |
taskification_results | MCP やタスク化の中間結果 |
execution_run_records | 実行ファクト |
response_branch_records | 最終レスポンスとトレース |
各ノードは branch_kind(root・retry・tool_variant・sampling_variant・prompt_variant・taskify_variant)と variant_config を JSON で保存します。
依存グラフは次の直列構造です。
seed_themes
→ conversation_sessions
→ conversation_branch_nodes
→ taskification_results
→ execution_run_records
→ response_branch_records
リプレイルール
taskify_rule_versionが変わった場合はtaskification_results以降を再生成model_snapshot_idが変わった場合はexecution_run_records以降を再生成- 特定ノードの条件変更はそのサブツリーのみ再生成
Evaluation and Distribution Labeling Assets
会話の処理経路を評価単位へ変換するグループです。
trajectory_units # ルートからターミナルノードまでのパスを正規化
→ trajectory_embeddings # ベクトル化
→ trajectory_clusters # クラスタリング
→ cluster_category_labels # カテゴリラベル付与
→ human_sample_queue # 人間レビューキュー
→ human_eval_labels # 評価結果の蓄積
→ candidate_preference_pairs
→ preference_labels # DPO ペア構築
カテゴリは tool_helped、tool_misuse、hallucination、good_recovery、overlong、underspecified、stable_high_quality などです。
設計原則
- embedding はラベルではなく表現
- クラスタは判断ではなく構造
- モデルによるカテゴリ化と人間の選好は別レイヤーとして明確に分離する
Training Dataset Generation Assets
評価済みの処理経路から学習データを作るグループです。
positive_topk_trajectories # 層別 top-K 抽出
negative_bottomk_trajectories # 層別 bottom-K 抽出
→ sft_messages_dataset # SFT 用データセット
→ dpo_preference_dataset # DPO 用データセット
→ suppression_dataset # 抑制用データセット
→ unlearning_dataset # アンラーニング用データセット
→ dataset_snapshot_registry # 凍結済みスナップショットのメタデータ管理
negative type は hard_negative・tool_misuse・hallucination・noise_reject・suppression_candidate です。すべてを一律にアンラーニングへ使わず、種類ごとに分けます。
設計上の主要な意思決定
能動的オーケストレーターか受動的キャプチャか
seed_themes → conversation_sessions → branch_nodes が会話を駆動するのか、既存の会話を整理するのかを検討しました。
まず手動でテーマを決め、セッション ID を発行して対話を繰り返します。その ID で会話をマテリアライズし、任意のノードから再実行できる土台を作ります。会話を自動で進めるオーケストレーターは、その上に後から載せる設計です。
分岐処理の位置
既存の session / session_turn / session_lineage は実行ログとして残し、その上に実験用の asset を作ります。
親子関係を持つターンに、次の情報を追加します。
branch_kindvariant_confignode_statuscontext_snapshot_idtaskification replay point
taskification_results は FlowLineageEvent と分け、prompt 分解と orchestration plan の中間表現として保存します。
ストレージ三層構造
| レイヤー | ストレージ | 用途 |
|---|---|---|
| Truth source | PostgreSQL | replay 必須のリネージュデータ |
| 分析ワークベンチ | DuckDB | 評価・クラスタリング向け派生データ |
| Immutable artifact | Iceberg / Parquet | 凍結済みデータセットスナップショット |
リネージュの保存方式、評価処理、データセット形式をそれぞれ変更できるように分けています。
実装優先順位
実装は次の5段階で進めます。
- replay 可能なリネージュモデルの確立
- taskification と branch generation の導入
- trajectory 化
- embedding / clustering / labeling
- dataset generation と snapshot freeze
最初に OSS ローカル LLM の replay と branch node の正規化を固めます。その結果を評価とデータセット生成で使います。
実装計画
ディレクトリ構造
model-foundry/src/loftllc/ に pipeline_runtime と並列のモジュールを追加します。既存の domain / infra / presentation の構成に合わせます。
loftllc/
domain/
conversation_lineage.py
infra/
lineage_store.py
presentation/
dagster/
assets/
conversation_lineage.py
アセット実行パターンの違い
pipeline_runtime は NATS メッセージごとに1つの RunRequest を作ります。リネージュアセットは蓄積済みのデータをまとめて処理します。
- NATS イベントはすでに
pipeline_runtimeで処理されllm_chat_turnsやflow_nodesに蓄積済み - リネージュアセットはそこから読み取る 下流コンシューマ として動作
- セッションの蓄積データをバッチ的にマテリアライズする方式
マテリアライズの単位は次のとおりです。
| グループ | トリガー |
|---|---|
| Group 1(Conversation Lineage) | 手動または 60 秒間隔センサー |
| Group 2(Evaluation) | Group 1 への依存による自動トリガー |
| Group 3(Dataset Generation) | 完全手動トリガー |
パーティション戦略
DynamicPartitionsDefinition を使います。セッション ID は Go ゲートウェイが実行時に作るため、センサーが llm_chat_turns から発見し、動的パーティションへ追加して RunRequest を発行します。Group 3 はスナップショット ID を別の動的パーティションにします。
リプレイメカニズム
ノード X からの replay は次の処理です。
- ノード X の会話状態(凍結済みメッセージ履歴)を取得
- 1 つ以上の LLM にそのメッセージ履歴を送信
- 新しいレスポンスを
response_branch_recordsとして記録
infra リソース LLMClientResource から /v1/chat/completions へ HTTP リクエストを送ります。replay もゲートウェイを通り、通常と同じルーティング、ログ、NATS イベントを記録します。
PostgreSQL スキーマ設計
リネージュ用テーブルの定義です。
| テーブル | 役割 |
|---|---|
seed_themes | 実験テーマの参照テーブル |
conversation_sessions | llm_chat_turns から集約されたセッションビュー |
conversation_branch_nodes | 会話内の分岐点(branch_node_id・parent_node_id・turn_index・message_history_json・branch_reason) |
taskification_results | 分解されたサブプロンプト |
execution_run_records | LLM 実行ファクト(original または replay) |
response_branch_records | 同一プロンプトに対する比較レコード(context_anchor_turn_id で分岐元を追跡) |
lineage_watermark | 増分処理用のウォーターマーク |
実装レビューで発見した設計問題
Phase 1 の会話リネージュ基盤を実装後、レビューで4つの問題が見つかりました。
問題1: response_branch_records のグルーピング破綻
replay ごとに新しい turn_id を付け、その ID でグループ分けしていました。このため各 replay が別グループになり、branch_index は実質0のままでした。同じ親から条件を変えた分岐を比較できず、応答の比較や preference pair の元データになりません。
修正方針:context_anchor_turn_id でグループ分けします。
問題2: original turn の execution_run_records メタデータ欠落
chat_history から取得するのは role / content / created_at だけでしたが、ビルダーは model / backend / assistant_message / token_usage / duration_ms も読む前提でした。元の実行の response_text と usage が空になります。
メタデータは llm_chat_turns(DuckDB)または pipeline_event_inbox(PostgreSQL)から取得する必要がありました。
問題3: replay branch の turn_index 衝突
replay branch に固定で turn_index = parent_turn_index + 1 を入れ、既存の次ターンと衝突していました。読み出しも ORDER BY turn_index ASC だけなので、兄弟分岐と直列ターンの順序が不安定です。任意ノードを replay root にする前提が崩れます。
修正方針:turn_index による直列モデルから、branch node table の親子ツリーへ変えます。
問題4: 仕様の中核 asset の未実装
conversation_sessions と conversation_branch_nodes のマテリアライズがなく、depth / branch_kind / context_snapshot_id / variant_config / node_status が保存されていませんでした。分岐グラフに必要な情報が不足していました。
修正後の残課題
1回目の修正後も、次の2点が残りました。
- 通常マテリアライズで replay node を除外するか、
build_execution_record_from_event()側で original source のない turn を skip するかの選択 conversation_branch_nodesの親参照 2-pass 構築が必要
追加の修正方針です。
upsert_branch_node()の conflict update に topology 系カラム(parent_node_id・depth・branch_kind)も含める- original run でも
status / error / tool_traceをpipeline_eventから埋める
再実行から学習データへつなぐ
既存の session / session_turn / session_lineage の上に、実験用の5つのテーブルを作ります。任意ノードを context anchor にして条件を制御した replay を行い、応答分布を評価して学習用 dataset へつなぐ設計です。
結果
- Conversation Lineage・Evaluation & Distribution Labeling・Training Dataset Generation の3アセットグループを設計し、合計約30個のアセットを追加
- PostgreSQL / DuckDB / Iceberg の三層ストレージ戦略を確定
- Phase 1(会話リネージュ基盤)の実装を完了
- コードレビューで4つの重大な設計問題を発見・修正
- 残課題(branch node の 2-pass 構築・context_anchor ベースのグルーピング・replay node の除外ロジック)を明確化
- 実装優先順位を「replay 基盤 → taskification → trajectory 化 → evaluation → dataset generation」の5段階に確定
その後の方針転換
運用中に、次の方針へ変更しました。
テレメトリの追加と可視化パイプラインへの転換
Dagster のアセットグラフだけでは、会話の状態をリアルタイムに把握しにくい問題がありました。Vector → Prometheus → Grafana のテレメトリ経路を追加し、会話の状態、分岐、評価結果をダッシュボードで表示するようにしました。
Dagster アセットのフラット化
可視化をテレメトリへ移し、Dagster のアセットはフラットにする方針へ変えました。会話ツリーを深い依存グラフで表す必要がなくなったためです。
agent-gateway のオーケストレーターモデルへの再設計
agent-gateway も、単一ゲートウェイで全リクエストを処理する構成から、オーケストレーターモデルがワーカー LLM 群を操作する構成へ変えました。LLM 同士の対話フローをアセットにし、Grafana でリアルタイムに表示します。
この再設計は別の記事で記録する予定です。
