RWM Console cluster: risingwave-ai-hub-int.ai-hub-rwm.svc.cluster.local

← cluster insights objects holding_values_journal_density_mv explain
Overview Objects Graph History
materialized view · insights.holding_values_journal_density_mv profiled over 5s
seconds (1–30)

Job is idle — throughput ~0; structure shown.

Stateful hash join (4 state tables) — consider a temporal join for dimension lookupsAggregation state — unbounded unless keyed or temporally filteredWindow state — add a WHERE rank <= N to bound it
64 operators
Materialize · insights.holding_values_journal_density_mv
0% idle 2 actors
Project
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · LeftOuter · settled_position_series_mv.account_id = settled_cost_basis_…
2 actors
HashJoin · LeftOuter · settled_position_series_mv.account_id = settled_cost_basis_… Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · settled_cost_basis_series_mv
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
LocalityProvider
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · (asset_prices_eod_ft.date >= settled_position_series_mv.dim…
2 actors
Filter · (asset_prices_eod_ft.date >= settled_position_series_mv.dim…
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · settled_position_series_mv.asset_id = asset_prices_eod_ft.a…
2 actors
HashJoin · Inner · settled_position_series_mv.asset_id = asset_prices_eod_ft.a… Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
LocalityProvider
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Filter · asset_prices_eod_ft
0% idle 2 actors
StreamScan · asset_prices_eod_ft
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
LocalityProvider
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · ($expr1 <= $expr2) AND (IsNull(first_value) OR ($expr2 <= f…
2 actors
Filter · ($expr1 <= $expr2) AND (IsNull(first_value) OR ($expr2 <= f…
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · settled_position_series_mv.asset_id = asset_prices_eod_ft.a…
2 actors
HashJoin · Inner · settled_position_series_mv.asset_id = asset_prices_eod_ft.a… Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
2 actors
HashAgg Aggregation state — unbounded unless keyed or temporally filtered
0% idle 2 actors
LocalityProvider
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
LocalityProvider
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · settled_position_series_mv
2 actors
OverWindow · settled_position_series_mv Window state — add a WHERE rank <= N to bound it
0% idle 2 actors
StreamScan · settled_position_series_mv
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Heat = the operator's output-buffer backpressure over the sampling window. Click a node to fold its subtree.
Materialize · insights.holding_values_journal_density_mv Materialize insights.holding_values… idle · 2 actors Project Project — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · LeftOuter · settled_position_series_mv.account_id = settled_cost_basis_… SyncLogStore LeftOuter · settled_pos… — · 2 actors HashJoin · LeftOuter · settled_position_series_mv.account_id = settled_cost_basis_… HashJoin LeftOuter · settled_pos… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · settled_cost_basis_series_mv StreamScan settled_cost_basis_seri… idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors LocalityProvider LocalityProvider idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · (asset_prices_eod_ft.date >= settled_position_series_mv.dim… Project (asset_prices_eod_ft.da… — · 2 actors Filter · (asset_prices_eod_ft.date >= settled_position_series_mv.dim… Filter (asset_prices_eod_ft.da… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · settled_position_series_mv.asset_id = asset_prices_eod_ft.a… SyncLogStore Inner · settled_positio… — · 2 actors HashJoin · Inner · settled_position_series_mv.asset_id = asset_prices_eod_ft.a… HashJoin Inner · settled_positio… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors LocalityProvider LocalityProvider idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Filter · asset_prices_eod_ft Filter asset_prices_eod_ft idle · 2 actors StreamScan · asset_prices_eod_ft StreamScan asset_prices_eod_ft idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors LocalityProvider LocalityProvider idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · ($expr1 <= $expr2) AND (IsNull(first_value) OR ($expr2 <= f… Project ($expr1 <= $expr2) AND … — · 2 actors Filter · ($expr1 <= $expr2) AND (IsNull(first_value) OR ($expr2 <= f… Filter ($expr1 <= $expr2) AND … idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · settled_position_series_mv.asset_id = asset_prices_eod_ft.a… SyncLogStore Inner · settled_positio… — · 2 actors HashJoin · Inner · settled_position_series_mv.asset_id = asset_prices_eod_ft.a… HashJoin Inner · settled_positio… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 2 actors HashAgg HashAgg idle · 2 actors LocalityProvider LocalityProvider idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors LocalityProvider LocalityProvider idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · settled_position_series_mv Project settled_position_series… — · 2 actors OverWindow · settled_position_series_mv OverWindow settled_position_series… idle · 2 actors StreamScan · settled_position_series_mv StreamScan settled_position_series… idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors
Streaming operator plan from EXPLAIN ANALYZE. Node heat = backpressure. Drag to pan, scroll to zoom.
Fragments (DESCRIBE FRAGMENTS) — click to expand
Fragment 19464 (Actor 164490,164491)
StreamMaterialize { columns: [account_id, asset_id, dim_value_date, type, currency_code, market_value, average_cost_per_unit, average_cost_per_unit_system_currency, total_cost_system_currency, cost_fx_provenance, purchased_quantity, $expr2(hidden), settled_position_series_mv.dim_settlement_date(hidden), settled_cost_basis_series_mv.settled_cost_basis_carried_mv.sum(hidden), settled_cost_basis_series_mv.effective_from(hidden)], stream_key: [account_id, asset_id, currency_code, $expr2, settled_position_series_mv.dim_settlement_date, dim_value_date, settled_cost_basis_series_mv.settled_cost_basis_carried_mv.sum, settled_cost_basis_series_mv.effective_from], pk_columns: [account_id, asset_id, currency_code, $expr2, settled_position_series_mv.dim_settlement_date, dim_value_date, settled_cost_basis_series_mv.settled_cost_basis_carried_mv.sum, settled_cost_basis_series_mv.effective_from], pk_conflict: NoCheck }
├── output: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, asset_prices_eod_ft.date, 'ASSET':Varchar, settled_position_series_mv.currency_code, $expr4, settled_cost_basis_series_mv.average_cost_per_unit, settled_cost_basis_series_mv.average_cost_per_unit_system_currency, settled_cost_basis_series_mv.total_cost_system_currency, settled_cost_basis_series_mv.cost_fx_provenance, settled_position_series_mv.settled_quantity, $expr2, settled_position_series_mv.dim_settlement_date, settled_cost_basis_series_mv.settled_cost_basis_carried_mv.sum, settled_cost_basis_series_mv.effective_from ]
├── stream key: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, $expr2, settled_position_series_mv.dim_settlement_date, asset_prices_eod_ft.date, settled_cost_basis_series_mv.settled_cost_basis_carried_mv.sum, settled_cost_basis_series_mv.effective_from ]
└── StreamProject { exprs: [settled_position_series_mv.account_id, settled_position_series_mv.asset_id, asset_prices_eod_ft.date, 'ASSET':Varchar, settled_position_series_mv.currency_code, (settled_position_series_mv.settled_quantity * asset_prices_eod_ft.reference_price) as $expr4, settled_cost_basis_series_mv.average_cost_per_unit, settled_cost_basis_series_mv.average_cost_per_unit_system_currency, settled_cost_basis_series_mv.total_cost_system_currency, settled_cost_basis_series_mv.cost_fx_provenance, settled_position_series_mv.settled_quantity, $expr2, settled_position_series_mv.dim_settlement_date, settled_cost_basis_series_mv.settled_cost_basis_carried_mv.sum, settled_cost_basis_series_mv.effective_from] }
    ├── output: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, asset_prices_eod_ft.date, 'ASSET':Varchar, settled_position_series_mv.currency_code, $expr4, settled_cost_basis_series_mv.average_cost_per_unit, settled_cost_basis_series_mv.average_cost_per_unit_system_currency, settled_cost_basis_series_mv.total_cost_system_currency, settled_cost_basis_series_mv.cost_fx_provenance, settled_position_series_mv.settled_quantity, $expr2, settled_position_series_mv.dim_settlement_date, settled_cost_basis_series_mv.settled_cost_basis_carried_mv.sum, settled_cost_basis_series_mv.effective_from ]
    ├── stream key: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, $expr2, settled_position_series_mv.dim_settlement_date, asset_prices_eod_ft.date, settled_cost_basis_series_mv.settled_cost_basis_carried_mv.sum, settled_cost_basis_series_mv.effective_from ]
    └── MergeExecutor
        ├── output: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.settled_quantity, asset_prices_eod_ft.date, asset_prices_eod_ft.reference_price, settled_cost_basis_series_mv.average_cost_per_unit, settled_cost_basis_series_mv.average_cost_per_unit_system_currency, settled_cost_basis_series_mv.total_cost_system_currency, settled_cost_basis_series_mv.cost_fx_provenance, $expr2, settled_position_series_mv.dim_settlement_date, settled_cost_basis_series_mv.account_id, settled_cost_basis_series_mv.asset_id, settled_cost_basis_series_mv.currency_code, settled_cost_basis_series_mv.settled_cost_basis_carried_mv.sum, settled_cost_basis_series_mv.effective_from ]
        └── stream key: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, $expr2, settled_position_series_mv.dim_settlement_date, asset_prices_eod_ft.date, settled_cost_basis_series_mv.settled_cost_basis_carried_mv.sum, settled_cost_basis_series_mv.effective_from ]

Fragment 19465 (Actor 164478,164479)
StreamSyncLogStore
├── output: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.settled_quantity, asset_prices_eod_ft.date, asset_prices_eod_ft.reference_price, settled_cost_basis_series_mv.average_cost_per_unit, settled_cost_basis_series_mv.average_cost_per_unit_system_currency, settled_cost_basis_series_mv.total_cost_system_currency, settled_cost_basis_series_mv.cost_fx_provenance, $expr2, settled_position_series_mv.dim_settlement_date, settled_cost_basis_series_mv.account_id, settled_cost_basis_series_mv.asset_id, settled_cost_basis_series_mv.currency_code, settled_cost_basis_series_mv.settled_cost_basis_carried_mv.sum, settled_cost_basis_series_mv.effective_from ]
├── stream key: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, $expr2, settled_position_series_mv.dim_settlement_date, asset_prices_eod_ft.date, settled_cost_basis_series_mv.settled_cost_basis_carried_mv.sum, settled_cost_basis_series_mv.effective_from ]
└── StreamHashJoin { type: LeftOuter, predicate: settled_position_series_mv.account_id = settled_cost_basis_series_mv.account_id AND settled_position_series_mv.asset_id = settled_cost_basis_series_mv.asset_id AND settled_position_series_mv.currency_code = settled_cost_basis_series_mv.currency_code AND (settled_cost_basis_series_mv.effective_from <= asset_prices_eod_ft.date) AND (IsNull(settled_cost_basis_series_mv.effective_to) OR (asset_prices_eod_ft.date < settled_cost_basis_series_mv.effective_to)) }
    ├── output: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.settled_quantity, asset_prices_eod_ft.date, asset_prices_eod_ft.reference_price, settled_cost_basis_series_mv.average_cost_per_unit, settled_cost_basis_series_mv.average_cost_per_unit_system_currency, settled_cost_basis_series_mv.total_cost_system_currency, settled_cost_basis_series_mv.cost_fx_provenance, $expr2, settled_position_series_mv.dim_settlement_date, settled_cost_basis_series_mv.account_id, settled_cost_basis_series_mv.asset_id, settled_cost_basis_series_mv.currency_code, settled_cost_basis_series_mv.settled_cost_basis_carried_mv.sum, settled_cost_basis_series_mv.effective_from ]
    ├── stream key: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, $expr2, settled_position_series_mv.dim_settlement_date, asset_prices_eod_ft.date, settled_cost_basis_series_mv.settled_cost_basis_carried_mv.sum, settled_cost_basis_series_mv.effective_from ]
    ├── MergeExecutor { output: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.settled_quantity, asset_prices_eod_ft.date, asset_prices_eod_ft.reference_price, $expr2, settled_position_series_mv.dim_settlement_date, asset_prices_eod_ft.asset_id, $expr3 ], stream key: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, $expr2, settled_position_series_mv.dim_settlement_date, asset_prices_eod_ft.date ] }
    └── MergeExecutor { output: [ settled_cost_basis_series_mv.account_id, settled_cost_basis_series_mv.asset_id, settled_cost_basis_series_mv.currency_code, settled_cost_basis_series_mv.effective_from, settled_cost_basis_series_mv.effective_to, settled_cost_basis_series_mv.average_cost_per_unit, settled_cost_basis_series_mv.average_cost_per_unit_system_currency, settled_cost_basis_series_mv.total_cost_system_currency, settled_cost_basis_series_mv.cost_fx_provenance, settled_cost_basis_series_mv.settled_cost_basis_carried_mv.sum ], stream key: [ settled_cost_basis_series_mv.account_id, settled_cost_basis_series_mv.asset_id, settled_cost_basis_series_mv.currency_code, settled_cost_basis_series_mv.settled_cost_basis_carried_mv.sum, settled_cost_basis_series_mv.effective_from ] }

Fragment 19466 (Actor 164470,164471)
StreamLocalityProvider { locality_columns: [settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code] } { output: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.settled_quantity, asset_prices_eod_ft.date, asset_prices_eod_ft.reference_price, $expr2, settled_position_series_mv.dim_settlement_date, asset_prices_eod_ft.asset_id, $expr3 ], stream key: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, $expr2, settled_position_series_mv.dim_settlement_date, asset_prices_eod_ft.date ] }
└── MergeExecutor { output: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.settled_quantity, asset_prices_eod_ft.date, asset_prices_eod_ft.reference_price, $expr2, settled_position_series_mv.dim_settlement_date, asset_prices_eod_ft.asset_id, $expr3 ], stream key: [ settled_position_series_mv.asset_id, $expr2, settled_position_series_mv.account_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, asset_prices_eod_ft.date ] }

