RWM Console cluster: risingwave-alinma.alinma-rw.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
107 operators
Materialize · authz.team_to_tasks_mv
0% idle 2 actors
Project
2 actors
GroupTopN
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
Exchange
0% idle 0 actors
SyncLogStore · Inner · tasks_dm.resource_party_id = team_to_parties_mv_next.party_…
2 actors
HashJoin · Inner · tasks_dm.resource_party_id = team_to_parties_mv_next.party_… 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
LocalityProvider
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · team_to_parties_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
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
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · 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
LocalityProvider
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
LocalityProvider
0% idle 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
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · 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
LocalityProvider
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
LocalityProvider
0% idle 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
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · 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
LocalityProvider
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
LocalityProvider
0% idle 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 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 Exchange Exchange idle · 0 actors SyncLogStore · Inner · tasks_dm.resource_party_id = team_to_parties_mv_next.party_… SyncLogStore Inner · tasks_dm.resour… — · 2 actors HashJoin · Inner · tasks_dm.resource_party_id = team_to_parties_mv_next.party_… HashJoin Inner · tasks_dm.resour… idle · 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 StreamScan · team_to_parties_mv_next StreamScan team_to_parties_mv_next 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 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 Project — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · tasks_dm.resource_portfolio_id = team_to_portfolios_mv_next… SyncLogStore 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 LocalityProvider LocalityProvider 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 LocalityProvider LocalityProvider idle · 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 Project — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · tasks_dm.resource_account_id = team_to_accounts_mv_next.acc… SyncLogStore 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 LocalityProvider LocalityProvider 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 LocalityProvider LocalityProvider idle · 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 Project — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · tasks_dm.resource_client_id = team_to_clients_mv.client_id SyncLogStore 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 LocalityProvider LocalityProvider 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 LocalityProvider LocalityProvider idle · 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 63398 (Actor 746049,746048)
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.resource_client_id
        │   ├── tasks_dm.task_id
        │   ├── team_to_clients_mv.team_id
        │   └── $src
        ├── stream key: [ team_to_clients_mv.team_id, tasks_dm.task_id ]
        └── StreamLocalityProvider { locality_columns: [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.resource_client_id
            │   ├── tasks_dm.task_id
            │   ├── team_to_clients_mv.team_id
            │   └── $src
            ├── stream key:
            │   ┌── team_to_clients_mv.team_id
            │   ├── tasks_dm.task_id
            │   ├── tasks_dm.resource_client_id
            │   ├── tasks_dm.task_id
            │   ├── team_to_clients_mv.team_id
            │   └── $src
            └── MergeExecutor
                ├── output:
                │   ┌── team_to_clients_mv.team_id
                │   ├── tasks_dm.task_id
                │   ├── tasks_dm.updated_at
                │   ├── tasks_dm.resource_client_id
                │   ├── tasks_dm.task_id
                │   ├── team_to_clients_mv.team_id
                │   └── $src
                └── stream key: [ tasks_dm.resource_client_id, tasks_dm.task_id, team_to_clients_mv.team_id, $src ]

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

Fragment 63400 (Actor 746053,746052)
StreamProject { exprs: [team_to_clients_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_client_id, tasks_dm.task_id, team_to_clients_mv.team_id, 0:Int32] }
├── output: [ team_to_clients_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_client_id, tasks_dm.task_id, team_to_clients_mv.team_id, 0:Int32 ]
├── stream key: [ tasks_dm.resource_client_id, tasks_dm.task_id, team_to_clients_mv.team_id ]
└── MergeExecutor
    ├── 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.resource_client_id, tasks_dm.task_id, team_to_clients_mv.team_id ]

Fragment 63401 (Actor 746054,746055)
StreamSyncLogStore
├── 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.resource_client_id, tasks_dm.task_id, team_to_clients_mv.team_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.resource_client_id, tasks_dm.task_id, team_to_clients_mv.team_id ]
    ├── MergeExecutor { output: [ tasks_dm.task_id, tasks_dm.resource_client_id, tasks_dm.updated_at ], stream key: [ tasks_dm.resource_client_id, tasks_dm.task_id ] }
    └── MergeExecutor { output: [ team_to_clients_mv.team_id, team_to_clients_mv.client_id ], stream key: [ team_to_clients_mv.client_id, team_to_clients_mv.team_id ] }

