Job is idle — throughput ~0; structure shown.
Fragment 25264 (Actor 122760,122759)
StreamMaterialize { columns: [asset_id, last_price, value_timestamp], stream_key: [asset_id], pk_columns: [asset_id], pk_conflict: NoCheck } { output: [ intraday_asset_prices_5min_mv.asset_id, $expr1, $expr2 ], stream key: [ intraday_asset_prices_5min_mv.asset_id ] }
└── StreamProject { exprs: [intraday_asset_prices_5min_mv.asset_id, Coalesce(first_value(intraday_asset_prices_5min_mv.last_price order_by(intraday_asset_prices_5min_mv.window_start DESC)) filter(IsNotNull(intraday_asset_prices_5min_mv.last_price) AND IsNotNull(intraday_asset_prices_5min_mv.window_start)), asset_prices_snapshot_mv_next.close) as $expr1, Case(Not(IsNull(first_value(intraday_asset_prices_5min_mv.last_price order_by(intraday_asset_prices_5min_mv.window_start DESC)) filter(IsNotNull(intraday_asset_prices_5min_mv.last_price) AND IsNotNull(intraday_asset_prices_5min_mv.window_start)))), max(intraday_asset_prices_5min_mv.value_timestamp), asset_prices_snapshot_mv_next.value_timestamp) as $expr2] }
├── output: [ intraday_asset_prices_5min_mv.asset_id, $expr1, $expr2 ]
├── stream key: [ intraday_asset_prices_5min_mv.asset_id ]
└── MergeExecutor { output: [ intraday_asset_prices_5min_mv.asset_id, first_value(intraday_asset_prices_5min_mv.last_price order_by(intraday_asset_prices_5min_mv.window_start DESC)) filter(IsNotNull(intraday_asset_prices_5min_mv.last_price) AND IsNotNull(intraday_asset_prices_5min_mv.window_start)), max(intraday_asset_prices_5min_mv.value_timestamp), asset_prices_snapshot_mv_next.close, asset_prices_snapshot_mv_next.value_timestamp, asset_prices_snapshot_mv_next.asset_id ], stream key: [ intraday_asset_prices_5min_mv.asset_id ] }
Fragment 25265 (Actor 122753,122754)
StreamSyncLogStore { output: [ intraday_asset_prices_5min_mv.asset_id, first_value(intraday_asset_prices_5min_mv.last_price order_by(intraday_asset_prices_5min_mv.window_start DESC)) filter(IsNotNull(intraday_asset_prices_5min_mv.last_price) AND IsNotNull(intraday_asset_prices_5min_mv.window_start)), max(intraday_asset_prices_5min_mv.value_timestamp), asset_prices_snapshot_mv_next.close, asset_prices_snapshot_mv_next.value_timestamp, asset_prices_snapshot_mv_next.asset_id ], stream key: [ intraday_asset_prices_5min_mv.asset_id ] }
└── StreamHashJoin { type: LeftOuter, predicate: intraday_asset_prices_5min_mv.asset_id = asset_prices_snapshot_mv_next.asset_id } { output: [ intraday_asset_prices_5min_mv.asset_id, first_value(intraday_asset_prices_5min_mv.last_price order_by(intraday_asset_prices_5min_mv.window_start DESC)) filter(IsNotNull(intraday_asset_prices_5min_mv.last_price) AND IsNotNull(intraday_asset_prices_5min_mv.window_start)), max(intraday_asset_prices_5min_mv.value_timestamp), asset_prices_snapshot_mv_next.close, asset_prices_snapshot_mv_next.value_timestamp, asset_prices_snapshot_mv_next.asset_id ], stream key: [ intraday_asset_prices_5min_mv.asset_id ] }
├── MergeExecutor { output: [ intraday_asset_prices_5min_mv.asset_id, first_value(intraday_asset_prices_5min_mv.last_price order_by(intraday_asset_prices_5min_mv.window_start DESC)) filter(IsNotNull(intraday_asset_prices_5min_mv.last_price) AND IsNotNull(intraday_asset_prices_5min_mv.window_start)), max(intraday_asset_prices_5min_mv.value_timestamp), intraday_asset_prices_5min_mv.asset_id ], stream key: [ intraday_asset_prices_5min_mv.asset_id ] }
└── MergeExecutor { output: [ asset_prices_snapshot_mv_next.asset_id, asset_prices_snapshot_mv_next.close, asset_prices_snapshot_mv_next.value_timestamp ], stream key: [ asset_prices_snapshot_mv_next.asset_id ] }
Fragment 25266 (Actor 122755,122756)
StreamSyncLogStore { output: [ intraday_asset_prices_5min_mv.asset_id, first_value(intraday_asset_prices_5min_mv.last_price order_by(intraday_asset_prices_5min_mv.window_start DESC)) filter(IsNotNull(intraday_asset_prices_5min_mv.last_price) AND IsNotNull(intraday_asset_prices_5min_mv.window_start)), max(intraday_asset_prices_5min_mv.value_timestamp), intraday_asset_prices_5min_mv.asset_id ], stream key: [ intraday_asset_prices_5min_mv.asset_id ] }
└── StreamHashJoin { type: LeftOuter, predicate: intraday_asset_prices_5min_mv.asset_id = intraday_asset_prices_5min_mv.asset_id } { output: [ intraday_asset_prices_5min_mv.asset_id, first_value(intraday_asset_prices_5min_mv.last_price order_by(intraday_asset_prices_5min_mv.window_start DESC)) filter(IsNotNull(intraday_asset_prices_5min_mv.last_price) AND IsNotNull(intraday_asset_prices_5min_mv.window_start)), max(intraday_asset_prices_5min_mv.value_timestamp), intraday_asset_prices_5min_mv.asset_id ], stream key: [ intraday_asset_prices_5min_mv.asset_id ] }
├── StreamProject { exprs: [intraday_asset_prices_5min_mv.asset_id] } { output: [ intraday_asset_prices_5min_mv.asset_id ], stream key: [ intraday_asset_prices_5min_mv.asset_id ] }
│ └── StreamHashAgg { group_key: [intraday_asset_prices_5min_mv.asset_id], aggs: [count] } { output: [ intraday_asset_prices_5min_mv.asset_id, count ], stream key: [ intraday_asset_prices_5min_mv.asset_id ] }
│ └── MergeExecutor { output: [ intraday_asset_prices_5min_mv.asset_id, $src ], stream key: [ intraday_asset_prices_5min_mv.asset_id, $src ] }
└── MergeExecutor { output: [ intraday_asset_prices_5min_mv.asset_id, first_value(intraday_asset_prices_5min_mv.last_price order_by(intraday_asset_prices_5min_mv.window_start DESC)) filter(IsNotNull(intraday_asset_prices_5min_mv.last_price) AND IsNotNull(intraday_asset_prices_5min_mv.window_start)), max(intraday_asset_prices_5min_mv.value_timestamp) ], stream key: [ intraday_asset_prices_5min_mv.asset_id ] }
Fragment 25267 (Actor 122750,122749)
StreamUnion { all: true } { output: [ intraday_asset_prices_5min_mv.asset_id, $src ], stream key: [ intraday_asset_prices_5min_mv.asset_id, $src ] }
├── MergeExecutor { output: [ intraday_asset_prices_5min_mv.asset_id, 0:Int32 ], stream key: [ intraday_asset_prices_5min_mv.asset_id ] }
└── MergeExecutor { output: [ asset_prices_snapshot_mv_next.asset_id, 1:Int32 ], stream key: [ asset_prices_snapshot_mv_next.asset_id ] }
Fragment 25268 (Actor 122758,122757)
StreamProject { exprs: [intraday_asset_prices_5min_mv.asset_id, 0:Int32] } { output: [ intraday_asset_prices_5min_mv.asset_id, 0:Int32 ], stream key: [ intraday_asset_prices_5min_mv.asset_id ] }
└── MergeExecutor { output: [ intraday_asset_prices_5min_mv.asset_id, first_value(intraday_asset_prices_5min_mv.last_price order_by(intraday_asset_prices_5min_mv.window_start DESC)) filter(IsNotNull(intraday_asset_prices_5min_mv.last_price) AND IsNotNull(intraday_asset_prices_5min_mv.window_start)), max(intraday_asset_prices_5min_mv.value_timestamp) ], stream key: [ intraday_asset_prices_5min_mv.asset_id ] }
Fragment 25269 (Actor 122751,122752)
StreamProject { exprs: [intraday_asset_prices_5min_mv.asset_id, first_value(intraday_asset_prices_5min_mv.last_price order_by(intraday_asset_prices_5min_mv.window_start DESC)) filter(IsNotNull(intraday_asset_prices_5min_mv.last_price) AND IsNotNull(intraday_asset_prices_5min_mv.window_start)), max(intraday_asset_prices_5min_mv.value_timestamp)] }
├── output: [ intraday_asset_prices_5min_mv.asset_id, first_value(intraday_asset_prices_5min_mv.last_price order_by(intraday_asset_prices_5min_mv.window_start DESC)) filter(IsNotNull(intraday_asset_prices_5min_mv.last_price) AND IsNotNull(intraday_asset_prices_5min_mv.window_start)), max(intraday_asset_prices_5min_mv.value_timestamp) ]
├── stream key: [ intraday_asset_prices_5min_mv.asset_id ]
└── StreamHashAgg [append_only] { group_key: [intraday_asset_prices_5min_mv.asset_id], aggs: [first_value(intraday_asset_prices_5min_mv.last_price order_by(intraday_asset_prices_5min_mv.window_start DESC)) filter(IsNotNull(intraday_asset_prices_5min_mv.last_price) AND IsNotNull(intraday_asset_prices_5min_mv.window_start)), max(intraday_asset_prices_5min_mv.value_timestamp), count] }
├── output: [ intraday_asset_prices_5min_mv.asset_id, first_value(intraday_asset_prices_5min_mv.last_price order_by(intraday_asset_prices_5min_mv.window_start DESC)) filter(IsNotNull(intraday_asset_prices_5min_mv.last_price) AND IsNotNull(intraday_asset_prices_5min_mv.window_start)), max(intraday_asset_prices_5min_mv.value_timestamp), count ]
├── stream key: [ intraday_asset_prices_5min_mv.asset_id ]
└── MergeExecutor { output: [ intraday_asset_prices_5min_mv.window_start, intraday_asset_prices_5min_mv.asset_id, intraday_asset_prices_5min_mv.last_price, intraday_asset_prices_5min_mv.value_timestamp ], stream key: [ intraday_asset_prices_5min_mv.window_start, intraday_asset_prices_5min_mv.asset_id ] }
Fragment 25270 (Actor 122748,122747)
StreamTableScan { table: intraday_asset_prices_5min_mv, columns: [window_start, asset_id, last_price, value_timestamp] } { output: [ intraday_asset_prices_5min_mv.window_start, intraday_asset_prices_5min_mv.asset_id, intraday_asset_prices_5min_mv.last_price, intraday_asset_prices_5min_mv.value_timestamp ], stream key: [ intraday_asset_prices_5min_mv.window_start, intraday_asset_prices_5min_mv.asset_id ] }
├── Upstream { output: [ window_start, asset_id, last_price, value_timestamp ], stream key: [] }
└── BatchPlanNode { output: [ window_start, asset_id, last_price, value_timestamp ], stream key: [] }
Fragment 25271 (Actor 122762,122761)
StreamProject { exprs: [asset_prices_snapshot_mv_next.asset_id, 1:Int32] } { output: [ asset_prices_snapshot_mv_next.asset_id, 1:Int32 ], stream key: [ asset_prices_snapshot_mv_next.asset_id ] }
└── StreamTableScan { table: asset_prices_snapshot_mv_next, columns: [asset_id] } { output: [ asset_prices_snapshot_mv_next.asset_id ], stream key: [ asset_prices_snapshot_mv_next.asset_id ] }
├── Upstream { output: [ asset_id ], stream key: [] }
└── BatchPlanNode { output: [ asset_id ], stream key: [] }
Fragment 25272 (Actor 122764,122763)
StreamTableScan { table: asset_prices_snapshot_mv_next, columns: [asset_id, close, value_timestamp] } { output: [ asset_prices_snapshot_mv_next.asset_id, asset_prices_snapshot_mv_next.close, asset_prices_snapshot_mv_next.value_timestamp ], stream key: [ asset_prices_snapshot_mv_next.asset_id ] }
├── Upstream { output: [ asset_id, close, value_timestamp ], stream key: [] }
└── BatchPlanNode { output: [ asset_id, close, value_timestamp ], stream key: [] }