Job is idle — throughput ~0; structure shown.
Fragment 62436 (Actor 742128,742127)
StreamMaterialize { columns: [asset_id, close, value_timestamp], stream_key: [asset_id], pk_columns: [asset_id], pk_conflict: NoCheck }
├── output: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.close, $expr1 ]
├── stream key: [ asset_prices_eod_ft.asset_id ]
└── StreamProject { exprs: [asset_prices_eod_ft.asset_id, asset_prices_eod_ft.close, AtTimeZone(asset_prices_eod_ft.date::Timestamp, 'UTC':Varchar) as $expr1] }
├── output: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.close, $expr1 ]
├── stream key: [ asset_prices_eod_ft.asset_id ]
└── StreamGroupTopN { order: [asset_prices_eod_ft.date DESC], limit: 1, offset: 0, group_key: [asset_prices_eod_ft.asset_id] }
├── 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 ]
└── StreamLocalityProvider { locality_columns: [asset_prices_eod_ft.asset_id] }
├── 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_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 ]
Fragment 62437 (Actor 742142,742141)
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: [] }