Fragment 19467 (Actor 164485,164484)
StreamProject { exprs: [settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.settled_quantity, asset_prices_eod_ft.date, asset_prices_eod_ft.reference_price, $expr2, settled_position_series_mv.dim_settlement_date, asset_prices_eod_ft.asset_id, $expr3] } { output: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.settled_quantity, asset_prices_eod_ft.date, asset_prices_eod_ft.reference_price, $expr2, settled_position_series_mv.dim_settlement_date, asset_prices_eod_ft.asset_id, $expr3 ], stream key: [ settled_position_series_mv.asset_id, $expr2, settled_position_series_mv.account_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, asset_prices_eod_ft.date ] }
└── StreamFilter { predicate: (asset_prices_eod_ft.date >= settled_position_series_mv.dim_settlement_date) AND (IsNull(first_value) OR (asset_prices_eod_ft.date < first_value)) } { output: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, settled_position_series_mv.settled_quantity, first_value, $expr2, asset_prices_eod_ft.asset_id, asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date, asset_prices_eod_ft.reference_price, $expr3 ], stream key: [ settled_position_series_mv.asset_id, $expr2, settled_position_series_mv.account_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, asset_prices_eod_ft.date ] }
    └── MergeExecutor { output: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, settled_position_series_mv.settled_quantity, first_value, $expr2, asset_prices_eod_ft.asset_id, asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date, asset_prices_eod_ft.reference_price, $expr3 ], stream key: [ settled_position_series_mv.asset_id, $expr2, settled_position_series_mv.account_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, asset_prices_eod_ft.date ] }

