RWM Console cluster: risingwave.platform.svc.cluster.local

← cluster authz objects team_to_tasks_mv explain
Overview Objects Graph History
materialized view · authz.team_to_tasks_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 lookups
70 operators
Materialize · authz.team_to_tasks_mv
0% idle 2 actors
Project
2 actors
GroupTopN
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 · Inner · tasks_dm.resource_party_id = team_to_parties_mv.party_id
2 actors
HashJoin · Inner · tasks_dm.resource_party_id = team_to_parties_mv.party_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
StreamScan · team_to_parties_mv
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · tasks_dm
2 actors
Filter · tasks_dm
0% idle 2 actors
StreamScan · tasks_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · Inner · tasks_dm.resource_portfolio_id = team_to_portfolios_mv_next…
2 actors
HashJoin · Inner · tasks_dm.resource_portfolio_id = team_to_portfolios_mv_next… Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · team_to_portfolios_mv_next
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · tasks_dm
2 actors
Filter · tasks_dm
0% idle 2 actors
StreamScan · tasks_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · Inner · tasks_dm.resource_account_id = team_to_accounts_mv_next.acc…
2 actors
HashJoin · Inner · tasks_dm.resource_account_id = team_to_accounts_mv_next.acc… Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · team_to_accounts_mv_next
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · tasks_dm
2 actors
Filter · tasks_dm
0% idle 2 actors
StreamScan · tasks_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · Inner · tasks_dm.resource_client_id = team_to_clients_mv.client_id
2 actors
HashJoin · Inner · tasks_dm.resource_client_id = team_to_clients_mv.client_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
StreamScan · team_to_clients_mv
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · tasks_dm
2 actors
Filter · tasks_dm
0% idle 2 actors
StreamScan · tasks_dm
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 · authz.team_to_tasks_mv Materialize authz.team_to_tasks_mv idle · 2 actors Project Project — · 2 actors GroupTopN GroupTopN 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 · Inner · tasks_dm.resource_party_id = team_to_parties_mv.party_id Project Inner · tasks_dm.resour… — · 2 actors HashJoin · Inner · tasks_dm.resource_party_id = team_to_parties_mv.party_id HashJoin Inner · tasks_dm.resour… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · team_to_parties_mv StreamScan team_to_parties_mv idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · tasks_dm Project tasks_dm — · 2 actors Filter · tasks_dm Filter tasks_dm idle · 2 actors StreamScan · tasks_dm StreamScan tasks_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · Inner · tasks_dm.resource_portfolio_id = team_to_portfolios_mv_next… Project Inner · tasks_dm.resour… — · 2 actors HashJoin · Inner · tasks_dm.resource_portfolio_id = team_to_portfolios_mv_next… HashJoin Inner · tasks_dm.resour… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · team_to_portfolios_mv_next StreamScan team_to_portfolios_mv_n… idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · tasks_dm Project tasks_dm — · 2 actors Filter · tasks_dm Filter tasks_dm idle · 2 actors StreamScan · tasks_dm StreamScan tasks_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · Inner · tasks_dm.resource_account_id = team_to_accounts_mv_next.acc… Project Inner · tasks_dm.resour… — · 2 actors HashJoin · Inner · tasks_dm.resource_account_id = team_to_accounts_mv_next.acc… HashJoin Inner · tasks_dm.resour… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · team_to_accounts_mv_next StreamScan team_to_accounts_mv_next idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · tasks_dm Project tasks_dm — · 2 actors Filter · tasks_dm Filter tasks_dm idle · 2 actors StreamScan · tasks_dm StreamScan tasks_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · Inner · tasks_dm.resource_client_id = team_to_clients_mv.client_id Project Inner · tasks_dm.resour… — · 2 actors HashJoin · Inner · tasks_dm.resource_client_id = team_to_clients_mv.client_id HashJoin Inner · tasks_dm.resour… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · team_to_clients_mv StreamScan team_to_clients_mv idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · tasks_dm Project tasks_dm — · 2 actors Filter · tasks_dm Filter tasks_dm idle · 2 actors StreamScan · tasks_dm StreamScan tasks_dm 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 12337 (Actor 33581,33580)
StreamMaterialize { columns: [team_id, task_id, updated_at], stream_key: [team_id, task_id], pk_columns: [team_id, task_id], pk_conflict: NoCheck }
├── output: [ team_to_clients_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at ]
├── stream key: [ team_to_clients_mv.team_id, tasks_dm.task_id ]
└── StreamProject { exprs: [team_to_clients_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at] }
    ├── output: [ team_to_clients_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at ]
    ├── stream key: [ team_to_clients_mv.team_id, tasks_dm.task_id ]
    └── StreamGroupTopN { order: [tasks_dm.updated_at DESC], limit: 1, offset: 0, group_key: [team_to_clients_mv.team_id, tasks_dm.task_id] }
        ├── output:
        │   ┌── team_to_clients_mv.team_id
        │   ├── tasks_dm.task_id
        │   ├── tasks_dm.updated_at
        │   ├── tasks_dm.task_id
        │   ├── team_to_clients_mv.team_id
        │   ├── tasks_dm.resource_client_id
        │   └── $src
        ├── stream key: [ team_to_clients_mv.team_id, tasks_dm.task_id ]
        └── MergeExecutor
            ├── output:
            │   ┌── team_to_clients_mv.team_id
            │   ├── tasks_dm.task_id
            │   ├── tasks_dm.updated_at
            │   ├── tasks_dm.task_id
            │   ├── team_to_clients_mv.team_id
            │   ├── tasks_dm.resource_client_id
            │   └── $src
            └── stream key: [ tasks_dm.task_id, team_to_clients_mv.team_id, tasks_dm.resource_client_id, $src ]

