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

← cluster alinma_bff objects advisor_kpi_counts_mv explain
Overview Objects Graph History
materialized view · alinma_bff.advisor_kpi_counts_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 filtered
121 operators
Materialize · alinma_bff.advisor_kpi_counts_mv
0% idle 2 actors
Project
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · LeftOuter · advisor_kpi_user_clients_mv.user_id = advisor_kpi_user_acco…
2 actors
HashJoin · LeftOuter · advisor_kpi_user_clients_mv.user_id = advisor_kpi_user_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
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
SyncLogStore · Inner · advisor_kpi_user_accounts_mv_next.account_id = accounts_dm.…
2 actors
HashJoin · Inner · advisor_kpi_user_accounts_mv_next.account_id = accounts_dm.… 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
LocalityProvider
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
StreamScan · advisor_kpi_user_accounts_mv_next
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
SyncLogStore · LeftOuter · advisor_kpi_user_clients_mv.user_id = advisor_kpi_user_port…
2 actors
HashJoin · LeftOuter · advisor_kpi_user_clients_mv.user_id = advisor_kpi_user_port… 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
SyncLogStore · Inner · advisor_kpi_user_portfolios_mv_next.portfolio_id = portfoli…
2 actors
HashJoin · Inner · advisor_kpi_user_portfolios_mv_next.portfolio_id = portfoli… 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 · portfolios_dm
2 actors
Filter · portfolios_dm
0% idle 2 actors
StreamScan · portfolios_dm
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
NoOp
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · advisor_kpi_user_portfolios_mv_next
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
SyncLogStore · LeftOuter · advisor_kpi_user_clients_mv.user_id = advisor_kpi_user_clie…
2 actors
HashJoin · LeftOuter · advisor_kpi_user_clients_mv.user_id = advisor_kpi_user_clie… 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
SyncLogStore · Inner · advisor_kpi_user_clients_mv.client_id = clients_dm.id
2 actors
HashJoin · Inner · advisor_kpi_user_clients_mv.client_id = clients_dm.id 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 · clients_dm
2 actors
Filter · clients_dm
0% idle 2 actors
StreamScan · clients_dm
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
NoOp
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · advisor_kpi_user_clients_mv
0% idle 2 actors
BatchPlan
2 actors
Merge
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
Union
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
Project
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 · alinma_bff.advisor_kpi_counts_mv Materialize alinma_bff.advisor_kpi_… idle · 2 actors Project Project — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · LeftOuter · advisor_kpi_user_clients_mv.user_id = advisor_kpi_user_acco… SyncLogStore LeftOuter · advisor_kpi… — · 2 actors HashJoin · LeftOuter · advisor_kpi_user_clients_mv.user_id = advisor_kpi_user_acco… HashJoin LeftOuter · advisor_kpi… 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 SyncLogStore · Inner · advisor_kpi_user_accounts_mv_next.account_id = accounts_dm.… SyncLogStore Inner · advisor_kpi_use… — · 2 actors HashJoin · Inner · advisor_kpi_user_accounts_mv_next.account_id = accounts_dm.… HashJoin Inner · advisor_kpi_use… 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 LocalityProvider LocalityProvider 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 StreamScan · advisor_kpi_user_accounts_mv_next StreamScan advisor_kpi_user_accoun… 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 SyncLogStore · LeftOuter · advisor_kpi_user_clients_mv.user_id = advisor_kpi_user_port… SyncLogStore LeftOuter · advisor_kpi… — · 2 actors HashJoin · LeftOuter · advisor_kpi_user_clients_mv.user_id = advisor_kpi_user_port… HashJoin LeftOuter · advisor_kpi… 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 SyncLogStore · Inner · advisor_kpi_user_portfolios_mv_next.portfolio_id = portfoli… SyncLogStore Inner · advisor_kpi_use… — · 2 actors HashJoin · Inner · advisor_kpi_user_portfolios_mv_next.portfolio_id = portfoli… HashJoin Inner · advisor_kpi_use… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · portfolios_dm Project portfolios_dm — · 2 actors Filter · portfolios_dm Filter portfolios_dm idle · 2 actors StreamScan · portfolios_dm StreamScan portfolios_dm 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 NoOp NoOp idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · advisor_kpi_user_portfolios_mv_next StreamScan advisor_kpi_user_portfo… 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 SyncLogStore · LeftOuter · advisor_kpi_user_clients_mv.user_id = advisor_kpi_user_clie… SyncLogStore LeftOuter · advisor_kpi… — · 2 actors HashJoin · LeftOuter · advisor_kpi_user_clients_mv.user_id = advisor_kpi_user_clie… HashJoin LeftOuter · advisor_kpi… 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 SyncLogStore · Inner · advisor_kpi_user_clients_mv.client_id = clients_dm.id SyncLogStore Inner · advisor_kpi_use… — · 2 actors HashJoin · Inner · advisor_kpi_user_clients_mv.client_id = clients_dm.id HashJoin Inner · advisor_kpi_use… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · clients_dm Project clients_dm — · 2 actors Filter · clients_dm Filter clients_dm idle · 2 actors StreamScan · clients_dm StreamScan clients_dm 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 NoOp NoOp idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · advisor_kpi_user_clients_mv StreamScan advisor_kpi_user_client… idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 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 Union Union 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 Project Project — · 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 63765 (Actor 746593,746594)
StreamMaterialize { columns: [user_id, clients_count, portfolios_count, accounts_count], stream_key: [user_id], pk_columns: [user_id], pk_conflict: NoCheck }
├── output: [ advisor_kpi_user_clients_mv.user_id, $expr4, $expr5, $expr6 ]
├── stream key: [ advisor_kpi_user_clients_mv.user_id ]
└── StreamProject { exprs: [advisor_kpi_user_clients_mv.user_id, Coalesce($expr1, 0:Int32) as $expr4, Coalesce($expr2, 0:Int32) as $expr5, Coalesce($expr3, 0:Int32) as $expr6] }
    ├── output: [ advisor_kpi_user_clients_mv.user_id, $expr4, $expr5, $expr6 ]
    ├── stream key: [ advisor_kpi_user_clients_mv.user_id ]
    └── MergeExecutor
        ├── output: [ advisor_kpi_user_clients_mv.user_id, $expr1, $expr2, $expr3, advisor_kpi_user_accounts_mv_next.user_id ]
        └── stream key: [ advisor_kpi_user_clients_mv.user_id ]

