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.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
LocalityProvider
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
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.port…
2 actors
HashJoin · Inner · tasks_dm.resource_portfolio_id = team_to_portfolios_mv.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
LocalityProvider
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · team_to_portfolios_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
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.account_…
2 actors
HashJoin · Inner · tasks_dm.resource_account_id = team_to_accounts_mv.account_… 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
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.party_id SyncLogStore 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 LocalityProvider LocalityProvider 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 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.port… SyncLogStore Inner · tasks_dm.resour… — · 2 actors HashJoin · Inner · tasks_dm.resource_portfolio_id = team_to_portfolios_mv.port… 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 StreamScan team_to_portfolios_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 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.account_… SyncLogStore Inner · tasks_dm.resour… — · 2 actors HashJoin · Inner · tasks_dm.resource_account_id = team_to_accounts_mv.account_… 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 StreamScan team_to_accounts_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 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 52676 (Actor 738170,738169)
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 52677 (Actor 738171,738172)
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.team_id
│   │   ├── tasks_dm.task_id
│   │   ├── tasks_dm.updated_at
│   │   ├── tasks_dm.resource_account_id
│   │   ├── tasks_dm.task_id
│   │   ├── team_to_accounts_mv.team_id
│   │   └── 1:Int32
│   └── stream key: [ tasks_dm.resource_account_id, tasks_dm.task_id, team_to_accounts_mv.team_id ]
├── MergeExecutor
│   ├── output:
│   │   ┌── team_to_portfolios_mv.team_id
│   │   ├── tasks_dm.task_id
│   │   ├── tasks_dm.updated_at
│   │   ├── tasks_dm.resource_portfolio_id
│   │   ├── tasks_dm.task_id
│   │   ├── team_to_portfolios_mv.team_id
│   │   └── 2:Int32
│   └── stream key: [ tasks_dm.resource_portfolio_id, tasks_dm.task_id, team_to_portfolios_mv.team_id ]
└── MergeExecutor
    ├── output:
    │   ┌── team_to_parties_mv.team_id
    │   ├── tasks_dm.task_id
    │   ├── tasks_dm.updated_at
    │   ├── tasks_dm.resource_party_id
    │   ├── tasks_dm.task_id
    │   ├── team_to_parties_mv.team_id
    │   └── 3:Int32
    └── stream key: [ tasks_dm.resource_party_id, tasks_dm.task_id, team_to_parties_mv.team_id ]

Fragment 52678 (Actor 738174,738173)
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 52679 (Actor 738175,738176)
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 52680 (Actor 738178,738177)
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 52681 (Actor 738228,738227)
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 52682 (Actor 738179,738180)
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 52683 (Actor 738230,738229)
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 52684 (Actor 738185,738186)
StreamProject { exprs: [team_to_accounts_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_account_id, tasks_dm.task_id, team_to_accounts_mv.team_id, 1:Int32] }
├── output: [ team_to_accounts_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_account_id, tasks_dm.task_id, team_to_accounts_mv.team_id, 1:Int32 ]
├── stream key: [ tasks_dm.resource_account_id, tasks_dm.task_id, team_to_accounts_mv.team_id ]
└── MergeExecutor
    ├── output: [ team_to_accounts_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_account_id, team_to_accounts_mv.account_id ]
    └── stream key: [ tasks_dm.resource_account_id, tasks_dm.task_id, team_to_accounts_mv.team_id ]

Fragment 52685 (Actor 738184,738183)
StreamSyncLogStore
├── output: [ team_to_accounts_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_account_id, team_to_accounts_mv.account_id ]
├── stream key: [ tasks_dm.resource_account_id, tasks_dm.task_id, team_to_accounts_mv.team_id ]
└── StreamHashJoin { type: Inner, predicate: tasks_dm.resource_account_id = team_to_accounts_mv.account_id }
    ├── output: [ team_to_accounts_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_account_id, team_to_accounts_mv.account_id ]
    ├── stream key: [ tasks_dm.resource_account_id, tasks_dm.task_id, team_to_accounts_mv.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.team_id, team_to_accounts_mv.account_id ], stream key: [ team_to_accounts_mv.account_id, team_to_accounts_mv.team_id ] }

Fragment 52686 (Actor 738188,738187)
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 52687 (Actor 738232,738231)
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 52688 (Actor 738192,738191)
StreamLocalityProvider { locality_columns: [team_to_accounts_mv.account_id] }
├── output: [ team_to_accounts_mv.team_id, team_to_accounts_mv.account_id ]
├── stream key: [ team_to_accounts_mv.account_id, team_to_accounts_mv.team_id ]
└── MergeExecutor { output: [ team_to_accounts_mv.team_id, team_to_accounts_mv.account_id ], stream key: [ team_to_accounts_mv.team_id, team_to_accounts_mv.account_id ] }

