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

← cluster adib_rm objects client_top_allocations_mv explain
Overview Objects Graph History
materialized view · adib_rm.client_top_allocations_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
24 operators
Materialize · adib_rm.client_top_allocations_mv
0% idle 2 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
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · client_to_account_groups_mv.account_group_id = intraday_pos…
2 actors
HashJoin · Inner · client_to_account_groups_mv.account_group_id = intraday_pos… Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
OverWindow · Inner · client_to_account_groups_mv.account_group_id = intraday_pos… Window state — add a WHERE rank <= N to bound it
0% idle 2 actors
GroupTopN · Inner · client_to_account_groups_mv.account_group_id = intraday_pos…
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · intraday_position_by_distribution_mv_next
2 actors
Filter · intraday_position_by_distribution_mv_next
0% idle 2 actors
StreamScan · intraday_position_by_distribution_mv_next
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · client_to_account_groups_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 · adib_rm.client_top_allocations_mv Materialize adib_rm.client_top_allo… idle · 2 actors Project Project — · 2 actors HashAgg HashAgg 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 SyncLogStore · Inner · client_to_account_groups_mv.account_group_id = intraday_pos… SyncLogStore Inner · client_to_accou… — · 2 actors HashJoin · Inner · client_to_account_groups_mv.account_group_id = intraday_pos… HashJoin Inner · client_to_accou… idle · 2 actors OverWindow · Inner · client_to_account_groups_mv.account_group_id = intraday_pos… OverWindow Inner · client_to_accou… idle · 2 actors GroupTopN · Inner · client_to_account_groups_mv.account_group_id = intraday_pos… GroupTopN Inner · client_to_accou… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · intraday_position_by_distribution_mv_next Project intraday_position_by_di… — · 2 actors Filter · intraday_position_by_distribution_mv_next Filter intraday_position_by_di… idle · 2 actors StreamScan · intraday_position_by_distribution_mv_next StreamScan intraday_position_by_di… idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · client_to_account_groups_mv StreamScan client_to_account_group… 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 25952 (Actor 131068,131067)
StreamMaterialize { columns: [client_id, account_group_type, top_allocations], stream_key: [client_id, account_group_type], pk_columns: [client_id, account_group_type], pk_conflict: NoCheck }
├── output: [ client_to_account_groups_mv.client_id, client_to_account_groups_mv.type, jsonb_agg($expr1 order_by(row_number ASC)) ]
├── stream key: [ client_to_account_groups_mv.client_id, client_to_account_groups_mv.type ]
└── StreamProject { exprs: [client_to_account_groups_mv.client_id, client_to_account_groups_mv.type, jsonb_agg($expr1 order_by(row_number ASC))] }
    ├── output: [ client_to_account_groups_mv.client_id, client_to_account_groups_mv.type, jsonb_agg($expr1 order_by(row_number ASC)) ]
    ├── stream key: [ client_to_account_groups_mv.client_id, client_to_account_groups_mv.type ]
    └── StreamHashAgg { group_key: [client_to_account_groups_mv.client_id, client_to_account_groups_mv.type], aggs: [jsonb_agg($expr1 order_by(row_number ASC)), count] }
        ├── output: [ client_to_account_groups_mv.client_id, client_to_account_groups_mv.type, jsonb_agg($expr1 order_by(row_number ASC)), count ]
        ├── stream key: [ client_to_account_groups_mv.client_id, client_to_account_groups_mv.type ]
        └── MergeExecutor
            ├── output:
            │   ┌── client_to_account_groups_mv.client_id
            │   ├── client_to_account_groups_mv.type
            │   ├── $expr1
            │   ├── row_number
            │   ├── client_to_account_groups_mv.$src
            │   ├── intraday_position_by_distribution_mv_next.currency_code
            │   ├── intraday_position_by_distribution_mv_next.position_type
            │   ├── intraday_position_by_distribution_mv_next.distribution_type
            │   ├── intraday_position_by_distribution_mv_next.taxonomy_node_id
            │   └── client_to_account_groups_mv.account_group_id
            └── stream key:
                ┌── client_to_account_groups_mv.client_id
                ├── client_to_account_groups_mv.$src
                ├── intraday_position_by_distribution_mv_next.currency_code
                ├── intraday_position_by_distribution_mv_next.position_type
                ├── intraday_position_by_distribution_mv_next.distribution_type
                ├── intraday_position_by_distribution_mv_next.taxonomy_node_id
                └── client_to_account_groups_mv.account_group_id