Fragment 63766 (Actor 746595,746596)
StreamSyncLogStore
├── output: [ advisor_kpi_user_clients_mv.user_id, $expr1, $expr2, $expr3, advisor_kpi_user_accounts_mv_next.user_id ]
├── stream key: [ advisor_kpi_user_clients_mv.user_id ]
└── StreamHashJoin { type: LeftOuter, predicate: advisor_kpi_user_clients_mv.user_id = advisor_kpi_user_accounts_mv_next.user_id }
    ├── output: [ advisor_kpi_user_clients_mv.user_id, $expr1, $expr2, $expr3, advisor_kpi_user_accounts_mv_next.user_id ]
    ├── stream key: [ advisor_kpi_user_clients_mv.user_id ]
    ├── MergeExecutor
    │   ├── output: [ advisor_kpi_user_clients_mv.user_id, $expr1, $expr2, advisor_kpi_user_portfolios_mv_next.user_id ]
    │   └── stream key: [ advisor_kpi_user_clients_mv.user_id ]
    └── MergeExecutor { output: [ advisor_kpi_user_accounts_mv_next.user_id, $expr3 ], stream key: [ advisor_kpi_user_accounts_mv_next.user_id ] }

Fragment 63767 (Actor 746600,746599)
StreamLocalityProvider { locality_columns: [advisor_kpi_user_clients_mv.user_id] }
├── output: [ advisor_kpi_user_clients_mv.user_id, $expr1, $expr2, advisor_kpi_user_portfolios_mv_next.user_id ]
├── stream key: [ advisor_kpi_user_clients_mv.user_id ]
└── MergeExecutor
    ├── output: [ advisor_kpi_user_clients_mv.user_id, $expr1, $expr2, advisor_kpi_user_portfolios_mv_next.user_id ]
    └── stream key: [ advisor_kpi_user_clients_mv.user_id ]