Fragment 52689 (Actor 738234,738233)
StreamTableScan { table: team_to_accounts_mv, columns: [team_id, account_id] }
├── output: [ team_to_accounts_mv.team_id, team_to_accounts_mv.account_id ]
├── stream key: [ team_to_accounts_mv.team_id, team_to_accounts_mv.account_id ]
├── Upstream { output: [ team_id, account_id ], stream key: [] }
└── BatchPlanNode { output: [ team_id, account_id ], stream key: [] }

Fragment 52690 (Actor 738194,738195)
StreamProject { exprs: [team_to_portfolios_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_portfolio_id, tasks_dm.task_id, team_to_portfolios_mv.team_id, 2:Int32] }
├── output: [ team_to_portfolios_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_portfolio_id, tasks_dm.task_id, team_to_portfolios_mv.team_id, 2:Int32 ]
├── stream key: [ tasks_dm.resource_portfolio_id, tasks_dm.task_id, team_to_portfolios_mv.team_id ]
└── MergeExecutor
    ├── output: [ team_to_portfolios_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_portfolio_id, team_to_portfolios_mv.portfolio_id ]
    └── stream key: [ tasks_dm.resource_portfolio_id, tasks_dm.task_id, team_to_portfolios_mv.team_id ]

Fragment 52691 (Actor 738197,738196)
StreamSyncLogStore
├── output: [ team_to_portfolios_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_portfolio_id, team_to_portfolios_mv.portfolio_id ]
├── stream key: [ tasks_dm.resource_portfolio_id, tasks_dm.task_id, team_to_portfolios_mv.team_id ]
└── StreamHashJoin { type: Inner, predicate: tasks_dm.resource_portfolio_id = team_to_portfolios_mv.portfolio_id }
    ├── output: [ team_to_portfolios_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_portfolio_id, team_to_portfolios_mv.portfolio_id ]
    ├── stream key: [ tasks_dm.resource_portfolio_id, tasks_dm.task_id, team_to_portfolios_mv.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.team_id, team_to_portfolios_mv.portfolio_id ]
        └── stream key: [ team_to_portfolios_mv.portfolio_id, team_to_portfolios_mv.team_id ]

Fragment 52692 (Actor 738199,738198)
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 52693 (Actor 738236,738235)
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 52694 (Actor 738204,738203)
StreamLocalityProvider { locality_columns: [team_to_portfolios_mv.portfolio_id] }
├── output: [ team_to_portfolios_mv.team_id, team_to_portfolios_mv.portfolio_id ]
├── stream key: [ team_to_portfolios_mv.portfolio_id, team_to_portfolios_mv.team_id ]
└── MergeExecutor { output: [ team_to_portfolios_mv.team_id, team_to_portfolios_mv.portfolio_id ], stream key: [ team_to_portfolios_mv.team_id, team_to_portfolios_mv.portfolio_id ] }

Fragment 52695 (Actor 738237,738238)
StreamTableScan { table: team_to_portfolios_mv, columns: [team_id, portfolio_id] }
├── output: [ team_to_portfolios_mv.team_id, team_to_portfolios_mv.portfolio_id ]
├── stream key: [ team_to_portfolios_mv.team_id, team_to_portfolios_mv.portfolio_id ]
├── Upstream { output: [ team_id, portfolio_id ], stream key: [] }
└── BatchPlanNode { output: [ team_id, portfolio_id ], stream key: [] }

Fragment 52696 (Actor 738207,738208)
StreamProject { exprs: [team_to_parties_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_party_id, tasks_dm.task_id, team_to_parties_mv.team_id, 3:Int32] }
├── output: [ team_to_parties_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_party_id, tasks_dm.task_id, team_to_parties_mv.team_id, 3:Int32 ]
├── stream key: [ tasks_dm.resource_party_id, tasks_dm.task_id, team_to_parties_mv.team_id ]
└── MergeExecutor
    ├── 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.resource_party_id, tasks_dm.task_id, team_to_parties_mv.team_id ]

Fragment 52697 (Actor 738209,738210)
StreamSyncLogStore
├── 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.resource_party_id, tasks_dm.task_id, team_to_parties_mv.team_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.resource_party_id, tasks_dm.task_id, team_to_parties_mv.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.team_id, team_to_parties_mv.party_id ], stream key: [ team_to_parties_mv.party_id, team_to_parties_mv.team_id ] }

Fragment 52698 (Actor 738211,738212)
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 52699 (Actor 738239,738240)
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 52700 (Actor 738219,738220)
StreamLocalityProvider { locality_columns: [team_to_parties_mv.party_id] }
├── output: [ team_to_parties_mv.team_id, team_to_parties_mv.party_id ]
├── stream key: [ team_to_parties_mv.party_id, team_to_parties_mv.team_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 52701 (Actor 738221,738222)
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: [] }