Fragment 12338 (Actor 33583,33582)
StreamUnion { all: true }
├── output:
│   ┌── team_to_clients_mv.team_id
│   ├── tasks_dm.task_id
│   ├── tasks_dm.updated_at
│   ├── tasks_dm.task_id
│   ├── team_to_clients_mv.team_id
│   ├── tasks_dm.resource_client_id
│   └── $src
├── stream key: [ tasks_dm.task_id, team_to_clients_mv.team_id, tasks_dm.resource_client_id, $src ]
├── MergeExecutor
│   ├── output:
│   │   ┌── team_to_clients_mv.team_id
│   │   ├── tasks_dm.task_id
│   │   ├── tasks_dm.updated_at
│   │   ├── tasks_dm.task_id
│   │   ├── team_to_clients_mv.team_id
│   │   ├── tasks_dm.resource_client_id
│   │   └── 0:Int32
│   └── stream key: [ tasks_dm.task_id, team_to_clients_mv.team_id, tasks_dm.resource_client_id ]
├── MergeExecutor
│   ├── output:
│   │   ┌── team_to_accounts_mv_next.team_id
│   │   ├── tasks_dm.task_id
│   │   ├── tasks_dm.updated_at
│   │   ├── tasks_dm.task_id
│   │   ├── team_to_accounts_mv_next.team_id
│   │   ├── tasks_dm.resource_account_id
│   │   └── 1:Int32
│   └── stream key: [ tasks_dm.task_id, team_to_accounts_mv_next.team_id, tasks_dm.resource_account_id ]
├── MergeExecutor
│   ├── output:
│   │   ┌── team_to_portfolios_mv_next.team_id
│   │   ├── tasks_dm.task_id
│   │   ├── tasks_dm.updated_at
│   │   ├── tasks_dm.task_id
│   │   ├── team_to_portfolios_mv_next.team_id
│   │   ├── tasks_dm.resource_portfolio_id
│   │   └── 2:Int32
│   └── stream key: [ tasks_dm.task_id, team_to_portfolios_mv_next.team_id, tasks_dm.resource_portfolio_id ]
└── MergeExecutor
    ├── output:
    │   ┌── team_to_parties_mv.team_id
    │   ├── tasks_dm.task_id
    │   ├── tasks_dm.updated_at
    │   ├── tasks_dm.task_id
    │   ├── team_to_parties_mv.team_id
    │   ├── tasks_dm.resource_party_id
    │   └── 3:Int32
    └── stream key: [ tasks_dm.task_id, team_to_parties_mv.team_id, tasks_dm.resource_party_id ]