Fragment 63402 (Actor 746119,746120)
StreamLocalityProvider { locality_columns: [tasks_dm.resource_client_id] }
├── output: [ tasks_dm.task_id, tasks_dm.resource_client_id, tasks_dm.updated_at ]
├── stream key: [ tasks_dm.resource_client_id, tasks_dm.task_id ]
└── MergeExecutor { output: [ tasks_dm.task_id, tasks_dm.resource_client_id, tasks_dm.updated_at ], stream key: [ tasks_dm.task_id ] }

Fragment 63403 (Actor 745948,745949)
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 63404 (Actor 746125,746126)
StreamLocalityProvider { locality_columns: [team_to_clients_mv.client_id] }
├── output: [ team_to_clients_mv.team_id, team_to_clients_mv.client_id ]
├── stream key: [ team_to_clients_mv.client_id, team_to_clients_mv.team_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 63405 (Actor 746163,746164)
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 63406 (Actor 746128,746127)
StreamProject { exprs: [team_to_accounts_mv_next.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_account_id, tasks_dm.task_id, team_to_accounts_mv_next.team_id, 1:Int32] }
├── output: [ team_to_accounts_mv_next.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_account_id, tasks_dm.task_id, team_to_accounts_mv_next.team_id, 1:Int32 ]
├── stream key: [ tasks_dm.resource_account_id, tasks_dm.task_id, team_to_accounts_mv_next.team_id ]
└── MergeExecutor
    ├── 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.resource_account_id, tasks_dm.task_id, team_to_accounts_mv_next.team_id ]