Fragment 63768 (Actor 746598,746597)
StreamSyncLogStore
├── output: [ advisor_kpi_user_clients_mv.user_id, $expr1, $expr2, advisor_kpi_user_portfolios_mv_next.user_id ]
├── stream key: [ advisor_kpi_user_clients_mv.user_id ]
└── StreamHashJoin { type: LeftOuter, predicate: advisor_kpi_user_clients_mv.user_id = advisor_kpi_user_portfolios_mv_next.user_id }
    ├── output: [ advisor_kpi_user_clients_mv.user_id, $expr1, $expr2, advisor_kpi_user_portfolios_mv_next.user_id ]
    ├── stream key: [ advisor_kpi_user_clients_mv.user_id ]
    ├── MergeExecutor { output: [ advisor_kpi_user_clients_mv.user_id, $expr1, advisor_kpi_user_clients_mv.user_id ], stream key: [ advisor_kpi_user_clients_mv.user_id ] }
    └── MergeExecutor { output: [ advisor_kpi_user_portfolios_mv_next.user_id, $expr2 ], stream key: [ advisor_kpi_user_portfolios_mv_next.user_id ] }

Fragment 63769 (Actor 744624,744625)
StreamLocalityProvider { locality_columns: [advisor_kpi_user_clients_mv.user_id] }
├── output: [ advisor_kpi_user_clients_mv.user_id, $expr1, advisor_kpi_user_clients_mv.user_id ]
├── stream key: [ advisor_kpi_user_clients_mv.user_id ]
└── MergeExecutor { output: [ advisor_kpi_user_clients_mv.user_id, $expr1, advisor_kpi_user_clients_mv.user_id ], stream key: [ advisor_kpi_user_clients_mv.user_id ] }

Fragment 63770 (Actor 744627,744626)
StreamSyncLogStore { output: [ advisor_kpi_user_clients_mv.user_id, $expr1, advisor_kpi_user_clients_mv.user_id ], stream key: [ advisor_kpi_user_clients_mv.user_id ] }
└── StreamHashJoin { type: LeftOuter, predicate: advisor_kpi_user_clients_mv.user_id = advisor_kpi_user_clients_mv.user_id }
    ├── output: [ advisor_kpi_user_clients_mv.user_id, $expr1, advisor_kpi_user_clients_mv.user_id ]
    ├── stream key: [ advisor_kpi_user_clients_mv.user_id ]
    ├── MergeExecutor { output: [ advisor_kpi_user_clients_mv.user_id ], stream key: [ advisor_kpi_user_clients_mv.user_id ] }
    └── MergeExecutor { output: [ advisor_kpi_user_clients_mv.user_id, $expr1 ], stream key: [ advisor_kpi_user_clients_mv.user_id ] }

Fragment 63771 (Actor 744630,744631)
StreamProject { exprs: [advisor_kpi_user_clients_mv.user_id] } { output: [ advisor_kpi_user_clients_mv.user_id ], stream key: [ advisor_kpi_user_clients_mv.user_id ] }
└── StreamHashAgg { group_key: [advisor_kpi_user_clients_mv.user_id], aggs: [count] }
    ├── output: [ advisor_kpi_user_clients_mv.user_id, count ]
    ├── stream key: [ advisor_kpi_user_clients_mv.user_id ]
    └── StreamLocalityProvider { locality_columns: [advisor_kpi_user_clients_mv.user_id] }
        ├── output: [ advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.client_id, $src ]
        ├── stream key: [ advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.client_id, $src ]
        └── MergeExecutor
            ├── output: [ advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.client_id, $src ]
            └── stream key: [ advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.client_id, $src ]

Fragment 63772 (Actor 744676,744677)
StreamUnion { all: true }
├── output: [ advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.client_id, $src ]
├── stream key: [ advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.client_id, $src ]
├── MergeExecutor
│   ├── output: [ advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.client_id, 0:Int32 ]
│   └── stream key: [ advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.client_id ]
├── MergeExecutor
│   ├── output: [ advisor_kpi_user_portfolios_mv_next.user_id, advisor_kpi_user_portfolios_mv_next.user_id, advisor_kpi_user_portfolios_mv_next.portfolio_id, 1:Int32 ]
│   └── stream key: [ advisor_kpi_user_portfolios_mv_next.user_id, advisor_kpi_user_portfolios_mv_next.portfolio_id ]
└── MergeExecutor
    ├── output: [ advisor_kpi_user_accounts_mv_next.user_id, advisor_kpi_user_accounts_mv_next.user_id, advisor_kpi_user_accounts_mv_next.account_id, 2:Int32 ]
    └── stream key: [ advisor_kpi_user_accounts_mv_next.user_id, advisor_kpi_user_accounts_mv_next.account_id ]

