Job is idle — throughput ~0; structure shown.
Fragment 54566 (Actor 740065,740066)
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.asset_id, asset_prices_eod_ft.date, $expr1, first_value ]
├── stream key: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date ]
└── StreamOverWindow { window_functions: [first_value($expr1) OVER(PARTITION BY asset_prices_eod_ft.asset_id ORDER BY asset_prices_eod_ft.date ASC ROWS BETWEEN 1 PRECEDING AND 1 PRECEDING)] }
├── output: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date, $expr1, first_value ]
├── stream key: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date ]
└── StreamLocalityProvider { locality_columns: [asset_prices_eod_ft.asset_id] }
├── output: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date, $expr1 ]
├── stream key: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date ]
└── MergeExecutor { output: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date, $expr1 ], stream key: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date ] }
Fragment 54567 (Actor 740069,740070)
StreamProject { exprs: [asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date, Coalesce(asset_adjusted_prices_ft.adjusted_close, asset_prices_eod_ft.close) as $expr1], output_watermarks: [[asset_prices_eod_ft.date]] }
├── output: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date, $expr1 ]
├── stream key: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date ]
└── MergeExecutor
├── output: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date, asset_prices_eod_ft.close, asset_adjusted_prices_ft.adjusted_close, asset_adjusted_prices_ft.asset_id, asset_adjusted_prices_ft.date ]
└── stream key: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date ]
Fragment 54568 (Actor 740068,740067)
StreamSyncLogStore { output: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date, asset_prices_eod_ft.close, asset_adjusted_prices_ft.adjusted_close, asset_adjusted_prices_ft.asset_id, asset_adjusted_prices_ft.date ], stream key: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date ] }
└── StreamHashJoin [window] { type: LeftOuter, predicate: asset_prices_eod_ft.date = asset_adjusted_prices_ft.date AND asset_prices_eod_ft.asset_id = asset_adjusted_prices_ft.asset_id, conditions_to_clean_state_in_join_key: [(asset_prices_eod_ft.date = asset_adjusted_prices_ft.date)], output_watermarks: [[asset_prices_eod_ft.date], [asset_adjusted_prices_ft.date]] }
├── output: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date, asset_prices_eod_ft.close, asset_adjusted_prices_ft.adjusted_close, asset_adjusted_prices_ft.asset_id, asset_adjusted_prices_ft.date ]
├── stream key: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date ]
├── MergeExecutor { output: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date, asset_prices_eod_ft.close ], stream key: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.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 54569 (Actor 740083,740084)
StreamFilter { predicate: Not(IsNull(asset_prices_eod_ft.close)) } { output: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date, asset_prices_eod_ft.close ], stream key: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date ] }
└── StreamTableScan { table: asset_prices_eod_ft, columns: [asset_id, date, close] } { output: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date, asset_prices_eod_ft.close ], stream key: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date ] }
├── Upstream { output: [ asset_id, date, close ], stream key: [] }
└── BatchPlanNode { output: [ asset_id, date, close ], stream key: [] }
Fragment 54570 (Actor 740073,740074)
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: [] }