Fragment 19468 (Actor 164488,164489)
StreamSyncLogStore { output: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, settled_position_series_mv.settled_quantity, first_value, $expr2, asset_prices_eod_ft.asset_id, asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date, asset_prices_eod_ft.reference_price, $expr3 ], stream key: [ settled_position_series_mv.asset_id, $expr2, settled_position_series_mv.account_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, asset_prices_eod_ft.date ] }
└── StreamHashJoin { type: Inner, predicate: settled_position_series_mv.asset_id = asset_prices_eod_ft.asset_id AND $expr2 = $expr3 } { output: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, settled_position_series_mv.settled_quantity, first_value, $expr2, asset_prices_eod_ft.asset_id, asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date, asset_prices_eod_ft.reference_price, $expr3 ], stream key: [ settled_position_series_mv.asset_id, $expr2, settled_position_series_mv.account_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, asset_prices_eod_ft.date ] }
    ├── MergeExecutor { output: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, settled_position_series_mv.settled_quantity, first_value, $expr2, asset_prices_eod_ft.asset_id ], stream key: [ settled_position_series_mv.asset_id, $expr2, settled_position_series_mv.account_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date ] }
    └── MergeExecutor { output: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date, asset_prices_eod_ft.reference_price, $expr3 ], stream key: [ asset_prices_eod_ft.asset_id, $expr3, asset_prices_eod_ft.date ] }

