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

← cluster adib_rm objects party_top_allocations_mv explain
Overview Objects Graph History
materialized view · adib_rm.party_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.party_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 · party_to_account_groups_mv_next.account_group_id = intraday…
2 actors
HashJoin · Inner · party_to_account_groups_mv_next.account_group_id = intraday… Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
OverWindow · Inner · party_to_account_groups_mv_next.account_group_id = intraday… Window state — add a WHERE rank <= N to bound it
0% idle 2 actors
GroupTopN · Inner · party_to_account_groups_mv_next.account_group_id = intraday…
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 · party_to_account_groups_mv_next
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.party_top_allocations_mv Materialize adib_rm.party_top_alloc… 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 · party_to_account_groups_mv_next.account_group_id = intraday… SyncLogStore Inner · party_to_accoun… — · 2 actors HashJoin · Inner · party_to_account_groups_mv_next.account_group_id = intraday… HashJoin Inner · party_to_accoun… idle · 2 actors OverWindow · Inner · party_to_account_groups_mv_next.account_group_id = intraday… OverWindow Inner · party_to_accoun… idle · 2 actors GroupTopN · Inner · party_to_account_groups_mv_next.account_group_id = intraday… GroupTopN Inner · party_to_accoun… 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 · party_to_account_groups_mv_next StreamScan party_to_account_groups… 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 25940 (Actor 130901,130902)
StreamMaterialize { columns: [party_id, account_group_type, top_allocations], stream_key: [party_id, account_group_type], pk_columns: [party_id, account_group_type], pk_conflict: NoCheck }
├── output: [ party_to_account_groups_mv_next.party_id, party_to_account_groups_mv_next.type, jsonb_agg($expr1 order_by(row_number ASC)) ]
├── stream key: [ party_to_account_groups_mv_next.party_id, party_to_account_groups_mv_next.type ]
└── StreamProject { exprs: [party_to_account_groups_mv_next.party_id, party_to_account_groups_mv_next.type, jsonb_agg($expr1 order_by(row_number ASC))] }
    ├── output: [ party_to_account_groups_mv_next.party_id, party_to_account_groups_mv_next.type, jsonb_agg($expr1 order_by(row_number ASC)) ]
    ├── stream key: [ party_to_account_groups_mv_next.party_id, party_to_account_groups_mv_next.type ]
    └── StreamHashAgg { group_key: [party_to_account_groups_mv_next.party_id, party_to_account_groups_mv_next.type], aggs: [jsonb_agg($expr1 order_by(row_number ASC)), count] }
        ├── output: [ party_to_account_groups_mv_next.party_id, party_to_account_groups_mv_next.type, jsonb_agg($expr1 order_by(row_number ASC)), count ]
        ├── stream key: [ party_to_account_groups_mv_next.party_id, party_to_account_groups_mv_next.type ]
        └── MergeExecutor
            ├── output:
            │   ┌── party_to_account_groups_mv_next.party_id
            │   ├── party_to_account_groups_mv_next.type
            │   ├── $expr1
            │   ├── row_number
            │   ├── party_to_account_groups_mv_next.parties.id
            │   ├── party_to_account_groups_mv_next.null:Varchar
            │   ├── party_to_account_groups_mv_next.$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
            │   └── party_to_account_groups_mv_next.account_group_id
            └── stream key:
                ┌── party_to_account_groups_mv_next.parties.id
                ├── party_to_account_groups_mv_next.null:Varchar
                ├── party_to_account_groups_mv_next.$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
                └── party_to_account_groups_mv_next.account_group_id

Fragment 25941 (Actor 130903,130904)
StreamProject { exprs: [party_to_account_groups_mv_next.party_id, party_to_account_groups_mv_next.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, party_to_account_groups_mv_next.parties.id, party_to_account_groups_mv_next.null:Varchar, party_to_account_groups_mv_next.$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, party_to_account_groups_mv_next.account_group_id] }
├── output: [ party_to_account_groups_mv_next.party_id, party_to_account_groups_mv_next.type, $expr1, row_number, party_to_account_groups_mv_next.parties.id, party_to_account_groups_mv_next.null:Varchar, party_to_account_groups_mv_next.$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, party_to_account_groups_mv_next.account_group_id ]
├── stream key: [ party_to_account_groups_mv_next.parties.id, party_to_account_groups_mv_next.null:Varchar, party_to_account_groups_mv_next.$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, party_to_account_groups_mv_next.account_group_id ]
└── MergeExecutor
    ├── output: [ party_to_account_groups_mv_next.party_id, party_to_account_groups_mv_next.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, party_to_account_groups_mv_next.parties.id, party_to_account_groups_mv_next.null:Varchar, party_to_account_groups_mv_next.$src, party_to_account_groups_mv_next.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: [ party_to_account_groups_mv_next.parties.id, party_to_account_groups_mv_next.null:Varchar, party_to_account_groups_mv_next.$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, party_to_account_groups_mv_next.account_group_id ]

Fragment 25942 (Actor 130905,130906)
StreamSyncLogStore
├── output: [ party_to_account_groups_mv_next.party_id, party_to_account_groups_mv_next.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, party_to_account_groups_mv_next.parties.id, party_to_account_groups_mv_next.null:Varchar, party_to_account_groups_mv_next.$src, party_to_account_groups_mv_next.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: [ party_to_account_groups_mv_next.parties.id, party_to_account_groups_mv_next.null:Varchar, party_to_account_groups_mv_next.$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, party_to_account_groups_mv_next.account_group_id ]
└── StreamHashJoin { type: Inner, predicate: party_to_account_groups_mv_next.account_group_id = intraday_position_by_distribution_mv_next.account_group_id }
    ├── output: [ party_to_account_groups_mv_next.party_id, party_to_account_groups_mv_next.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, party_to_account_groups_mv_next.parties.id, party_to_account_groups_mv_next.null:Varchar, party_to_account_groups_mv_next.$src, party_to_account_groups_mv_next.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: [ party_to_account_groups_mv_next.parties.id, party_to_account_groups_mv_next.null:Varchar, party_to_account_groups_mv_next.$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, party_to_account_groups_mv_next.account_group_id ]
    ├── MergeExecutor { output: [ party_to_account_groups_mv_next.party_id, party_to_account_groups_mv_next.account_group_id, party_to_account_groups_mv_next.type, party_to_account_groups_mv_next.parties.id, party_to_account_groups_mv_next.null:Varchar, party_to_account_groups_mv_next.$src ], stream key: [ party_to_account_groups_mv_next.parties.id, party_to_account_groups_mv_next.null:Varchar, party_to_account_groups_mv_next.$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 25943 (Actor 130908,130907)
StreamTableScan { table: party_to_account_groups_mv_next, columns: [party_id, account_group_id, type, parties.id, null:Varchar, $src] } { output: [ party_to_account_groups_mv_next.party_id, party_to_account_groups_mv_next.account_group_id, party_to_account_groups_mv_next.type, party_to_account_groups_mv_next.parties.id, party_to_account_groups_mv_next.null:Varchar, party_to_account_groups_mv_next.$src ], stream key: [ party_to_account_groups_mv_next.parties.id, party_to_account_groups_mv_next.null:Varchar, party_to_account_groups_mv_next.$src ] }
├── Upstream { output: [ party_id, account_group_id, type, parties.id, null:Varchar, $src ], stream key: [] }
└── BatchPlanNode { output: [ party_id, account_group_id, type, parties.id, null:Varchar, $src ], stream key: [] }

Fragment 25944 (Actor 130900,130899)
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: [] }