Fragment 63773 (Actor 744353,744354)
StreamProject { exprs: [advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.client_id, 0:Int32] }
├── output: [ advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.client_id, 0:Int32 ]
├── stream key: [ advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.client_id ]
└── MergeExecutor
    ├── output: [ advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.client_id ]
    └── stream key: [ advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.client_id ]

Fragment 63774 (Actor 744356,744355)
StreamTableScan { table: advisor_kpi_user_clients_mv, columns: [user_id, client_id] }
├── output: [ advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.client_id ]
├── stream key: [ advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.client_id ]
├── Upstream { output: [ user_id, client_id ], stream key: [] }
└── BatchPlanNode { output: [ user_id, client_id ], stream key: [] }

Fragment 63775 (Actor 745870,745871)
StreamProject { exprs: [advisor_kpi_user_portfolios_mv_next.user_id, advisor_kpi_user_portfolios_mv_next.user_id, advisor_kpi_user_portfolios_mv_next.portfolio_id, 1:Int32] }
├── output: [ advisor_kpi_user_portfolios_mv_next.user_id, advisor_kpi_user_portfolios_mv_next.user_id, advisor_kpi_user_portfolios_mv_next.portfolio_id, 1:Int32 ]
├── stream key: [ advisor_kpi_user_portfolios_mv_next.user_id, advisor_kpi_user_portfolios_mv_next.portfolio_id ]
└── MergeExecutor
    ├── output: [ advisor_kpi_user_portfolios_mv_next.user_id, advisor_kpi_user_portfolios_mv_next.portfolio_id ]
    └── stream key: [ advisor_kpi_user_portfolios_mv_next.user_id, advisor_kpi_user_portfolios_mv_next.portfolio_id ]

Fragment 63776 (Actor 745868,745869)
StreamTableScan { table: advisor_kpi_user_portfolios_mv_next, columns: [user_id, portfolio_id] }
├── output: [ advisor_kpi_user_portfolios_mv_next.user_id, advisor_kpi_user_portfolios_mv_next.portfolio_id ]
├── stream key: [ advisor_kpi_user_portfolios_mv_next.user_id, advisor_kpi_user_portfolios_mv_next.portfolio_id ]
├── Upstream { output: [ user_id, portfolio_id ], stream key: [] }
└── BatchPlanNode { output: [ user_id, portfolio_id ], stream key: [] }

Fragment 63777 (Actor 745876,745877)
StreamProject { exprs: [advisor_kpi_user_accounts_mv_next.user_id, advisor_kpi_user_accounts_mv_next.user_id, advisor_kpi_user_accounts_mv_next.account_id, 2:Int32] }
├── output: [ advisor_kpi_user_accounts_mv_next.user_id, advisor_kpi_user_accounts_mv_next.user_id, advisor_kpi_user_accounts_mv_next.account_id, 2:Int32 ]
├── stream key: [ advisor_kpi_user_accounts_mv_next.user_id, advisor_kpi_user_accounts_mv_next.account_id ]
└── MergeExecutor
    ├── output: [ advisor_kpi_user_accounts_mv_next.user_id, advisor_kpi_user_accounts_mv_next.account_id ]
    └── stream key: [ advisor_kpi_user_accounts_mv_next.user_id, advisor_kpi_user_accounts_mv_next.account_id ]

Fragment 63778 (Actor 745874,745875)
StreamTableScan { table: advisor_kpi_user_accounts_mv_next, columns: [user_id, account_id] }
├── output: [ advisor_kpi_user_accounts_mv_next.user_id, advisor_kpi_user_accounts_mv_next.account_id ]
├── stream key: [ advisor_kpi_user_accounts_mv_next.user_id, advisor_kpi_user_accounts_mv_next.account_id ]
├── Upstream { output: [ user_id, account_id ], stream key: [] }
└── BatchPlanNode { output: [ user_id, account_id ], stream key: [] }

Fragment 63779 (Actor 745680,745681)
StreamProject { exprs: [advisor_kpi_user_clients_mv.user_id, count::Int32 as $expr1] }
├── output: [ advisor_kpi_user_clients_mv.user_id, $expr1 ]
├── stream key: [ advisor_kpi_user_clients_mv.user_id ]
└── StreamHashAgg { group_key: [advisor_kpi_user_clients_mv.user_id], aggs: [count] }
    ├── output: [ advisor_kpi_user_clients_mv.user_id, count ]
    ├── stream key: [ advisor_kpi_user_clients_mv.user_id ]
    └── StreamLocalityProvider { locality_columns: [advisor_kpi_user_clients_mv.user_id] }
        ├── output: [ advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.client_id, clients_dm.id ]
        ├── stream key: [ advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.client_id ]
        └── MergeExecutor
            ├── output: [ advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.client_id, clients_dm.id ]
            └── stream key: [ advisor_kpi_user_clients_mv.client_id, advisor_kpi_user_clients_mv.user_id ]