Fragment 19469 (Actor 164492,164493)
StreamLocalityProvider { locality_columns: [settled_position_series_mv.asset_id, $expr2] } { output: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, settled_position_series_mv.settled_quantity, first_value, $expr2, asset_prices_eod_ft.asset_id ], stream key: [ settled_position_series_mv.asset_id, $expr2, settled_position_series_mv.account_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date ] }
└── MergeExecutor { output: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, settled_position_series_mv.settled_quantity, first_value, $expr2, asset_prices_eod_ft.asset_id ], stream key: [ settled_position_series_mv.asset_id, settled_position_series_mv.account_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, $expr2 ] }

Fragment 19470 (Actor 164481,164480)
StreamProject { exprs: [settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, settled_position_series_mv.settled_quantity, first_value, $expr2, asset_prices_eod_ft.asset_id] } { output: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, settled_position_series_mv.settled_quantity, first_value, $expr2, asset_prices_eod_ft.asset_id ], stream key: [ settled_position_series_mv.asset_id, settled_position_series_mv.account_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, $expr2 ] }
└── StreamFilter { predicate: ($expr1 <= $expr2) AND (IsNull(first_value) OR ($expr2 <= first_value)) } { output: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, settled_position_series_mv.settled_quantity, first_value, $expr1, asset_prices_eod_ft.asset_id, $expr2 ], stream key: [ settled_position_series_mv.asset_id, settled_position_series_mv.account_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, $expr2 ] }
    └── MergeExecutor { output: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, settled_position_series_mv.settled_quantity, first_value, $expr1, asset_prices_eod_ft.asset_id, $expr2 ], stream key: [ settled_position_series_mv.asset_id, settled_position_series_mv.account_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, $expr2 ] }

