RWM Console cluster: risingwave-adib.adib-rw.svc.cluster.local

← cluster adib_rm objects investment_account_balance_snapshot_journal_mv explain
Overview Objects Graph History
materialized view · adib_rm.investment_account_balance_snapshot_journal_mv profiled over 5s
seconds (1–30)
Stateful hash join (4 state tables) — consider a temporal join for dimension lookupsAggregation state — unbounded unless keyed or temporally filteredDynamic filter — verify it pairs with a temporal condition to clean state
73 operators
Materialize · adib_rm.investment_account_balance_snapshot_journal_mv
0% idle 2 actors
Project
2 actors
HashAgg Aggregation state — unbounded unless keyed or temporally filtered
0% idle 2 actors
Project
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · settled_cash_balances_mv_next.account_id = accounts_dm.acco…
2 actors
HashJoin · Inner · settled_cash_balances_mv_next.account_id = accounts_dm.acco… 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 · accounts_dm
2 actors
Filter · accounts_dm
0% idle 2 actors
StreamScan · accounts_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · LeftOuter · settled_cash_balances_mv_next.account_id = holding_values_i…
2 actors
HashJoin · LeftOuter · settled_cash_balances_mv_next.account_id = holding_values_i… Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · LeftOuter · settled_cash_balances_mv_next.account_id = settled_cash_bal…
2 actors
HashJoin · LeftOuter · settled_cash_balances_mv_next.account_id = settled_cash_bal… 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
NoOp
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · settled_cash_balances_mv_next
2 actors
GroupTopN · settled_cash_balances_mv_next
0% idle 2 actors
Project · settled_cash_balances_mv_next
2 actors
DynamicFilter · settled_cash_balances_mv_next Dynamic filter — verify it pairs with a temporal condition to clean state
0% idle 2 actors
Merge
2 actors
Exchange
0% 2/s 0 actors
Now
0% 2/s 1 actor
Project · settled_cash_balances_mv_next
2 actors
StreamScan · settled_cash_balances_mv_next
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Project · LeftOuter · settled_cash_balances_mv_next.account_id = settled_cash_bal…
2 actors
HashAgg · LeftOuter · settled_cash_balances_mv_next.account_id = settled_cash_bal… Aggregation state — unbounded unless keyed or temporally filtered
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Union
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
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
Merge
2 actors
Exchange
0% idle 0 actors
Project · Inner · holding_values_intraday_ft.asset_id = assets_dm_next.id AND…
2 actors
GroupTopN · Inner · holding_values_intraday_ft.asset_id = assets_dm_next.id AND…
0% idle 2 actors
TemporalJoin · Inner · holding_values_intraday_ft.asset_id = assets_dm_next.id AND…
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · assets_dm_next
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · holding_values_intraday_ft
2 actors
Filter · holding_values_intraday_ft
0% idle 2 actors
StreamScan · holding_values_intraday_ft
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
2 actors
Merge
2 actors
Heat = the operator's output-buffer backpressure over the sampling window. Click a node to fold its subtree.
Materialize · adib_rm.investment_account_balance_snapshot_journal_mv Materialize adib_rm.investment_acco… idle · 2 actors Project Project — · 2 actors HashAgg HashAgg idle · 2 actors Project Project — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · settled_cash_balances_mv_next.account_id = accounts_dm.acco… SyncLogStore Inner · settled_cash_ba… — · 2 actors HashJoin · Inner · settled_cash_balances_mv_next.account_id = accounts_dm.acco… HashJoin Inner · settled_cash_ba… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · accounts_dm Project accounts_dm — · 2 actors Filter · accounts_dm Filter accounts_dm idle · 2 actors StreamScan · accounts_dm StreamScan accounts_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · LeftOuter · settled_cash_balances_mv_next.account_id = holding_values_i… SyncLogStore LeftOuter · settled_cas… — · 2 actors HashJoin · LeftOuter · settled_cash_balances_mv_next.account_id = holding_values_i… HashJoin LeftOuter · settled_cas… idle · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · LeftOuter · settled_cash_balances_mv_next.account_id = settled_cash_bal… SyncLogStore LeftOuter · settled_cas… — · 2 actors HashJoin · LeftOuter · settled_cash_balances_mv_next.account_id = settled_cash_bal… HashJoin LeftOuter · settled_cas… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors NoOp NoOp idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · settled_cash_balances_mv_next Project settled_cash_balances_m… — · 2 actors GroupTopN · settled_cash_balances_mv_next GroupTopN settled_cash_balances_m… idle · 2 actors Project · settled_cash_balances_mv_next Project settled_cash_balances_m… — · 2 actors DynamicFilter · settled_cash_balances_mv_next DynamicFilter settled_cash_balances_m… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange 2/s · 0 actors Now Now 2/s · 1 actor Project · settled_cash_balances_mv_next Project settled_cash_balances_m… — · 2 actors StreamScan · settled_cash_balances_mv_next StreamScan settled_cash_balances_m… idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Project · LeftOuter · settled_cash_balances_mv_next.account_id = settled_cash_bal… Project LeftOuter · settled_cas… — · 2 actors HashAgg · LeftOuter · settled_cash_balances_mv_next.account_id = settled_cash_bal… HashAgg LeftOuter · settled_cas… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Union Union 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 Project Project — · 2 actors HashAgg HashAgg idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · Inner · holding_values_intraday_ft.asset_id = assets_dm_next.id AND… Project Inner · holding_values_… — · 2 actors GroupTopN · Inner · holding_values_intraday_ft.asset_id = assets_dm_next.id AND… GroupTopN Inner · holding_values_… idle · 2 actors TemporalJoin · Inner · holding_values_intraday_ft.asset_id = assets_dm_next.id AND… TemporalJoin Inner · holding_values_… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · assets_dm_next StreamScan assets_dm_next idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · holding_values_intraday_ft Project holding_values_intraday… — · 2 actors Filter · holding_values_intraday_ft Filter holding_values_intraday… idle · 2 actors StreamScan · holding_values_intraday_ft StreamScan holding_values_intraday… idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 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 25377 (Actor 124005,124004)
StreamMaterialize { columns: [account_id, available_balance, fact_date], stream_key: [account_id], pk_columns: [account_id], pk_conflict: NoCheck } { output: [ settled_cash_balances_mv_next.account_id, sum($expr3), max($expr4) ], stream key: [ settled_cash_balances_mv_next.account_id ] }
└── StreamProject { exprs: [settled_cash_balances_mv_next.account_id, sum($expr3), max($expr4)] } { output: [ settled_cash_balances_mv_next.account_id, sum($expr3), max($expr4) ], stream key: [ settled_cash_balances_mv_next.account_id ] }
    └── StreamHashAgg { group_key: [settled_cash_balances_mv_next.account_id], aggs: [sum($expr3), max($expr4), count] } { output: [ settled_cash_balances_mv_next.account_id, sum($expr3), max($expr4), count ], stream key: [ settled_cash_balances_mv_next.account_id ] }
        └── StreamProject { exprs: [settled_cash_balances_mv_next.account_id, Case((Not(IsNull(holding_values_intraday_ft.account_id)) AND (IsNull(settled_cash_balances_mv_next.account_id) OR (max($expr2) >= settled_cash_balances_mv_next.dim_settlement_date))), sum(holding_values_intraday_ft.market_value), settled_cash_balances_mv_next.settled_cash_balance) as $expr3, Case((Not(IsNull(holding_values_intraday_ft.account_id)) AND (IsNull(settled_cash_balances_mv_next.account_id) OR (max($expr2) >= settled_cash_balances_mv_next.dim_settlement_date))), max($expr2), settled_cash_balances_mv_next.dim_settlement_date) as $expr4, settled_cash_balances_mv_next.currency_code] }
            ├── output: [ settled_cash_balances_mv_next.account_id, $expr3, $expr4, settled_cash_balances_mv_next.currency_code ]
            ├── stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ]
            └── MergeExecutor { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.settled_cash_balance, settled_cash_balances_mv_next.dim_settlement_date, holding_values_intraday_ft.account_id, sum(holding_values_intraday_ft.market_value), max($expr2), accounts_dm.account_id, settled_cash_balances_mv_next.currency_code ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }

