Job is idle — throughput ~0; structure shown.
Fragment 17892 (Actor 157829,157828)
StreamMaterialize { columns: [asset_id, fact_date, adjusted_last_close_price, prev_price], stream_key: [asset_id, fact_date], pk_columns: [asset_id, fact_date], pk_conflict: NoCheck }
├── output: [ asset_prices_eod_ft_next.asset_id, asset_prices_eod_ft_next.date, $expr1, first_value ]
├── stream key: [ asset_prices_eod_ft_next.asset_id, asset_prices_eod_ft_next.date ]
└── StreamOverWindow { window_functions: [first_value($expr1) OVER(PARTITION BY asset_prices_eod_ft_next.asset_id ORDER BY asset_prices_eod_ft_next.date ASC ROWS BETWEEN 1 PRECEDING AND 1 PRECEDING)] }
├── output: [ asset_prices_eod_ft_next.asset_id, asset_prices_eod_ft_next.date, $expr1, first_value ]
├── stream key: [ asset_prices_eod_ft_next.asset_id, asset_prices_eod_ft_next.date ]
└── StreamLocalityProvider { locality_columns: [asset_prices_eod_ft_next.asset_id] }
├── output: [ asset_prices_eod_ft_next.asset_id, asset_prices_eod_ft_next.date, $expr1 ]
├── stream key: [ asset_prices_eod_ft_next.asset_id, asset_prices_eod_ft_next.date ]
└── MergeExecutor { output: [ asset_prices_eod_ft_next.asset_id, asset_prices_eod_ft_next.date, $expr1 ], stream key: [ asset_prices_eod_ft_next.asset_id, asset_prices_eod_ft_next.date ] }
Fragment 17893 (Actor 157830,157831)
StreamProject { exprs: [asset_prices_eod_ft_next.asset_id, asset_prices_eod_ft_next.date, Coalesce(asset_adjusted_prices_ft.adjusted_close, asset_prices_eod_ft_next.close) as $expr1], output_watermarks: [[asset_prices_eod_ft_next.date]] }
├── output: [ asset_prices_eod_ft_next.asset_id, asset_prices_eod_ft_next.date, $expr1 ]
├── stream key: [ asset_prices_eod_ft_next.asset_id, asset_prices_eod_ft_next.date ]
└── MergeExecutor
├── output: [ asset_prices_eod_ft_next.asset_id, asset_prices_eod_ft_next.date, asset_prices_eod_ft_next.close, asset_adjusted_prices_ft.adjusted_close, asset_adjusted_prices_ft.asset_id, asset_adjusted_prices_ft.date ]
└── stream key: [ asset_prices_eod_ft_next.asset_id, asset_prices_eod_ft_next.date ]
Fragment 17894 (Actor 157832,157833)
StreamSyncLogStore { output: [ asset_prices_eod_ft_next.asset_id, asset_prices_eod_ft_next.date, asset_prices_eod_ft_next.close, asset_adjusted_prices_ft.adjusted_close, asset_adjusted_prices_ft.asset_id, asset_adjusted_prices_ft.date ], stream key: [ asset_prices_eod_ft_next.asset_id, asset_prices_eod_ft_next.date ] }
└── StreamHashJoin [window] { type: LeftOuter, predicate: asset_prices_eod_ft_next.date = asset_adjusted_prices_ft.date AND asset_prices_eod_ft_next.asset_id = asset_adjusted_prices_ft.asset_id, conditions_to_clean_state_in_join_key: [(asset_prices_eod_ft_next.date = asset_adjusted_prices_ft.date)], output_watermarks: [[asset_prices_eod_ft_next.date], [asset_adjusted_prices_ft.date]] }
├── output: [ asset_prices_eod_ft_next.asset_id, asset_prices_eod_ft_next.date, asset_prices_eod_ft_next.close, asset_adjusted_prices_ft.adjusted_close, asset_adjusted_prices_ft.asset_id, asset_adjusted_prices_ft.date ]
├── stream key: [ asset_prices_eod_ft_next.asset_id, asset_prices_eod_ft_next.date ]
├── MergeExecutor { output: [ asset_prices_eod_ft_next.asset_id, asset_prices_eod_ft_next.date, asset_prices_eod_ft_next.close ], stream key: [ asset_prices_eod_ft_next.asset_id, asset_prices_eod_ft_next.date ] }
└── MergeExecutor { output: [ asset_adjusted_prices_ft.asset_id, asset_adjusted_prices_ft.date, asset_adjusted_prices_ft.adjusted_close ], stream key: [ asset_adjusted_prices_ft.asset_id, asset_adjusted_prices_ft.date ] }
Fragment 17895 (Actor 157838,157839)
StreamFilter { predicate: Not(IsNull(asset_prices_eod_ft_next.close)) } { output: [ asset_prices_eod_ft_next.asset_id, asset_prices_eod_ft_next.date, asset_prices_eod_ft_next.close ], stream key: [ asset_prices_eod_ft_next.asset_id, asset_prices_eod_ft_next.date ] }
└── StreamTableScan { table: asset_prices_eod_ft_next, columns: [asset_id, date, close] } { output: [ asset_prices_eod_ft_next.asset_id, asset_prices_eod_ft_next.date, asset_prices_eod_ft_next.close ], stream key: [ asset_prices_eod_ft_next.asset_id, asset_prices_eod_ft_next.date ] }
├── Upstream { output: [ asset_id, date, close ], stream key: [] }
└── BatchPlanNode { output: [ asset_id, date, close ], stream key: [] }
Fragment 17896 (Actor 157840,157841)
StreamTableScan { table: asset_adjusted_prices_ft, columns: [asset_id, date, adjusted_close] } { output: [ asset_adjusted_prices_ft.asset_id, asset_adjusted_prices_ft.date, asset_adjusted_prices_ft.adjusted_close ], stream key: [ asset_adjusted_prices_ft.asset_id, asset_adjusted_prices_ft.date ] }
├── Upstream { output: [ asset_id, date, adjusted_close ], stream key: [] }
└── BatchPlanNode { output: [ asset_id, date, adjusted_close ], stream key: [] }