Fragment 19471 (Actor 164472,164473)
StreamSyncLogStore { output: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, settled_position_series_mv.settled_quantity, first_value, $expr1, asset_prices_eod_ft.asset_id, $expr2 ], stream key: [ settled_position_series_mv.asset_id, settled_position_series_mv.account_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, $expr2 ] }
└── StreamHashJoin { type: Inner, predicate: settled_position_series_mv.asset_id = asset_prices_eod_ft.asset_id } { output: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, settled_position_series_mv.settled_quantity, first_value, $expr1, asset_prices_eod_ft.asset_id, $expr2 ], stream key: [ settled_position_series_mv.asset_id, settled_position_series_mv.account_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, $expr2 ] }
    ├── MergeExecutor { output: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, settled_position_series_mv.settled_quantity, first_value, $expr1 ], stream key: [ settled_position_series_mv.asset_id, settled_position_series_mv.account_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date ] }
    └── MergeExecutor { output: [ asset_prices_eod_ft.asset_id, $expr2 ], stream key: [ asset_prices_eod_ft.asset_id, $expr2 ] }

Fragment 19472 (Actor 164474,164475)
StreamLocalityProvider { locality_columns: [settled_position_series_mv.asset_id] } { output: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, settled_position_series_mv.settled_quantity, first_value, $expr1 ], stream key: [ settled_position_series_mv.asset_id, settled_position_series_mv.account_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date ] }
└── MergeExecutor { output: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, settled_position_series_mv.settled_quantity, first_value, $expr1 ], stream key: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date ] }

Fragment 19473 (Actor 164483,164482)
StreamProject { exprs: [settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, settled_position_series_mv.settled_quantity, first_value, AtTimeZone(DateTrunc('MONTH':Varchar, AtTimeZone(settled_position_series_mv.dim_settlement_date::Timestamp, 'UTC':Varchar), 'UTC':Varchar), 'UTC':Varchar)::Date as $expr1] } { output: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, settled_position_series_mv.settled_quantity, first_value, $expr1 ], stream key: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date ] }
└── StreamOverWindow { window_functions: [first_value(settled_position_series_mv.dim_settlement_date) OVER(PARTITION BY settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code ORDER BY settled_position_series_mv.dim_settlement_date ASC ROWS BETWEEN 1 FOLLOWING AND 1 FOLLOWING)] } { output: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, settled_position_series_mv.settled_quantity, first_value ], stream key: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date ] }
    └── StreamTableScan { table: settled_position_series_mv, columns: [account_id, asset_id, currency_code, dim_settlement_date, settled_quantity] } { output: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date, settled_position_series_mv.settled_quantity ], stream key: [ settled_position_series_mv.account_id, settled_position_series_mv.asset_id, settled_position_series_mv.currency_code, settled_position_series_mv.dim_settlement_date ] }
        ├── Upstream { output: [ account_id, asset_id, currency_code, dim_settlement_date, settled_quantity ], stream key: [] }
        └── BatchPlanNode { output: [ account_id, asset_id, currency_code, dim_settlement_date, settled_quantity ], stream key: [] }

Fragment 19474 (Actor 164487,164486)
StreamProject { exprs: [asset_prices_eod_ft.asset_id, $expr2] } { output: [ asset_prices_eod_ft.asset_id, $expr2 ], stream key: [ asset_prices_eod_ft.asset_id, $expr2 ] }
└── StreamHashAgg { group_key: [asset_prices_eod_ft.asset_id, $expr2], aggs: [count] } { output: [ asset_prices_eod_ft.asset_id, $expr2, count ], stream key: [ asset_prices_eod_ft.asset_id, $expr2 ] }
    └── StreamLocalityProvider { locality_columns: [asset_prices_eod_ft.asset_id, $expr2] } { output: [ asset_prices_eod_ft.asset_id, $expr2, asset_prices_eod_ft.date ], stream key: [ asset_prices_eod_ft.asset_id, $expr2, asset_prices_eod_ft.date ] }
        └── MergeExecutor { output: [ asset_prices_eod_ft.asset_id, $expr2, asset_prices_eod_ft.date ], stream key: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date ] }