Fragment 12339 (Actor 33584,33585)
StreamProject { exprs: [team_to_clients_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.task_id, team_to_clients_mv.team_id, tasks_dm.resource_client_id, 0:Int32] }
├── output: [ team_to_clients_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.task_id, team_to_clients_mv.team_id, tasks_dm.resource_client_id, 0:Int32 ]
├── stream key: [ tasks_dm.task_id, team_to_clients_mv.team_id, tasks_dm.resource_client_id ]
└── StreamHashJoin { type: Inner, predicate: tasks_dm.resource_client_id = team_to_clients_mv.client_id }
    ├── output: [ team_to_clients_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_client_id, team_to_clients_mv.client_id ]
    ├── stream key: [ tasks_dm.task_id, team_to_clients_mv.team_id, tasks_dm.resource_client_id ]
    ├── MergeExecutor { output: [ tasks_dm.task_id, tasks_dm.resource_client_id, tasks_dm.updated_at ], stream key: [ tasks_dm.task_id ] }
    └── MergeExecutor { output: [ team_to_clients_mv.team_id, team_to_clients_mv.client_id ], stream key: [ team_to_clients_mv.team_id, team_to_clients_mv.client_id ] }

Fragment 12340 (Actor 33567,33566)
StreamProject { exprs: [tasks_dm.task_id, tasks_dm.resource_client_id, tasks_dm.updated_at] }
├── output: [ tasks_dm.task_id, tasks_dm.resource_client_id, tasks_dm.updated_at ]
├── stream key: [ tasks_dm.task_id ]
└── StreamFilter { predicate: IsNull(tasks_dm.disabled_at) }
    ├── output: [ tasks_dm.task_id, tasks_dm.resource_client_id, tasks_dm.updated_at, tasks_dm.disabled_at ]
    ├── stream key: [ tasks_dm.task_id ]
    └── StreamTableScan { table: tasks_dm, columns: [task_id, resource_client_id, updated_at, disabled_at] }
        ├── output: [ tasks_dm.task_id, tasks_dm.resource_client_id, tasks_dm.updated_at, tasks_dm.disabled_at ]
        ├── stream key: [ tasks_dm.task_id ]
        ├── Upstream { output: [ task_id, resource_client_id, updated_at, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ task_id, resource_client_id, updated_at, disabled_at ], stream key: [] }

Fragment 12341 (Actor 33611,33610)
StreamTableScan { table: team_to_clients_mv, columns: [team_id, client_id] }
├── output: [ team_to_clients_mv.team_id, team_to_clients_mv.client_id ]
├── stream key: [ team_to_clients_mv.team_id, team_to_clients_mv.client_id ]
├── Upstream { output: [ team_id, client_id ], stream key: [] }
└── BatchPlanNode { output: [ team_id, client_id ], stream key: [] }

Fragment 12342 (Actor 33609,33608)
StreamProject { exprs: [team_to_accounts_mv_next.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.task_id, team_to_accounts_mv_next.team_id, tasks_dm.resource_account_id, 1:Int32] }
├── output: [ team_to_accounts_mv_next.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.task_id, team_to_accounts_mv_next.team_id, tasks_dm.resource_account_id, 1:Int32 ]
├── stream key: [ tasks_dm.task_id, team_to_accounts_mv_next.team_id, tasks_dm.resource_account_id ]
└── StreamHashJoin { type: Inner, predicate: tasks_dm.resource_account_id = team_to_accounts_mv_next.account_id }
    ├── output: [ team_to_accounts_mv_next.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_account_id, team_to_accounts_mv_next.account_id ]
    ├── stream key: [ tasks_dm.task_id, team_to_accounts_mv_next.team_id, tasks_dm.resource_account_id ]
    ├── MergeExecutor { output: [ tasks_dm.task_id, tasks_dm.resource_account_id, tasks_dm.updated_at ], stream key: [ tasks_dm.task_id ] }
    └── MergeExecutor
        ├── output: [ team_to_accounts_mv_next.team_id, team_to_accounts_mv_next.account_id ]
        └── stream key: [ team_to_accounts_mv_next.team_id, team_to_accounts_mv_next.account_id ]