Fragment 25378 (Actor 124006,124007)
StreamSyncLogStore { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.settled_cash_balance, settled_cash_balances_mv_next.dim_settlement_date, holding_values_intraday_ft.account_id, sum(holding_values_intraday_ft.market_value), max($expr2), accounts_dm.account_id, settled_cash_balances_mv_next.currency_code ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
└── StreamHashJoin { type: Inner, predicate: settled_cash_balances_mv_next.account_id = accounts_dm.account_id } { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.settled_cash_balance, settled_cash_balances_mv_next.dim_settlement_date, holding_values_intraday_ft.account_id, sum(holding_values_intraday_ft.market_value), max($expr2), accounts_dm.account_id, settled_cash_balances_mv_next.currency_code ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
    ├── MergeExecutor { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.settled_cash_balance, settled_cash_balances_mv_next.dim_settlement_date, holding_values_intraday_ft.account_id, sum(holding_values_intraday_ft.market_value), max($expr2), settled_cash_balances_mv_next.currency_code, holding_values_intraday_ft.currency_code ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
    └── MergeExecutor { output: [ accounts_dm.account_id ], stream key: [ accounts_dm.account_id ] }

Fragment 25379 (Actor 124010,124011)
StreamSyncLogStore { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.settled_cash_balance, settled_cash_balances_mv_next.dim_settlement_date, holding_values_intraday_ft.account_id, sum(holding_values_intraday_ft.market_value), max($expr2), settled_cash_balances_mv_next.currency_code, holding_values_intraday_ft.currency_code ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
└── StreamHashJoin { type: LeftOuter, predicate: settled_cash_balances_mv_next.account_id = holding_values_intraday_ft.account_id AND settled_cash_balances_mv_next.currency_code = holding_values_intraday_ft.currency_code }
    ├── output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.settled_cash_balance, settled_cash_balances_mv_next.dim_settlement_date, holding_values_intraday_ft.account_id, sum(holding_values_intraday_ft.market_value), max($expr2), settled_cash_balances_mv_next.currency_code, holding_values_intraday_ft.currency_code ]
    ├── stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ]
    ├── MergeExecutor { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.settled_cash_balance, settled_cash_balances_mv_next.dim_settlement_date, settled_cash_balances_mv_next.currency_code ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
    └── MergeExecutor { output: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code, sum(holding_values_intraday_ft.market_value), max($expr2) ], stream key: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code ] }

Fragment 25380 (Actor 124013,124012)
StreamSyncLogStore { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.settled_cash_balance, settled_cash_balances_mv_next.dim_settlement_date, settled_cash_balances_mv_next.currency_code ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
└── StreamHashJoin { type: LeftOuter, predicate: settled_cash_balances_mv_next.account_id = settled_cash_balances_mv_next.account_id AND settled_cash_balances_mv_next.currency_code = settled_cash_balances_mv_next.currency_code } { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.settled_cash_balance, settled_cash_balances_mv_next.dim_settlement_date, settled_cash_balances_mv_next.currency_code ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
    ├── StreamProject { exprs: [settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code] } { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
    │   └── StreamHashAgg { group_key: [settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code], aggs: [count] } { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, count ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
    │       └── MergeExecutor { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, $src ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, $src ] }
    └── MergeExecutor { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.settled_cash_balance, settled_cash_balances_mv_next.dim_settlement_date ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }

Fragment 25381 (Actor 124016,124017)
StreamUnion { all: true } { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, $src ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, $src ] }
├── MergeExecutor { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, 0:Int32 ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
└── MergeExecutor { output: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code, 1:Int32 ], stream key: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code ] }

Fragment 25382 (Actor 124022,124023)
StreamProject { exprs: [settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, 0:Int32] } { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, 0:Int32 ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
└── MergeExecutor { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.settled_cash_balance, settled_cash_balances_mv_next.dim_settlement_date ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }

Fragment 25383 (Actor 124018,124019)
StreamProject { exprs: [settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.settled_cash_balance, settled_cash_balances_mv_next.dim_settlement_date] } { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.settled_cash_balance, settled_cash_balances_mv_next.dim_settlement_date ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
└── StreamGroupTopN { order: [settled_cash_balances_mv_next.dim_settlement_date DESC], limit: 1, offset: 0, group_key: [settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code] } { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.dim_settlement_date, settled_cash_balances_mv_next.settled_cash_balance ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
    └── StreamProject { exprs: [settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.dim_settlement_date, settled_cash_balances_mv_next.settled_cash_balance] } { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.dim_settlement_date, settled_cash_balances_mv_next.settled_cash_balance ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.dim_settlement_date ] }
        └── StreamDynamicFilter { predicate: ($expr1 <= now), output: [settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.dim_settlement_date, settled_cash_balances_mv_next.settled_cash_balance, $expr1], cleaned_by_watermark: true } { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.dim_settlement_date, settled_cash_balances_mv_next.settled_cash_balance, $expr1 ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.dim_settlement_date ] }
            ├── StreamProject { exprs: [settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.dim_settlement_date, settled_cash_balances_mv_next.settled_cash_balance, AtTimeZone(settled_cash_balances_mv_next.dim_settlement_date::Timestamp, 'UTC':Varchar) as $expr1] }
            │   ├── output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.dim_settlement_date, settled_cash_balances_mv_next.settled_cash_balance, $expr1 ]
            │   ├── stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.dim_settlement_date ]
            │   └── StreamTableScan { table: settled_cash_balances_mv_next, columns: [account_id, currency_code, dim_settlement_date, settled_cash_balance] } { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.dim_settlement_date, settled_cash_balances_mv_next.settled_cash_balance ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.dim_settlement_date ] }
            │       ├── Upstream { output: [ account_id, currency_code, dim_settlement_date, settled_cash_balance ], stream key: [] }
            │       └── BatchPlanNode { output: [ account_id, currency_code, dim_settlement_date, settled_cash_balance ], stream key: [] }
            └── MergeExecutor { output: [ now ], stream key: [] }