Fragment 19475 (Actor 164467,164466)
StreamProject { exprs: [asset_prices_eod_ft.asset_id, AtTimeZone(DateTrunc('MONTH':Varchar, AtTimeZone(asset_prices_eod_ft.date::Timestamp, 'UTC':Varchar), 'UTC':Varchar), 'UTC':Varchar)::Date as $expr2, asset_prices_eod_ft.date], output_watermarks: [[asset_prices_eod_ft.date]] } { output: [ asset_prices_eod_ft.asset_id, $expr2, asset_prices_eod_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.reference_price ], stream key: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date ] }

Fragment 19476 (Actor 164477,164476)
StreamFilter { predicate: Not(IsNull(asset_prices_eod_ft.reference_price)) } { output: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date, asset_prices_eod_ft.reference_price ], stream key: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date ] }
└── StreamTableScan { table: asset_prices_eod_ft, columns: [asset_id, date, reference_price] } { output: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date, asset_prices_eod_ft.reference_price ], stream key: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date ] }
    ├── Upstream { output: [ asset_id, date, reference_price ], stream key: [] }
    └── BatchPlanNode { output: [ asset_id, date, reference_price ], stream key: [] }

Fragment 19477 (Actor 164462,164463)
StreamLocalityProvider { locality_columns: [asset_prices_eod_ft.asset_id, $expr3] } { output: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date, asset_prices_eod_ft.reference_price, $expr3 ], stream key: [ asset_prices_eod_ft.asset_id, $expr3, asset_prices_eod_ft.date ] }
└── MergeExecutor { output: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date, asset_prices_eod_ft.reference_price, $expr3 ], stream key: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date ] }

Fragment 19478 (Actor 164464,164465)
StreamProject { exprs: [asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date, asset_prices_eod_ft.reference_price, AtTimeZone(DateTrunc('MONTH':Varchar, AtTimeZone(asset_prices_eod_ft.date::Timestamp, 'UTC':Varchar), 'UTC':Varchar), 'UTC':Varchar)::Date as $expr3], output_watermarks: [[asset_prices_eod_ft.date]] } { output: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date, asset_prices_eod_ft.reference_price, $expr3 ], 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.reference_price ], stream key: [ asset_prices_eod_ft.asset_id, asset_prices_eod_ft.date ] }

Fragment 19479 (Actor 164468,164469)
StreamTableScan { table: settled_cost_basis_series_mv, columns: [account_id, asset_id, currency_code, effective_from, effective_to, average_cost_per_unit, average_cost_per_unit_system_currency, total_cost_system_currency, cost_fx_provenance, settled_cost_basis_carried_mv.sum] }
├── output: [ settled_cost_basis_series_mv.account_id, settled_cost_basis_series_mv.asset_id, settled_cost_basis_series_mv.currency_code, settled_cost_basis_series_mv.effective_from, settled_cost_basis_series_mv.effective_to, settled_cost_basis_series_mv.average_cost_per_unit, settled_cost_basis_series_mv.average_cost_per_unit_system_currency, settled_cost_basis_series_mv.total_cost_system_currency, settled_cost_basis_series_mv.cost_fx_provenance, settled_cost_basis_series_mv.settled_cost_basis_carried_mv.sum ]
├── stream key: [ settled_cost_basis_series_mv.account_id, settled_cost_basis_series_mv.asset_id, settled_cost_basis_series_mv.currency_code, settled_cost_basis_series_mv.settled_cost_basis_carried_mv.sum, settled_cost_basis_series_mv.effective_from ]
├── Upstream { output: [ account_id, asset_id, currency_code, effective_from, effective_to, average_cost_per_unit, average_cost_per_unit_system_currency, total_cost_system_currency, cost_fx_provenance, settled_cost_basis_carried_mv.sum ], stream key: [] }
└── BatchPlanNode { output: [ account_id, asset_id, currency_code, effective_from, effective_to, average_cost_per_unit, average_cost_per_unit_system_currency, total_cost_system_currency, cost_fx_provenance, settled_cost_basis_carried_mv.sum ], stream key: [] }