Fragment 63780 (Actor 745683,745682)
StreamSyncLogStore
├── output: [ advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.client_id, clients_dm.id ]
├── stream key: [ advisor_kpi_user_clients_mv.client_id, advisor_kpi_user_clients_mv.user_id ]
└── StreamHashJoin { type: Inner, predicate: advisor_kpi_user_clients_mv.client_id = clients_dm.id }
    ├── output: [ advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.client_id, clients_dm.id ]
    ├── stream key: [ advisor_kpi_user_clients_mv.client_id, advisor_kpi_user_clients_mv.user_id ]
    ├── MergeExecutor
    │   ├── output: [ advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.client_id ]
    │   └── stream key: [ advisor_kpi_user_clients_mv.client_id, advisor_kpi_user_clients_mv.user_id ]
    └── MergeExecutor { output: [ clients_dm.id ], stream key: [ clients_dm.id ] }

Fragment 63781 (Actor 745686,745687)
StreamLocalityProvider { locality_columns: [advisor_kpi_user_clients_mv.client_id] }
├── output: [ advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.client_id ]
├── stream key: [ advisor_kpi_user_clients_mv.client_id, advisor_kpi_user_clients_mv.user_id ]
└── MergeExecutor
    ├── output: [ advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.client_id ]
    └── stream key: [ advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.client_id ]

Fragment 63782 (Actor 744357,744358)
StreamNoOp
├── output: [ advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.client_id ]
├── stream key: [ advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.client_id ]
└── MergeExecutor
    ├── output: [ advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.client_id ]
    └── stream key: [ advisor_kpi_user_clients_mv.user_id, advisor_kpi_user_clients_mv.client_id ]

Fragment 63783 (Actor 745883,745882)
StreamProject { exprs: [clients_dm.id] } { output: [ clients_dm.id ], stream key: [ clients_dm.id ] }
└── StreamFilter { predicate: IsNull(clients_dm.closing_date) } { output: [ clients_dm.id, clients_dm.closing_date ], stream key: [ clients_dm.id ] }
    └── StreamTableScan { table: clients_dm, columns: [id, closing_date] } { output: [ clients_dm.id, clients_dm.closing_date ], stream key: [ clients_dm.id ] }
        ├── Upstream { output: [ id, closing_date ], stream key: [] }
        └── BatchPlanNode { output: [ id, closing_date ], stream key: [] }

Fragment 63784 (Actor 745772,745773)
StreamProject { exprs: [advisor_kpi_user_portfolios_mv_next.user_id, count::Int32 as $expr2] }
├── output: [ advisor_kpi_user_portfolios_mv_next.user_id, $expr2 ]
├── stream key: [ advisor_kpi_user_portfolios_mv_next.user_id ]
└── StreamHashAgg { group_key: [advisor_kpi_user_portfolios_mv_next.user_id], aggs: [count] }
    ├── output: [ advisor_kpi_user_portfolios_mv_next.user_id, count ]
    ├── stream key: [ advisor_kpi_user_portfolios_mv_next.user_id ]
    └── StreamLocalityProvider { locality_columns: [advisor_kpi_user_portfolios_mv_next.user_id] }
        ├── output: [ advisor_kpi_user_portfolios_mv_next.user_id, advisor_kpi_user_portfolios_mv_next.portfolio_id, portfolios_dm.portfolio_id ]
        ├── stream key: [ advisor_kpi_user_portfolios_mv_next.user_id, advisor_kpi_user_portfolios_mv_next.portfolio_id ]
        └── MergeExecutor
            ├── output: [ advisor_kpi_user_portfolios_mv_next.user_id, advisor_kpi_user_portfolios_mv_next.portfolio_id, portfolios_dm.portfolio_id ]
            └── stream key: [ advisor_kpi_user_portfolios_mv_next.portfolio_id, advisor_kpi_user_portfolios_mv_next.user_id ]