Fragment 25384 (Actor 124024)
StreamNow { output: [ now ], stream key: [] }

Fragment 25385 (Actor 124014,124015)
StreamProject { exprs: [holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code, 1:Int32] } { output: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code, 1:Int32 ], stream key: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code ] }
└── MergeExecutor { output: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code, sum(holding_values_intraday_ft.market_value), max($expr2) ], stream key: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code ] }

Fragment 25386 (Actor 124008,124009)
StreamProject { exprs: [holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code, sum(holding_values_intraday_ft.market_value), max($expr2)] } { output: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code, sum(holding_values_intraday_ft.market_value), max($expr2) ], stream key: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code ] }
└── StreamHashAgg { group_key: [holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code], aggs: [sum(holding_values_intraday_ft.market_value), max($expr2), count] } { output: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code, sum(holding_values_intraday_ft.market_value), max($expr2), count ], stream key: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code ] }
    └── MergeExecutor { output: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code, holding_values_intraday_ft.market_value, $expr2, holding_values_intraday_ft.asset_id ], stream key: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.asset_id ] }

Fragment 25387 (Actor 123929,123930)
StreamProject { exprs: [holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code, holding_values_intraday_ft.market_value, AtTimeZone(holding_values_intraday_ft.holding_timestamp, 'UTC':Varchar)::Date as $expr2, holding_values_intraday_ft.asset_id] } { output: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code, holding_values_intraday_ft.market_value, $expr2, holding_values_intraday_ft.asset_id ], 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.market_value, holding_values_intraday_ft.id, assets_dm_next.id ], stream key: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.asset_id ] }
    └── StreamTemporalJoin { type: Inner, append_only: false, predicate: holding_values_intraday_ft.asset_id = assets_dm_next.id AND (assets_dm_next.type = 'CASH':Varchar), nested_loop: false, 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.market_value, holding_values_intraday_ft.id, assets_dm_next.id ], stream key: [ holding_values_intraday_ft.id, holding_values_intraday_ft.holding_timestamp, 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.market_value, holding_values_intraday_ft.id ], stream key: [ holding_values_intraday_ft.id, holding_values_intraday_ft.holding_timestamp ] }
        └── MergeExecutor { output: [ assets_dm_next.id, assets_dm_next.type ], stream key: [ assets_dm_next.id ] }

