Job is idle — throughput ~0; structure shown.
Fragment 5348 (Actor 96821,96822)
StreamMaterialize { columns: [account_id, asset_id, currency_code, purchased_quantity, average_cost_per_unit, market_value, holding_timestamp], stream_key: [account_id, asset_id], pk_columns: [account_id, asset_id], pk_conflict: NoCheck }
├── output: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.asset_id, holding_values_intraday_ft.currency_code, holding_values_intraday_ft.purchased_quantity, holding_values_intraday_ft.average_cost_per_unit, holding_values_intraday_ft.market_value, holding_values_intraday_ft.holding_timestamp ]
├── stream key: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.asset_id ]
└── StreamProject { exprs: [holding_values_intraday_ft.account_id, holding_values_intraday_ft.asset_id, holding_values_intraday_ft.currency_code, holding_values_intraday_ft.purchased_quantity, holding_values_intraday_ft.average_cost_per_unit, holding_values_intraday_ft.market_value, holding_values_intraday_ft.holding_timestamp] }
├── output: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.asset_id, holding_values_intraday_ft.currency_code, holding_values_intraday_ft.purchased_quantity, holding_values_intraday_ft.average_cost_per_unit, holding_values_intraday_ft.market_value, holding_values_intraday_ft.holding_timestamp ]
├── stream key: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.asset_id ]
└── StreamGroupTopN { order: [holding_values_intraday_ft.holding_timestamp DESC], limit: 1, offset: 0, group_key: [holding_values_intraday_ft.account_id, holding_values_intraday_ft.asset_id] }
├── output:
│ ┌── holding_values_intraday_ft.account_id
│ ├── holding_values_intraday_ft.asset_id
│ ├── holding_values_intraday_ft.currency_code
│ ├── holding_values_intraday_ft.holding_timestamp
│ ├── holding_values_intraday_ft.purchased_quantity
│ ├── holding_values_intraday_ft.market_value
│ ├── holding_values_intraday_ft.average_cost_per_unit
│ └── holding_values_intraday_ft.id
├── stream key: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.asset_id ]
└── MergeExecutor
├── output:
│ ┌── holding_values_intraday_ft.account_id
│ ├── holding_values_intraday_ft.asset_id
│ ├── holding_values_intraday_ft.currency_code
│ ├── holding_values_intraday_ft.holding_timestamp
│ ├── holding_values_intraday_ft.purchased_quantity
│ ├── holding_values_intraday_ft.market_value
│ ├── holding_values_intraday_ft.average_cost_per_unit
│ └── holding_values_intraday_ft.id
└── stream key: [ holding_values_intraday_ft.id, holding_values_intraday_ft.holding_timestamp ]
Fragment 5349 (Actor 98247,98248)
StreamProject { exprs: [holding_values_intraday_ft.account_id, holding_values_intraday_ft.asset_id, holding_values_intraday_ft.currency_code, holding_values_intraday_ft.holding_timestamp, holding_values_intraday_ft.purchased_quantity, holding_values_intraday_ft.market_value, holding_values_intraday_ft.average_cost_per_unit, holding_values_intraday_ft.id], output_watermarks: [[holding_values_intraday_ft.holding_timestamp]] }
├── output: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.asset_id, holding_values_intraday_ft.currency_code, holding_values_intraday_ft.holding_timestamp, holding_values_intraday_ft.purchased_quantity, holding_values_intraday_ft.market_value, holding_values_intraday_ft.average_cost_per_unit, holding_values_intraday_ft.id ]
├── stream key: [ holding_values_intraday_ft.id, holding_values_intraday_ft.holding_timestamp ]
└── StreamFilter { predicate: IsNull(holding_values_intraday_ft.disabled_at) }
├── output: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.asset_id, holding_values_intraday_ft.currency_code, holding_values_intraday_ft.holding_timestamp, holding_values_intraday_ft.purchased_quantity, holding_values_intraday_ft.market_value, holding_values_intraday_ft.average_cost_per_unit, holding_values_intraday_ft.id, holding_values_intraday_ft.disabled_at ]
├── stream key: [ holding_values_intraday_ft.id, holding_values_intraday_ft.holding_timestamp ]
└── StreamTableScan { table: holding_values_intraday_ft, columns: [account_id, asset_id, currency_code, holding_timestamp, purchased_quantity, market_value, average_cost_per_unit, id, disabled_at] }
├── output: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.asset_id, holding_values_intraday_ft.currency_code, holding_values_intraday_ft.holding_timestamp, holding_values_intraday_ft.purchased_quantity, holding_values_intraday_ft.market_value, holding_values_intraday_ft.average_cost_per_unit, holding_values_intraday_ft.id, holding_values_intraday_ft.disabled_at ]
├── stream key: [ holding_values_intraday_ft.id, holding_values_intraday_ft.holding_timestamp ]
├── Upstream { output: [ account_id, asset_id, currency_code, holding_timestamp, purchased_quantity, market_value, average_cost_per_unit, id, disabled_at ], stream key: [] }
└── BatchPlanNode { output: [ account_id, asset_id, currency_code, holding_timestamp, purchased_quantity, market_value, average_cost_per_unit, id, disabled_at ], stream key: [] }