Fragment 25953 (Actor 131072,131071)
StreamProject { exprs: [client_to_account_groups_mv.client_id, client_to_account_groups_mv.type, JsonbBuildObject('assetClass':Varchar, intraday_position_by_distribution_mv_next.taxonomy_node_id, 'taxonomyNodeId':Varchar, intraday_position_by_distribution_mv_next.taxonomy_node_id, 'marketValue':Varchar, JsonbBuildObject('amount':Varchar, intraday_position_by_distribution_mv_next.market_value::Varchar, 'currencyCode':Varchar, intraday_position_by_distribution_mv_next.currency_code), 'fairValue':Varchar, JsonbBuildObject('amount':Varchar, intraday_position_by_distribution_mv_next.fair_value::Varchar, 'currencyCode':Varchar, intraday_position_by_distribution_mv_next.currency_code), 'weight':Varchar, intraday_position_by_distribution_mv_next.weight, 'rank':Varchar, row_number) as $expr1, row_number, client_to_account_groups_mv.$src, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type, intraday_position_by_distribution_mv_next.taxonomy_node_id, client_to_account_groups_mv.account_group_id] }
├── output: [ client_to_account_groups_mv.client_id, client_to_account_groups_mv.type, $expr1, row_number, client_to_account_groups_mv.$src, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type, intraday_position_by_distribution_mv_next.taxonomy_node_id, client_to_account_groups_mv.account_group_id ]
├── stream key: [ client_to_account_groups_mv.client_id, client_to_account_groups_mv.$src, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type, intraday_position_by_distribution_mv_next.taxonomy_node_id, client_to_account_groups_mv.account_group_id ]
└── MergeExecutor { output: [ client_to_account_groups_mv.client_id, client_to_account_groups_mv.type, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.market_value, intraday_position_by_distribution_mv_next.fair_value, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.weight, row_number, client_to_account_groups_mv.$src, client_to_account_groups_mv.account_group_id, intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type ], stream key: [ client_to_account_groups_mv.client_id, client_to_account_groups_mv.$src, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type, intraday_position_by_distribution_mv_next.taxonomy_node_id, client_to_account_groups_mv.account_group_id ] }