Fragment 12343 (Actor 33568,33569)
StreamProject { exprs: [tasks_dm.task_id, tasks_dm.resource_account_id, tasks_dm.updated_at] }
├── output: [ tasks_dm.task_id, tasks_dm.resource_account_id, tasks_dm.updated_at ]
├── stream key: [ tasks_dm.task_id ]
└── StreamFilter { predicate: IsNull(tasks_dm.disabled_at) }
    ├── output: [ tasks_dm.task_id, tasks_dm.resource_account_id, tasks_dm.updated_at, tasks_dm.disabled_at ]
    ├── stream key: [ tasks_dm.task_id ]
    └── StreamTableScan { table: tasks_dm, columns: [task_id, resource_account_id, updated_at, disabled_at] }
        ├── output: [ tasks_dm.task_id, tasks_dm.resource_account_id, tasks_dm.updated_at, tasks_dm.disabled_at ]
        ├── stream key: [ tasks_dm.task_id ]
        ├── Upstream { output: [ task_id, resource_account_id, updated_at, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ task_id, resource_account_id, updated_at, disabled_at ], stream key: [] }

Fragment 12344 (Actor 33621,33620)
StreamTableScan { table: team_to_accounts_mv_next, columns: [team_id, account_id] }
├── output: [ team_to_accounts_mv_next.team_id, team_to_accounts_mv_next.account_id ]
├── stream key: [ team_to_accounts_mv_next.team_id, team_to_accounts_mv_next.account_id ]
├── Upstream { output: [ team_id, account_id ], stream key: [] }
└── BatchPlanNode { output: [ team_id, account_id ], stream key: [] }

Fragment 12345 (Actor 33612,33613)
StreamProject { exprs: [team_to_portfolios_mv_next.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.task_id, team_to_portfolios_mv_next.team_id, tasks_dm.resource_portfolio_id, 2:Int32] }
├── output: [ team_to_portfolios_mv_next.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.task_id, team_to_portfolios_mv_next.team_id, tasks_dm.resource_portfolio_id, 2:Int32 ]
├── stream key: [ tasks_dm.task_id, team_to_portfolios_mv_next.team_id, tasks_dm.resource_portfolio_id ]
└── StreamHashJoin { type: Inner, predicate: tasks_dm.resource_portfolio_id = team_to_portfolios_mv_next.portfolio_id }
    ├── output: [ team_to_portfolios_mv_next.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_portfolio_id, team_to_portfolios_mv_next.portfolio_id ]
    ├── stream key: [ tasks_dm.task_id, team_to_portfolios_mv_next.team_id, tasks_dm.resource_portfolio_id ]
    ├── MergeExecutor { output: [ tasks_dm.task_id, tasks_dm.resource_portfolio_id, tasks_dm.updated_at ], stream key: [ tasks_dm.task_id ] }
    └── MergeExecutor
        ├── output: [ team_to_portfolios_mv_next.team_id, team_to_portfolios_mv_next.portfolio_id ]
        └── stream key: [ team_to_portfolios_mv_next.team_id, team_to_portfolios_mv_next.portfolio_id ]