Fragment 25388 (Actor 124026,124025)
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.market_value, 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.market_value, 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.market_value, 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, market_value, 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.market_value, 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, market_value, id, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ account_id, asset_id, currency_code, holding_timestamp, market_value, id, disabled_at ], stream key: [] }

Fragment 25389 (Actor 123928,123927)
StreamTableScan { table: assets_dm_next, columns: [id, type] } { output: [ assets_dm_next.id, assets_dm_next.type ], stream key: [ assets_dm_next.id ] }
├── Upstream { output: [ id, type ], stream key: [] }
└── BatchPlanNode { output: [ id, type ], stream key: [] }

Fragment 25390 (Actor 124020,124021)
StreamNoOp { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.settled_cash_balance, settled_cash_balances_mv_next.dim_settlement_date ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
└── MergeExecutor { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.settled_cash_balance, settled_cash_balances_mv_next.dim_settlement_date ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }

Fragment 25391 (Actor 124027,124028)
StreamProject { exprs: [accounts_dm.account_id] } { output: [ accounts_dm.account_id ], stream key: [ accounts_dm.account_id ] }
└── StreamFilter { predicate: IsNull(accounts_dm.disabled_at) } { output: [ accounts_dm.account_id, accounts_dm.disabled_at ], stream key: [ accounts_dm.account_id ] }
    └── StreamTableScan { table: accounts_dm, columns: [account_id, disabled_at] } { output: [ accounts_dm.account_id, accounts_dm.disabled_at ], stream key: [ accounts_dm.account_id ] }
        ├── Upstream { output: [ account_id, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ account_id, disabled_at ], stream key: [] }