Fragment 25954 (Actor 131069,131070)
StreamSyncLogStore { output: [ client_to_account_groups_mv.client_id, client_to_account_groups_mv.type, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.market_value, intraday_position_by_distribution_mv_next.fair_value, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.weight, row_number, client_to_account_groups_mv.$src, client_to_account_groups_mv.account_group_id, intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type ], stream key: [ client_to_account_groups_mv.client_id, client_to_account_groups_mv.$src, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type, intraday_position_by_distribution_mv_next.taxonomy_node_id, client_to_account_groups_mv.account_group_id ] }
└── StreamHashJoin { type: Inner, predicate: client_to_account_groups_mv.account_group_id = intraday_position_by_distribution_mv_next.account_group_id }
    ├── output: [ client_to_account_groups_mv.client_id, client_to_account_groups_mv.type, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.market_value, intraday_position_by_distribution_mv_next.fair_value, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.weight, row_number, client_to_account_groups_mv.$src, client_to_account_groups_mv.account_group_id, intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type ]
    ├── stream key: [ client_to_account_groups_mv.client_id, client_to_account_groups_mv.$src, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type, intraday_position_by_distribution_mv_next.taxonomy_node_id, client_to_account_groups_mv.account_group_id ]
    ├── MergeExecutor { output: [ client_to_account_groups_mv.client_id, client_to_account_groups_mv.account_group_id, client_to_account_groups_mv.type, client_to_account_groups_mv.$src ], stream key: [ client_to_account_groups_mv.client_id, client_to_account_groups_mv.$src ] }
    └── StreamOverWindow { window_functions: [row_number() OVER(PARTITION BY intraday_position_by_distribution_mv_next.account_group_id ORDER BY intraday_position_by_distribution_mv_next.market_value DESC, intraday_position_by_distribution_mv_next.taxonomy_node_id ASC ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW)] }
        ├── output: [ intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.market_value, intraday_position_by_distribution_mv_next.fair_value, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.weight, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type, row_number ]
        ├── stream key: [ intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type, intraday_position_by_distribution_mv_next.taxonomy_node_id ]
        └── StreamGroupTopN { order: [intraday_position_by_distribution_mv_next.market_value DESC, intraday_position_by_distribution_mv_next.taxonomy_node_id ASC], limit: 5, offset: 0, group_key: [intraday_position_by_distribution_mv_next.account_group_id] } { output: [ intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.market_value, intraday_position_by_distribution_mv_next.fair_value, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.weight, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type ], stream key: [ intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type, intraday_position_by_distribution_mv_next.taxonomy_node_id ] }
            └── MergeExecutor { output: [ intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.market_value, intraday_position_by_distribution_mv_next.fair_value, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.weight, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type ], stream key: [ intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type, intraday_position_by_distribution_mv_next.taxonomy_node_id ] }

Fragment 25955 (Actor 131074,131073)
StreamTableScan { table: client_to_account_groups_mv, columns: [client_id, account_group_id, type, $src] } { output: [ client_to_account_groups_mv.client_id, client_to_account_groups_mv.account_group_id, client_to_account_groups_mv.type, client_to_account_groups_mv.$src ], stream key: [ client_to_account_groups_mv.client_id, client_to_account_groups_mv.$src ] }
├── Upstream { output: [ client_id, account_group_id, type, $src ], stream key: [] }
└── BatchPlanNode { output: [ client_id, account_group_id, type, $src ], stream key: [] }

Fragment 25956 (Actor 131075,131076)
StreamProject { exprs: [intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.market_value, intraday_position_by_distribution_mv_next.fair_value, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.weight, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type] }
├── output: [ intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.market_value, intraday_position_by_distribution_mv_next.fair_value, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.weight, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type ]
├── stream key: [ intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type, intraday_position_by_distribution_mv_next.taxonomy_node_id ]
└── StreamFilter { predicate: (intraday_position_by_distribution_mv_next.distribution_type = 'asset_classes':Varchar) AND (intraday_position_by_distribution_mv_next.position_type = 'POSITION':Varchar) } { output: [ intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.market_value, intraday_position_by_distribution_mv_next.fair_value, intraday_position_by_distribution_mv_next.weight, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type ], stream key: [ intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type, intraday_position_by_distribution_mv_next.taxonomy_node_id ] }
    └── StreamTableScan { table: intraday_position_by_distribution_mv_next, columns: [account_group_id, taxonomy_node_id, currency_code, market_value, fair_value, weight, position_type, distribution_type] } { output: [ intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.market_value, intraday_position_by_distribution_mv_next.fair_value, intraday_position_by_distribution_mv_next.weight, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type ], stream key: [ intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type, intraday_position_by_distribution_mv_next.taxonomy_node_id ] }
        ├── Upstream { output: [ account_group_id, taxonomy_node_id, currency_code, market_value, fair_value, weight, position_type, distribution_type ], stream key: [] }
        └── BatchPlanNode { output: [ account_group_id, taxonomy_node_id, currency_code, market_value, fair_value, weight, position_type, distribution_type ], stream key: [] }