Fragment 12346 (Actor 33622,33623)
StreamProject { exprs: [tasks_dm.task_id, tasks_dm.resource_portfolio_id, tasks_dm.updated_at] }
├── output: [ tasks_dm.task_id, tasks_dm.resource_portfolio_id, tasks_dm.updated_at ]
├── stream key: [ tasks_dm.task_id ]
└── StreamFilter { predicate: IsNull(tasks_dm.disabled_at) }
    ├── output: [ tasks_dm.task_id, tasks_dm.resource_portfolio_id, tasks_dm.updated_at, tasks_dm.disabled_at ]
    ├── stream key: [ tasks_dm.task_id ]
    └── StreamTableScan { table: tasks_dm, columns: [task_id, resource_portfolio_id, updated_at, disabled_at] }
        ├── output: [ tasks_dm.task_id, tasks_dm.resource_portfolio_id, tasks_dm.updated_at, tasks_dm.disabled_at ]
        ├── stream key: [ tasks_dm.task_id ]
        ├── Upstream { output: [ task_id, resource_portfolio_id, updated_at, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ task_id, resource_portfolio_id, updated_at, disabled_at ], stream key: [] }

Fragment 12347 (Actor 33625,33624)
StreamTableScan { table: team_to_portfolios_mv_next, columns: [team_id, portfolio_id] }
├── output: [ team_to_portfolios_mv_next.team_id, team_to_portfolios_mv_next.portfolio_id ]
├── stream key: [ team_to_portfolios_mv_next.team_id, team_to_portfolios_mv_next.portfolio_id ]
├── Upstream { output: [ team_id, portfolio_id ], stream key: [] }
└── BatchPlanNode { output: [ team_id, portfolio_id ], stream key: [] }

Fragment 12348 (Actor 33615,33614)
StreamProject { exprs: [team_to_parties_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.task_id, team_to_parties_mv.team_id, tasks_dm.resource_party_id, 3:Int32] }
├── output: [ team_to_parties_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.task_id, team_to_parties_mv.team_id, tasks_dm.resource_party_id, 3:Int32 ]
├── stream key: [ tasks_dm.task_id, team_to_parties_mv.team_id, tasks_dm.resource_party_id ]
└── StreamHashJoin { type: Inner, predicate: tasks_dm.resource_party_id = team_to_parties_mv.party_id }
    ├── output: [ team_to_parties_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_party_id, team_to_parties_mv.party_id ]
    ├── stream key: [ tasks_dm.task_id, team_to_parties_mv.team_id, tasks_dm.resource_party_id ]
    ├── MergeExecutor { output: [ tasks_dm.task_id, tasks_dm.resource_party_id, tasks_dm.updated_at ], stream key: [ tasks_dm.task_id ] }
    └── MergeExecutor { output: [ team_to_parties_mv.team_id, team_to_parties_mv.party_id ], stream key: [ team_to_parties_mv.team_id, team_to_parties_mv.party_id ] }

Fragment 12349 (Actor 33617,33616)
StreamProject { exprs: [tasks_dm.task_id, tasks_dm.resource_party_id, tasks_dm.updated_at] }
├── output: [ tasks_dm.task_id, tasks_dm.resource_party_id, tasks_dm.updated_at ]
├── stream key: [ tasks_dm.task_id ]
└── StreamFilter { predicate: IsNull(tasks_dm.disabled_at) }
    ├── output: [ tasks_dm.task_id, tasks_dm.resource_party_id, tasks_dm.updated_at, tasks_dm.disabled_at ]
    ├── stream key: [ tasks_dm.task_id ]
    └── StreamTableScan { table: tasks_dm, columns: [task_id, resource_party_id, updated_at, disabled_at] }
        ├── output: [ tasks_dm.task_id, tasks_dm.resource_party_id, tasks_dm.updated_at, tasks_dm.disabled_at ]
        ├── stream key: [ tasks_dm.task_id ]
        ├── Upstream { output: [ task_id, resource_party_id, updated_at, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ task_id, resource_party_id, updated_at, disabled_at ], stream key: [] }

Fragment 12350 (Actor 33627,33626)
StreamTableScan { table: team_to_parties_mv, columns: [team_id, party_id] }
├── output: [ team_to_parties_mv.team_id, team_to_parties_mv.party_id ]
├── stream key: [ team_to_parties_mv.team_id, team_to_parties_mv.party_id ]
├── Upstream { output: [ team_id, party_id ], stream key: [] }
└── BatchPlanNode { output: [ team_id, party_id ], stream key: [] }