Fragment 63407 (Actor 746130,746129)
StreamSyncLogStore
├── 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.resource_account_id, tasks_dm.task_id, team_to_accounts_mv_next.team_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.resource_account_id, tasks_dm.task_id, team_to_accounts_mv_next.team_id ]
    ├── MergeExecutor { output: [ tasks_dm.task_id, tasks_dm.resource_account_id, tasks_dm.updated_at ], stream key: [ tasks_dm.resource_account_id, 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.account_id, team_to_accounts_mv_next.team_id ]

Fragment 63408 (Actor 746131,746132)
StreamLocalityProvider { locality_columns: [tasks_dm.resource_account_id] }
├── output: [ tasks_dm.task_id, tasks_dm.resource_account_id, tasks_dm.updated_at ]
├── stream key: [ tasks_dm.resource_account_id, tasks_dm.task_id ]
└── MergeExecutor { output: [ tasks_dm.task_id, tasks_dm.resource_account_id, tasks_dm.updated_at ], stream key: [ tasks_dm.task_id ] }

Fragment 63409 (Actor 746165,746166)
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 63410 (Actor 746133,746134)
StreamLocalityProvider { locality_columns: [team_to_accounts_mv_next.account_id] }
├── output: [ team_to_accounts_mv_next.team_id, team_to_accounts_mv_next.account_id ]
├── stream key: [ team_to_accounts_mv_next.account_id, team_to_accounts_mv_next.team_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 63411 (Actor 746167,746168)
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 63412 (Actor 746137,746138)
StreamProject { exprs: [team_to_portfolios_mv_next.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_portfolio_id, tasks_dm.task_id, team_to_portfolios_mv_next.team_id, 2:Int32] }
├── output: [ team_to_portfolios_mv_next.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_portfolio_id, tasks_dm.task_id, team_to_portfolios_mv_next.team_id, 2:Int32 ]
├── stream key: [ tasks_dm.resource_portfolio_id, tasks_dm.task_id, team_to_portfolios_mv_next.team_id ]
└── MergeExecutor
    ├── 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.resource_portfolio_id, tasks_dm.task_id, team_to_portfolios_mv_next.team_id ]

Fragment 63413 (Actor 746135,746136)
StreamSyncLogStore
├── 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.resource_portfolio_id, tasks_dm.task_id, team_to_portfolios_mv_next.team_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.resource_portfolio_id, tasks_dm.task_id, team_to_portfolios_mv_next.team_id ]
    ├── MergeExecutor { output: [ tasks_dm.task_id, tasks_dm.resource_portfolio_id, tasks_dm.updated_at ], stream key: [ tasks_dm.resource_portfolio_id, 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.portfolio_id, team_to_portfolios_mv_next.team_id ]

Fragment 63414 (Actor 746140,746139)
StreamLocalityProvider { locality_columns: [tasks_dm.resource_portfolio_id] }
├── output: [ tasks_dm.task_id, tasks_dm.resource_portfolio_id, tasks_dm.updated_at ]
├── stream key: [ tasks_dm.resource_portfolio_id, tasks_dm.task_id ]
└── MergeExecutor { output: [ tasks_dm.task_id, tasks_dm.resource_portfolio_id, tasks_dm.updated_at ], stream key: [ tasks_dm.task_id ] }

Fragment 63415 (Actor 746170,746169)
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 63416 (Actor 746142,746141)
StreamLocalityProvider { locality_columns: [team_to_portfolios_mv_next.portfolio_id] }
├── output: [ team_to_portfolios_mv_next.team_id, team_to_portfolios_mv_next.portfolio_id ]
├── stream key: [ team_to_portfolios_mv_next.portfolio_id, team_to_portfolios_mv_next.team_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 63417 (Actor 746172,746171)
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 63418 (Actor 746146,746145)
StreamProject { exprs: [team_to_parties_mv_next.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_party_id, tasks_dm.task_id, team_to_parties_mv_next.team_id, 3:Int32] }
├── output: [ team_to_parties_mv_next.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_party_id, tasks_dm.task_id, team_to_parties_mv_next.team_id, 3:Int32 ]
├── stream key: [ tasks_dm.resource_party_id, tasks_dm.task_id, team_to_parties_mv_next.team_id ]
└── MergeExecutor
    ├── output: [ team_to_parties_mv_next.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_party_id, team_to_parties_mv_next.party_id ]
    └── stream key: [ tasks_dm.resource_party_id, tasks_dm.task_id, team_to_parties_mv_next.team_id ]

Fragment 63419 (Actor 746143,746144)
StreamSyncLogStore
├── output: [ team_to_parties_mv_next.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_party_id, team_to_parties_mv_next.party_id ]
├── stream key: [ tasks_dm.resource_party_id, tasks_dm.task_id, team_to_parties_mv_next.team_id ]
└── StreamHashJoin { type: Inner, predicate: tasks_dm.resource_party_id = team_to_parties_mv_next.party_id }
    ├── output: [ team_to_parties_mv_next.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_party_id, team_to_parties_mv_next.party_id ]
    ├── stream key: [ tasks_dm.resource_party_id, tasks_dm.task_id, team_to_parties_mv_next.team_id ]
    ├── MergeExecutor { output: [ tasks_dm.task_id, tasks_dm.resource_party_id, tasks_dm.updated_at ], stream key: [ tasks_dm.resource_party_id, tasks_dm.task_id ] }
    └── MergeExecutor { output: [ team_to_parties_mv_next.team_id, team_to_parties_mv_next.party_id ], stream key: [ team_to_parties_mv_next.party_id, team_to_parties_mv_next.team_id ] }

Fragment 63420 (Actor 746148,746147)
StreamLocalityProvider { locality_columns: [tasks_dm.resource_party_id] }
├── output: [ tasks_dm.task_id, tasks_dm.resource_party_id, tasks_dm.updated_at ]
├── stream key: [ tasks_dm.resource_party_id, tasks_dm.task_id ]
└── MergeExecutor { output: [ tasks_dm.task_id, tasks_dm.resource_party_id, tasks_dm.updated_at ], stream key: [ tasks_dm.task_id ] }

Fragment 63421 (Actor 746173,746174)
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 63422 (Actor 746149,746150)
StreamLocalityProvider { locality_columns: [team_to_parties_mv_next.party_id] }
├── output: [ team_to_parties_mv_next.team_id, team_to_parties_mv_next.party_id ]
├── stream key: [ team_to_parties_mv_next.party_id, team_to_parties_mv_next.team_id ]
└── MergeExecutor { output: [ team_to_parties_mv_next.team_id, team_to_parties_mv_next.party_id ], stream key: [ team_to_parties_mv_next.team_id, team_to_parties_mv_next.party_id ] }

Fragment 63423 (Actor 746152,746151)
StreamTableScan { table: team_to_parties_mv_next, columns: [team_id, party_id] }
├── output: [ team_to_parties_mv_next.team_id, team_to_parties_mv_next.party_id ]
├── stream key: [ team_to_parties_mv_next.team_id, team_to_parties_mv_next.party_id ]
├── Upstream { output: [ team_id, party_id ], stream key: [] }
└── BatchPlanNode { output: [ team_id, party_id ], stream key: [] }