Fragment 63785 (Actor 745786,745787)
StreamSyncLogStore
├── output: [ advisor_kpi_user_portfolios_mv_next.user_id, advisor_kpi_user_portfolios_mv_next.portfolio_id, portfolios_dm.portfolio_id ]
├── stream key: [ advisor_kpi_user_portfolios_mv_next.portfolio_id, advisor_kpi_user_portfolios_mv_next.user_id ]
└── StreamHashJoin { type: Inner, predicate: advisor_kpi_user_portfolios_mv_next.portfolio_id = portfolios_dm.portfolio_id }
    ├── output: [ advisor_kpi_user_portfolios_mv_next.user_id, advisor_kpi_user_portfolios_mv_next.portfolio_id, portfolios_dm.portfolio_id ]
    ├── stream key: [ advisor_kpi_user_portfolios_mv_next.portfolio_id, advisor_kpi_user_portfolios_mv_next.user_id ]
    ├── MergeExecutor
    │   ├── output: [ advisor_kpi_user_portfolios_mv_next.user_id, advisor_kpi_user_portfolios_mv_next.portfolio_id ]
    │   └── stream key: [ advisor_kpi_user_portfolios_mv_next.portfolio_id, advisor_kpi_user_portfolios_mv_next.user_id ]
    └── MergeExecutor { output: [ portfolios_dm.portfolio_id ], stream key: [ portfolios_dm.portfolio_id ] }

Fragment 63786 (Actor 745796,745797)
StreamLocalityProvider { locality_columns: [advisor_kpi_user_portfolios_mv_next.portfolio_id] }
├── output: [ advisor_kpi_user_portfolios_mv_next.user_id, advisor_kpi_user_portfolios_mv_next.portfolio_id ]
├── stream key: [ advisor_kpi_user_portfolios_mv_next.portfolio_id, advisor_kpi_user_portfolios_mv_next.user_id ]
└── MergeExecutor
    ├── output: [ advisor_kpi_user_portfolios_mv_next.user_id, advisor_kpi_user_portfolios_mv_next.portfolio_id ]
    └── stream key: [ advisor_kpi_user_portfolios_mv_next.user_id, advisor_kpi_user_portfolios_mv_next.portfolio_id ]

Fragment 63787 (Actor 745867,745866)
StreamNoOp
├── output: [ advisor_kpi_user_portfolios_mv_next.user_id, advisor_kpi_user_portfolios_mv_next.portfolio_id ]
├── stream key: [ advisor_kpi_user_portfolios_mv_next.user_id, advisor_kpi_user_portfolios_mv_next.portfolio_id ]
└── MergeExecutor
    ├── output: [ advisor_kpi_user_portfolios_mv_next.user_id, advisor_kpi_user_portfolios_mv_next.portfolio_id ]
    └── stream key: [ advisor_kpi_user_portfolios_mv_next.user_id, advisor_kpi_user_portfolios_mv_next.portfolio_id ]

Fragment 63788 (Actor 746232,746231)
StreamProject { exprs: [portfolios_dm.portfolio_id] } { output: [ portfolios_dm.portfolio_id ], stream key: [ portfolios_dm.portfolio_id ] }
└── StreamFilter { predicate: IsNull(portfolios_dm.disabled_at) }
    ├── output: [ portfolios_dm.portfolio_id, portfolios_dm.disabled_at ]
    ├── stream key: [ portfolios_dm.portfolio_id ]
    └── StreamTableScan { table: portfolios_dm, columns: [portfolio_id, disabled_at] }
        ├── output: [ portfolios_dm.portfolio_id, portfolios_dm.disabled_at ]
        ├── stream key: [ portfolios_dm.portfolio_id ]
        ├── Upstream { output: [ portfolio_id, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ portfolio_id, disabled_at ], stream key: [] }

