Job is idle — throughput ~0; structure shown.
Fragment 54814 (Actor 740417,740418)
StreamMaterialize { columns: [benchmark_id, fact_date, asset_id, daily_subperiod_return], stream_key: [benchmark_id, asset_id, fact_date], pk_columns: [benchmark_id, asset_id, fact_date], pk_conflict: NoCheck }
├── output: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, benchmark_constituents_ft.asset_id, $expr1 ]
├── stream key: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.date ]
└── StreamProject { exprs: [benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, benchmark_constituents_ft.asset_id, (asset_twrr_mv_next.daily_subperiod_return * benchmark_constituents_ft.weight) as $expr1] }
├── output: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, benchmark_constituents_ft.asset_id, $expr1 ]
├── stream key: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.date ]
└── MergeExecutor
├── output: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.weight, asset_twrr_mv_next.daily_subperiod_return, benchmarks_dm.id ]
└── stream key: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.date ]
Fragment 54815 (Actor 740419,740420)
StreamSyncLogStore
├── output: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.weight, asset_twrr_mv_next.daily_subperiod_return, benchmarks_dm.id ]
├── stream key: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.date ]
└── StreamHashJoin { type: Inner, predicate: benchmark_constituents_ft.benchmark_id = benchmarks_dm.id }
├── output: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.weight, asset_twrr_mv_next.daily_subperiod_return, benchmarks_dm.id ]
├── stream key: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.date ]
├── MergeExecutor
│ ├── output:
│ │ ┌── benchmark_constituents_ft.benchmark_id
│ │ ├── benchmark_constituents_ft.date
│ │ ├── benchmark_constituents_ft.asset_id
│ │ ├── benchmark_constituents_ft.weight
│ │ ├── asset_twrr_mv_next.daily_subperiod_return
│ │ ├── asset_twrr_mv_next.asset_id
│ │ └── asset_twrr_mv_next.fact_date
│ └── stream key: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.date ]
└── MergeExecutor { output: [ benchmarks_dm.id ], stream key: [ benchmarks_dm.id ] }
Fragment 54816 (Actor 740421,740422)
StreamLocalityProvider { locality_columns: [benchmark_constituents_ft.benchmark_id] }
├── output:
│ ┌── benchmark_constituents_ft.benchmark_id
│ ├── benchmark_constituents_ft.date
│ ├── benchmark_constituents_ft.asset_id
│ ├── benchmark_constituents_ft.weight
│ ├── asset_twrr_mv_next.daily_subperiod_return
│ ├── asset_twrr_mv_next.asset_id
│ └── asset_twrr_mv_next.fact_date
├── stream key: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.date ]
└── MergeExecutor
├── output:
│ ┌── benchmark_constituents_ft.benchmark_id
│ ├── benchmark_constituents_ft.date
│ ├── benchmark_constituents_ft.asset_id
│ ├── benchmark_constituents_ft.weight
│ ├── asset_twrr_mv_next.daily_subperiod_return
│ ├── asset_twrr_mv_next.asset_id
│ └── asset_twrr_mv_next.fact_date
└── stream key: [ benchmark_constituents_ft.asset_id, benchmark_constituents_ft.date, benchmark_constituents_ft.benchmark_id ]
Fragment 54817 (Actor 740424,740423)
StreamSyncLogStore
├── output:
│ ┌── benchmark_constituents_ft.benchmark_id
│ ├── benchmark_constituents_ft.date
│ ├── benchmark_constituents_ft.asset_id
│ ├── benchmark_constituents_ft.weight
│ ├── asset_twrr_mv_next.daily_subperiod_return
│ ├── asset_twrr_mv_next.asset_id
│ └── asset_twrr_mv_next.fact_date
├── stream key: [ benchmark_constituents_ft.asset_id, benchmark_constituents_ft.date, benchmark_constituents_ft.benchmark_id ]
└── StreamHashJoin { type: Inner, predicate: benchmark_constituents_ft.asset_id = asset_twrr_mv_next.asset_id AND benchmark_constituents_ft.date = asset_twrr_mv_next.fact_date }
├── output:
│ ┌── benchmark_constituents_ft.benchmark_id
│ ├── benchmark_constituents_ft.date
│ ├── benchmark_constituents_ft.asset_id
│ ├── benchmark_constituents_ft.weight
│ ├── asset_twrr_mv_next.daily_subperiod_return
│ ├── asset_twrr_mv_next.asset_id
│ └── asset_twrr_mv_next.fact_date
├── stream key: [ benchmark_constituents_ft.asset_id, benchmark_constituents_ft.date, benchmark_constituents_ft.benchmark_id ]
├── MergeExecutor
│ ├── output: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.weight ]
│ └── stream key: [ benchmark_constituents_ft.asset_id, benchmark_constituents_ft.date, benchmark_constituents_ft.benchmark_id ]
└── MergeExecutor { output: [ asset_twrr_mv_next.asset_id, asset_twrr_mv_next.fact_date, asset_twrr_mv_next.daily_subperiod_return ], stream key: [ asset_twrr_mv_next.asset_id, asset_twrr_mv_next.fact_date ] }
Fragment 54818 (Actor 740454,740453)
StreamLocalityProvider { locality_columns: [benchmark_constituents_ft.asset_id, benchmark_constituents_ft.date] }
├── output: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.weight ]
├── stream key: [ benchmark_constituents_ft.asset_id, benchmark_constituents_ft.date, benchmark_constituents_ft.benchmark_id ]
└── MergeExecutor
├── output: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.weight ]
└── stream key: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, benchmark_constituents_ft.asset_id ]
Fragment 54819 (Actor 740458,740457)
StreamTableScan { table: benchmark_constituents_ft, columns: [benchmark_id, date, asset_id, weight] }
├── output: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.weight ]
├── stream key: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, benchmark_constituents_ft.asset_id ]
├── Upstream { output: [ benchmark_id, date, asset_id, weight ], stream key: [] }
└── BatchPlanNode { output: [ benchmark_id, date, asset_id, weight ], stream key: [] }
Fragment 54820 (Actor 740440,740439)
StreamFilter { predicate: Not(IsNull(asset_twrr_mv_next.daily_subperiod_return)) }
├── output: [ asset_twrr_mv_next.asset_id, asset_twrr_mv_next.fact_date, asset_twrr_mv_next.daily_subperiod_return ]
├── stream key: [ asset_twrr_mv_next.asset_id, asset_twrr_mv_next.fact_date ]
└── StreamTableScan { table: asset_twrr_mv_next, columns: [asset_id, fact_date, daily_subperiod_return] }
├── output: [ asset_twrr_mv_next.asset_id, asset_twrr_mv_next.fact_date, asset_twrr_mv_next.daily_subperiod_return ]
├── stream key: [ asset_twrr_mv_next.asset_id, asset_twrr_mv_next.fact_date ]
├── Upstream { output: [ asset_id, fact_date, daily_subperiod_return ], stream key: [] }
└── BatchPlanNode { output: [ asset_id, fact_date, daily_subperiod_return ], stream key: [] }
Fragment 54821 (Actor 740460,740459)
StreamProject { exprs: [benchmarks_dm.id] } { output: [ benchmarks_dm.id ], stream key: [ benchmarks_dm.id ] }
└── StreamFilter { predicate: IsNull(benchmarks_dm.disabled_at) } { output: [ benchmarks_dm.id, benchmarks_dm.disabled_at ], stream key: [ benchmarks_dm.id ] }
└── StreamTableScan { table: benchmarks_dm, columns: [id, disabled_at] } { output: [ benchmarks_dm.id, benchmarks_dm.disabled_at ], stream key: [ benchmarks_dm.id ] }
├── Upstream { output: [ id, disabled_at ], stream key: [] }
└── BatchPlanNode { output: [ id, disabled_at ], stream key: [] }