Fragment 63789 (Actor 745842,745843)
StreamProject { exprs: [advisor_kpi_user_accounts_mv_next.user_id, count::Int32 as $expr3] }
├── output: [ advisor_kpi_user_accounts_mv_next.user_id, $expr3 ]
├── stream key: [ advisor_kpi_user_accounts_mv_next.user_id ]
└── StreamHashAgg { group_key: [advisor_kpi_user_accounts_mv_next.user_id], aggs: [count] }
    ├── output: [ advisor_kpi_user_accounts_mv_next.user_id, count ]
    ├── stream key: [ advisor_kpi_user_accounts_mv_next.user_id ]
    └── StreamLocalityProvider { locality_columns: [advisor_kpi_user_accounts_mv_next.user_id] }
        ├── output: [ advisor_kpi_user_accounts_mv_next.user_id, advisor_kpi_user_accounts_mv_next.account_id, accounts_dm.account_id ]
        ├── stream key: [ advisor_kpi_user_accounts_mv_next.user_id, advisor_kpi_user_accounts_mv_next.account_id ]
        └── MergeExecutor
            ├── output: [ advisor_kpi_user_accounts_mv_next.user_id, advisor_kpi_user_accounts_mv_next.account_id, accounts_dm.account_id ]
            └── stream key: [ advisor_kpi_user_accounts_mv_next.account_id, advisor_kpi_user_accounts_mv_next.user_id ]

Fragment 63790 (Actor 745844,745845)
StreamSyncLogStore
├── output: [ advisor_kpi_user_accounts_mv_next.user_id, advisor_kpi_user_accounts_mv_next.account_id, accounts_dm.account_id ]
├── stream key: [ advisor_kpi_user_accounts_mv_next.account_id, advisor_kpi_user_accounts_mv_next.user_id ]
└── StreamHashJoin { type: Inner, predicate: advisor_kpi_user_accounts_mv_next.account_id = accounts_dm.account_id }
    ├── output: [ advisor_kpi_user_accounts_mv_next.user_id, advisor_kpi_user_accounts_mv_next.account_id, accounts_dm.account_id ]
    ├── stream key: [ advisor_kpi_user_accounts_mv_next.account_id, advisor_kpi_user_accounts_mv_next.user_id ]
    ├── MergeExecutor
    │   ├── output: [ advisor_kpi_user_accounts_mv_next.user_id, advisor_kpi_user_accounts_mv_next.account_id ]
    │   └── stream key: [ advisor_kpi_user_accounts_mv_next.account_id, advisor_kpi_user_accounts_mv_next.user_id ]
    └── MergeExecutor { output: [ accounts_dm.account_id ], stream key: [ accounts_dm.account_id ] }

Fragment 63791 (Actor 745847,745846)
StreamLocalityProvider { locality_columns: [advisor_kpi_user_accounts_mv_next.account_id] }
├── output: [ advisor_kpi_user_accounts_mv_next.user_id, advisor_kpi_user_accounts_mv_next.account_id ]
├── stream key: [ advisor_kpi_user_accounts_mv_next.account_id, advisor_kpi_user_accounts_mv_next.user_id ]
└── MergeExecutor
    ├── output: [ advisor_kpi_user_accounts_mv_next.user_id, advisor_kpi_user_accounts_mv_next.account_id ]
    └── stream key: [ advisor_kpi_user_accounts_mv_next.user_id, advisor_kpi_user_accounts_mv_next.account_id ]

Fragment 63792 (Actor 745873,745872)
StreamNoOp
├── output: [ advisor_kpi_user_accounts_mv_next.user_id, advisor_kpi_user_accounts_mv_next.account_id ]
├── stream key: [ advisor_kpi_user_accounts_mv_next.user_id, advisor_kpi_user_accounts_mv_next.account_id ]
└── MergeExecutor
    ├── output: [ advisor_kpi_user_accounts_mv_next.user_id, advisor_kpi_user_accounts_mv_next.account_id ]
    └── stream key: [ advisor_kpi_user_accounts_mv_next.user_id, advisor_kpi_user_accounts_mv_next.account_id ]

Fragment 63793 (Actor 745851,745850)
StreamProject { exprs: [accounts_dm.account_id] } { output: [ accounts_dm.account_id ], stream key: [ accounts_dm.account_id ] }
└── StreamFilter { predicate: IsNull(accounts_dm.closing_date) AND IsNull(accounts_dm.disabled_at) }
    ├── output: [ accounts_dm.account_id, accounts_dm.closing_date, accounts_dm.disabled_at ]
    ├── stream key: [ accounts_dm.account_id ]
    └── StreamTableScan { table: accounts_dm, columns: [account_id, closing_date, disabled_at] }
        ├── output: [ accounts_dm.account_id, accounts_dm.closing_date, accounts_dm.disabled_at ]
        ├── stream key: [ accounts_dm.account_id ]
        ├── Upstream { output: [ account_id, closing_date, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ account_id, closing_date, disabled_at ], stream key: [] }