RWM Console cluster: risingwave-ai-hub-int.ai-hub-rwm.svc.cluster.local

← cluster alpheya_experience_bff objects party_task_unique_mv explain
Overview Objects Graph History
materialized view · alpheya_experience_bff.party_task_unique_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 lookupsDynamic filter — verify it pairs with a temporal condition to clean state
93 operators
Materialize · alpheya_experience_bff.party_task_unique_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_portfolio_id = party_active_portfolio_inv…
2 actors
HashJoin · Inner · tasks_dm.resource_portfolio_id = party_active_portfolio_inv… 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 · party_active_portfolio_involvements_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 · party_account_direct_mv_next.party_id = party_active_custom…
2 actors
HashJoin · Inner · party_account_direct_mv_next.party_id = party_active_custom… 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 · party_active_customer_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
SyncLogStore · Inner · tasks_dm.resource_account_id = party_account_direct_mv_next…
2 actors
HashJoin · Inner · tasks_dm.resource_account_id = party_account_direct_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
Union
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · IsNull(party_account_direct_mv_next.effective_end_date)
2 actors
Filter · IsNull(party_account_direct_mv_next.effective_end_date)
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · party_account_direct_mv_next
2 actors
DynamicFilter · party_account_direct_mv_next Dynamic filter — verify it pairs with a temporal condition to clean state
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Now
0% idle 1 actor
Project · party_account_direct_mv_next
2 actors
Filter · party_account_direct_mv_next
0% idle 2 actors
StreamScan · party_account_direct_mv_next
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · ($expr2 > now), output_watermarks: [[$expr2]], output: [par…
2 actors
DynamicFilter · ($expr2 > now), output_watermarks: [[$expr2]], output: [par… Dynamic filter — verify it pairs with a temporal condition to clean state
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Now
0% idle 1 actor
Project · ($expr2 > now), output_watermarks: [[$expr2]], output: [par…
2 actors
Filter · ($expr2 > now), output_watermarks: [[$expr2]], output: [par…
0% idle 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 · alpheya_experience_bff.party_task_unique_mv Materialize alpheya_experience_bff.… 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_portfolio_id = party_active_portfolio_inv… SyncLogStore Inner · tasks_dm.resour… — · 2 actors HashJoin · Inner · tasks_dm.resource_portfolio_id = party_active_portfolio_inv… 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 · party_active_portfolio_involvements_mv StreamScan party_active_portfolio_… 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 · party_account_direct_mv_next.party_id = party_active_custom… SyncLogStore Inner · party_account_d… — · 2 actors HashJoin · Inner · party_account_direct_mv_next.party_id = party_active_custom… HashJoin Inner · party_account_d… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · party_active_customer_parties_mv StreamScan party_active_customer_p… 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 SyncLogStore · Inner · tasks_dm.resource_account_id = party_account_direct_mv_next… SyncLogStore Inner · tasks_dm.resour… — · 2 actors HashJoin · Inner · tasks_dm.resource_account_id = party_account_direct_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 Union Union idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · IsNull(party_account_direct_mv_next.effective_end_date) Project IsNull(party_account_di… — · 2 actors Filter · IsNull(party_account_direct_mv_next.effective_end_date) Filter IsNull(party_account_di… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · party_account_direct_mv_next Project party_account_direct_mv… — · 2 actors DynamicFilter · party_account_direct_mv_next DynamicFilter party_account_direct_mv… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Now Now idle · 1 actor Project · party_account_direct_mv_next Project party_account_direct_mv… — · 2 actors Filter · party_account_direct_mv_next Filter party_account_direct_mv… idle · 2 actors StreamScan · party_account_direct_mv_next StreamScan party_account_direct_mv… idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · ($expr2 > now), output_watermarks: [[$expr2]], output: [par… Project ($expr2 > now), output_… — · 2 actors DynamicFilter · ($expr2 > now), output_watermarks: [[$expr2]], output: [par… DynamicFilter ($expr2 > now), output_… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Now Now idle · 1 actor Project · ($expr2 > now), output_watermarks: [[$expr2]], output: [par… Project ($expr2 > now), output_… — · 2 actors Filter · ($expr2 > now), output_watermarks: [[$expr2]], output: [par… Filter ($expr2 > now), output_… idle · 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 17963 (Actor 157815,157814)
StreamMaterialize { columns: [party_id, task_id, status, priority], stream_key: [party_id, task_id], pk_columns: [party_id, task_id], pk_conflict: NoCheck }
├── output: [ party_account_direct_mv_next.party_id, tasks_dm.task_id, tasks_dm.status, tasks_dm.priority ]
├── stream key: [ party_account_direct_mv_next.party_id, tasks_dm.task_id ]
└── StreamProject { exprs: [party_account_direct_mv_next.party_id, tasks_dm.task_id, tasks_dm.status, tasks_dm.priority] }
    ├── output: [ party_account_direct_mv_next.party_id, tasks_dm.task_id, tasks_dm.status, tasks_dm.priority ]
    ├── stream key: [ party_account_direct_mv_next.party_id, tasks_dm.task_id ]
    └── StreamGroupTopN { order: [tasks_dm.task_id ASC], limit: 1, offset: 0, group_key: [party_account_direct_mv_next.party_id, tasks_dm.task_id] }
        ├── output:
        │   ┌── party_account_direct_mv_next.party_id
        │   ├── tasks_dm.task_id
        │   ├── tasks_dm.status
        │   ├── tasks_dm.priority
        │   ├── party_account_direct_mv_next.party_id
        │   ├── tasks_dm.resource_account_id
        │   ├── tasks_dm.task_id
        │   ├── party_account_direct_mv_next.party_involvements_dm.customer_relationship_id
        │   ├── party_account_direct_mv_next.party_involvements_dm.entity_id
        │   ├── party_account_direct_mv_next.party_involvements_dm.id
        │   ├── party_account_direct_mv_next.$src
        │   ├── $src
        │   └── $src
        ├── stream key: [ party_account_direct_mv_next.party_id, tasks_dm.task_id ]
        └── StreamLocalityProvider { locality_columns: [party_account_direct_mv_next.party_id, tasks_dm.task_id] }
            ├── output:
            │   ┌── party_account_direct_mv_next.party_id
            │   ├── tasks_dm.task_id
            │   ├── tasks_dm.status
            │   ├── tasks_dm.priority
            │   ├── party_account_direct_mv_next.party_id
            │   ├── tasks_dm.resource_account_id
            │   ├── tasks_dm.task_id
            │   ├── party_account_direct_mv_next.party_involvements_dm.customer_relationship_id
            │   ├── party_account_direct_mv_next.party_involvements_dm.entity_id
            │   ├── party_account_direct_mv_next.party_involvements_dm.id
            │   ├── party_account_direct_mv_next.$src
            │   ├── $src
            │   └── $src
            ├── stream key:
            │   ┌── party_account_direct_mv_next.party_id
            │   ├── tasks_dm.task_id
            │   ├── party_account_direct_mv_next.party_id
            │   ├── tasks_dm.resource_account_id
            │   ├── tasks_dm.task_id
            │   ├── party_account_direct_mv_next.party_involvements_dm.customer_relationship_id
            │   ├── party_account_direct_mv_next.party_involvements_dm.entity_id
            │   ├── party_account_direct_mv_next.party_involvements_dm.id
            │   ├── party_account_direct_mv_next.$src
            │   ├── $src
            │   └── $src
            └── MergeExecutor
                ├── output:
                │   ┌── party_account_direct_mv_next.party_id
                │   ├── tasks_dm.task_id
                │   ├── tasks_dm.status
                │   ├── tasks_dm.priority
                │   ├── party_account_direct_mv_next.party_id
                │   ├── tasks_dm.resource_account_id
                │   ├── tasks_dm.task_id
                │   ├── party_account_direct_mv_next.party_involvements_dm.customer_relationship_id
                │   ├── party_account_direct_mv_next.party_involvements_dm.entity_id
                │   ├── party_account_direct_mv_next.party_involvements_dm.id
                │   ├── party_account_direct_mv_next.$src
                │   ├── $src
                │   └── $src
                └── stream key:
                    ┌── party_account_direct_mv_next.party_id
                    ├── tasks_dm.resource_account_id
                    ├── tasks_dm.task_id
                    ├── party_account_direct_mv_next.party_involvements_dm.customer_relationship_id
                    ├── party_account_direct_mv_next.party_involvements_dm.entity_id
                    ├── party_account_direct_mv_next.party_involvements_dm.id
                    ├── party_account_direct_mv_next.$src
                    ├── $src
                    └── $src

Fragment 17964 (Actor 158321,158320)
StreamUnion { all: true }
├── output:
│   ┌── party_account_direct_mv_next.party_id
│   ├── tasks_dm.task_id
│   ├── tasks_dm.status
│   ├── tasks_dm.priority
│   ├── party_account_direct_mv_next.party_id
│   ├── tasks_dm.resource_account_id
│   ├── tasks_dm.task_id
│   ├── party_account_direct_mv_next.party_involvements_dm.customer_relationship_id
│   ├── party_account_direct_mv_next.party_involvements_dm.entity_id
│   ├── party_account_direct_mv_next.party_involvements_dm.id
│   ├── party_account_direct_mv_next.$src
│   ├── $src
│   └── $src
├── stream key:
│   ┌── party_account_direct_mv_next.party_id
│   ├── tasks_dm.resource_account_id
│   ├── tasks_dm.task_id
│   ├── party_account_direct_mv_next.party_involvements_dm.customer_relationship_id
│   ├── party_account_direct_mv_next.party_involvements_dm.entity_id
│   ├── party_account_direct_mv_next.party_involvements_dm.id
│   ├── party_account_direct_mv_next.$src
│   ├── $src
│   └── $src
├── MergeExecutor
│   ├── output:
│   │   ┌── party_account_direct_mv_next.party_id
│   │   ├── tasks_dm.task_id
│   │   ├── tasks_dm.status
│   │   ├── tasks_dm.priority
│   │   ├── party_account_direct_mv_next.party_id
│   │   ├── tasks_dm.resource_account_id
│   │   ├── tasks_dm.task_id
│   │   ├── party_account_direct_mv_next.party_involvements_dm.customer_relationship_id
│   │   ├── party_account_direct_mv_next.party_involvements_dm.entity_id
│   │   ├── party_account_direct_mv_next.party_involvements_dm.id
│   │   ├── party_account_direct_mv_next.$src
│   │   ├── $src
│   │   └── 0:Int32
│   └── stream key:
│       ┌── party_account_direct_mv_next.party_id
│       ├── tasks_dm.resource_account_id
│       ├── tasks_dm.task_id
│       ├── party_account_direct_mv_next.party_involvements_dm.customer_relationship_id
│       ├── party_account_direct_mv_next.party_involvements_dm.entity_id
│       ├── party_account_direct_mv_next.party_involvements_dm.id
│       ├── party_account_direct_mv_next.$src
│       └── $src
└── MergeExecutor
    ├── output:
    │   ┌── party_active_portfolio_involvements_mv.party_id
    │   ├── tasks_dm.task_id
    │   ├── tasks_dm.status
    │   ├── tasks_dm.priority
    │   ├── tasks_dm.resource_portfolio_id
    │   ├── tasks_dm.task_id
    │   ├── party_active_portfolio_involvements_mv.party_id
    │   ├── party_active_portfolio_involvements_mv.involvement_type
    │   ├── null:Varchar
    │   ├── null:Varchar
    │   ├── null:Int32
    │   ├── null:Int32
    │   └── 1:Int32
    └── stream key:
        ┌── tasks_dm.resource_portfolio_id
        ├── tasks_dm.task_id
        ├── party_active_portfolio_involvements_mv.party_id
        └── party_active_portfolio_involvements_mv.involvement_type

Fragment 17965 (Actor 158323,158322)
StreamProject { exprs: [party_account_direct_mv_next.party_id, tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, party_account_direct_mv_next.party_id, tasks_dm.resource_account_id, tasks_dm.task_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, $src, 0:Int32] }
├── output: [ party_account_direct_mv_next.party_id, tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, party_account_direct_mv_next.party_id, tasks_dm.resource_account_id, tasks_dm.task_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, $src, 0:Int32 ]
├── stream key: [ party_account_direct_mv_next.party_id, tasks_dm.resource_account_id, tasks_dm.task_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, $src ]
└── MergeExecutor
    ├── output: [ party_account_direct_mv_next.party_id, tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_account_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, $src, party_active_customer_parties_mv.party_id ]
    └── stream key: [ party_account_direct_mv_next.party_id, tasks_dm.resource_account_id, tasks_dm.task_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, $src ]

Fragment 17966 (Actor 158325,158324)
StreamSyncLogStore
├── output: [ party_account_direct_mv_next.party_id, tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_account_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, $src, party_active_customer_parties_mv.party_id ]
├── stream key: [ party_account_direct_mv_next.party_id, tasks_dm.resource_account_id, tasks_dm.task_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, $src ]
└── StreamHashJoin { type: Inner, predicate: party_account_direct_mv_next.party_id = party_active_customer_parties_mv.party_id }
    ├── output: [ party_account_direct_mv_next.party_id, tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_account_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, $src, party_active_customer_parties_mv.party_id ]
    ├── stream key: [ party_account_direct_mv_next.party_id, tasks_dm.resource_account_id, tasks_dm.task_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, $src ]
    ├── MergeExecutor
    │   ├── output: [ tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, party_account_direct_mv_next.party_id, tasks_dm.resource_account_id, party_account_direct_mv_next.account_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, $src ]
    │   └── stream key: [ party_account_direct_mv_next.party_id, tasks_dm.resource_account_id, tasks_dm.task_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, $src ]
    └── MergeExecutor { output: [ party_active_customer_parties_mv.party_id ], stream key: [ party_active_customer_parties_mv.party_id ] }

Fragment 17967 (Actor 158329,158328)
StreamLocalityProvider { locality_columns: [party_account_direct_mv_next.party_id] }
├── output: [ tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, party_account_direct_mv_next.party_id, tasks_dm.resource_account_id, party_account_direct_mv_next.account_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, $src ]
├── stream key: [ party_account_direct_mv_next.party_id, tasks_dm.resource_account_id, tasks_dm.task_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, $src ]
└── MergeExecutor
    ├── output: [ tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, party_account_direct_mv_next.party_id, tasks_dm.resource_account_id, party_account_direct_mv_next.account_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, $src ]
    └── stream key: [ tasks_dm.resource_account_id, tasks_dm.task_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, $src ]

Fragment 17968 (Actor 158330,158331)
StreamSyncLogStore
├── output: [ tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, party_account_direct_mv_next.party_id, tasks_dm.resource_account_id, party_account_direct_mv_next.account_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, $src ]
├── stream key: [ tasks_dm.resource_account_id, tasks_dm.task_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, $src ]
└── StreamHashJoin { type: Inner, predicate: tasks_dm.resource_account_id = party_account_direct_mv_next.account_id }
    ├── output: [ tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, party_account_direct_mv_next.party_id, tasks_dm.resource_account_id, party_account_direct_mv_next.account_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, $src ]
    ├── stream key: [ tasks_dm.resource_account_id, tasks_dm.task_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, $src ]
    ├── MergeExecutor { output: [ tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_account_id ], stream key: [ tasks_dm.resource_account_id, tasks_dm.task_id ] }
    └── MergeExecutor
        ├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.account_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, $src ]
        └── stream key: [ party_account_direct_mv_next.account_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, $src ]

Fragment 17969 (Actor 158333,158332)
StreamLocalityProvider { locality_columns: [tasks_dm.resource_account_id] } { output: [ tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_account_id ], stream key: [ tasks_dm.resource_account_id, tasks_dm.task_id ] }
└── MergeExecutor { output: [ tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_account_id ], stream key: [ tasks_dm.task_id ] }

Fragment 17970 (Actor 158378,158379)
StreamProject { exprs: [tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_account_id] } { output: [ tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_account_id ], stream key: [ tasks_dm.task_id ] }
└── StreamFilter { predicate: IsNull(tasks_dm.disabled_at) AND In(tasks_dm.status, 'TO_DO':Varchar, 'IN_PROGRESS':Varchar) } { output: [ tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_account_id, tasks_dm.disabled_at ], stream key: [ tasks_dm.task_id ] }
    └── StreamTableScan { table: tasks_dm, columns: [task_id, status, priority, resource_account_id, disabled_at] } { output: [ tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_account_id, tasks_dm.disabled_at ], stream key: [ tasks_dm.task_id ] }
        ├── Upstream { output: [ task_id, status, priority, resource_account_id, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ task_id, status, priority, resource_account_id, disabled_at ], stream key: [] }

Fragment 17971 (Actor 158334,158335)
StreamLocalityProvider { locality_columns: [party_account_direct_mv_next.account_id] }
├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.account_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, $src ]
├── stream key: [ party_account_direct_mv_next.account_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, $src ]
└── MergeExecutor
    ├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.account_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, $src ]
    └── stream key: [ party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, $src ]

Fragment 17972 (Actor 158336,158337)
StreamUnion { all: true }
├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.account_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, $src ]
├── stream key: [ party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, $src ]
├── MergeExecutor
│   ├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.account_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, 0:Int32 ]
│   └── stream key: [ party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src ]
└── MergeExecutor
    ├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.account_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, 1:Int32 ]
    └── stream key: [ party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src ]

Fragment 17973 (Actor 158349,158348)
StreamProject { exprs: [party_account_direct_mv_next.party_id, party_account_direct_mv_next.account_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, 0:Int32] }
├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.account_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, 0:Int32 ]
├── stream key: [ party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src ]
└── StreamDynamicFilter { predicate: ($expr2 > now), output_watermarks: [[$expr2]], output: [party_account_direct_mv_next.party_id, party_account_direct_mv_next.account_id, $expr2, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src], cleaned_by_watermark: true }
    ├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.account_id, $expr2, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src ]
    ├── stream key: [ party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src ]
    ├── StreamProject { exprs: [party_account_direct_mv_next.party_id, party_account_direct_mv_next.account_id, AtTimeZone(party_account_direct_mv_next.effective_end_date::Timestamp, 'UTC':Varchar) as $expr2, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src] }
    │   ├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.account_id, $expr2, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src ]
    │   ├── stream key: [ party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src ]
    │   └── StreamFilter { predicate: IsNotTrue(IsNull(party_account_direct_mv_next.effective_end_date)) }
    │       ├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.account_id, party_account_direct_mv_next.effective_end_date, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src ]
    │       ├── stream key: [ party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src ]
    │       └── MergeExecutor
    │           ├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.account_id, party_account_direct_mv_next.effective_end_date, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src ]
    │           └── stream key: [ party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src ]
    └── MergeExecutor { output: [ now ], stream key: [] }

Fragment 17974 (Actor 158347,158346)
StreamProject { exprs: [party_account_direct_mv_next.party_id, party_account_direct_mv_next.account_id, party_account_direct_mv_next.effective_end_date, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src] }
├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.account_id, party_account_direct_mv_next.effective_end_date, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src ]
├── stream key: [ party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src ]
└── StreamDynamicFilter { predicate: ($expr1 <= now), output: [party_account_direct_mv_next.party_id, party_account_direct_mv_next.account_id, party_account_direct_mv_next.effective_end_date, $expr1, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src], cleaned_by_watermark: true }
    ├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.account_id, party_account_direct_mv_next.effective_end_date, $expr1, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src ]
    ├── stream key: [ party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src ]
    ├── StreamProject { exprs: [party_account_direct_mv_next.party_id, party_account_direct_mv_next.account_id, party_account_direct_mv_next.effective_end_date, AtTimeZone(party_account_direct_mv_next.effective_start_date::Timestamp, 'UTC':Varchar) as $expr1, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src] }
    │   ├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.account_id, party_account_direct_mv_next.effective_end_date, $expr1, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src ]
    │   ├── stream key: [ party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src ]
    │   └── StreamFilter { predicate: (IsNotTrue(IsNull(party_account_direct_mv_next.effective_end_date)) OR IsNull(party_account_direct_mv_next.effective_end_date)) }
    │       ├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.account_id, party_account_direct_mv_next.effective_start_date, party_account_direct_mv_next.effective_end_date, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src ]
    │       ├── stream key: [ party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src ]
    │       └── StreamTableScan { table: party_account_direct_mv_next, columns: [party_id, account_id, effective_start_date, effective_end_date, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.id, $src] }
    │           ├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.account_id, party_account_direct_mv_next.effective_start_date, party_account_direct_mv_next.effective_end_date, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src ]
    │           ├── stream key: [ party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src ]
    │           ├── Upstream { output: [ party_id, account_id, effective_start_date, effective_end_date, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.id, $src ], stream key: [] }
    │           └── BatchPlanNode { output: [ party_id, account_id, effective_start_date, effective_end_date, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.id, $src ], stream key: [] }
    └── MergeExecutor { output: [ now ], stream key: [] }

Fragment 17975 (Actor 158342)
StreamNow { output: [ now ], stream key: [] }

Fragment 17976 (Actor 158343)
StreamNow { output: [ now ], stream key: [] }

Fragment 17977 (Actor 158344,158345)
StreamProject { exprs: [party_account_direct_mv_next.party_id, party_account_direct_mv_next.account_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, 1:Int32] }
├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.account_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src, 1:Int32 ]
├── stream key: [ party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src ]
└── StreamFilter { predicate: IsNull(party_account_direct_mv_next.effective_end_date) }
    ├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.account_id, party_account_direct_mv_next.effective_end_date, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src ]
    ├── stream key: [ party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src ]
    └── MergeExecutor
        ├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.account_id, party_account_direct_mv_next.effective_end_date, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src ]
        └── stream key: [ party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src ]

Fragment 17978 (Actor 158385,158384)
StreamTableScan { table: party_active_customer_parties_mv, columns: [party_id] } { output: [ party_active_customer_parties_mv.party_id ], stream key: [ party_active_customer_parties_mv.party_id ] }
├── Upstream { output: [ party_id ], stream key: [] }
└── BatchPlanNode { output: [ party_id ], stream key: [] }

Fragment 17979 (Actor 158369,158368)
StreamProject { exprs: [party_active_portfolio_involvements_mv.party_id, tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_portfolio_id, tasks_dm.task_id, party_active_portfolio_involvements_mv.party_id, party_active_portfolio_involvements_mv.involvement_type, null:Varchar, null:Varchar, null:Int32, null:Int32, 1:Int32] }
├── output: [ party_active_portfolio_involvements_mv.party_id, tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_portfolio_id, tasks_dm.task_id, party_active_portfolio_involvements_mv.party_id, party_active_portfolio_involvements_mv.involvement_type, null:Varchar, null:Varchar, null:Int32, null:Int32, 1:Int32 ]
├── stream key: [ tasks_dm.resource_portfolio_id, tasks_dm.task_id, party_active_portfolio_involvements_mv.party_id, party_active_portfolio_involvements_mv.involvement_type ]
└── MergeExecutor { output: [ party_active_portfolio_involvements_mv.party_id, tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_portfolio_id, party_active_portfolio_involvements_mv.portfolio_id, party_active_portfolio_involvements_mv.involvement_type ], stream key: [ tasks_dm.resource_portfolio_id, tasks_dm.task_id, party_active_portfolio_involvements_mv.party_id, party_active_portfolio_involvements_mv.involvement_type ] }

Fragment 17980 (Actor 158370,158371)
StreamSyncLogStore { output: [ party_active_portfolio_involvements_mv.party_id, tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_portfolio_id, party_active_portfolio_involvements_mv.portfolio_id, party_active_portfolio_involvements_mv.involvement_type ], stream key: [ tasks_dm.resource_portfolio_id, tasks_dm.task_id, party_active_portfolio_involvements_mv.party_id, party_active_portfolio_involvements_mv.involvement_type ] }
└── StreamHashJoin { type: Inner, predicate: tasks_dm.resource_portfolio_id = party_active_portfolio_involvements_mv.portfolio_id }
    ├── output: [ party_active_portfolio_involvements_mv.party_id, tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_portfolio_id, party_active_portfolio_involvements_mv.portfolio_id, party_active_portfolio_involvements_mv.involvement_type ]
    ├── stream key: [ tasks_dm.resource_portfolio_id, tasks_dm.task_id, party_active_portfolio_involvements_mv.party_id, party_active_portfolio_involvements_mv.involvement_type ]
    ├── MergeExecutor { output: [ tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_portfolio_id ], stream key: [ tasks_dm.resource_portfolio_id, tasks_dm.task_id ] }
    └── MergeExecutor { output: [ party_active_portfolio_involvements_mv.party_id, party_active_portfolio_involvements_mv.portfolio_id, party_active_portfolio_involvements_mv.involvement_type ], stream key: [ party_active_portfolio_involvements_mv.portfolio_id, party_active_portfolio_involvements_mv.party_id, party_active_portfolio_involvements_mv.involvement_type ] }

Fragment 17981 (Actor 158373,158372)
StreamLocalityProvider { locality_columns: [tasks_dm.resource_portfolio_id] } { output: [ tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_portfolio_id ], stream key: [ tasks_dm.resource_portfolio_id, tasks_dm.task_id ] }
└── MergeExecutor { output: [ tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_portfolio_id ], stream key: [ tasks_dm.task_id ] }

Fragment 17982 (Actor 158387,158386)
StreamProject { exprs: [tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_portfolio_id] } { output: [ tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_portfolio_id ], stream key: [ tasks_dm.task_id ] }
└── StreamFilter { predicate: IsNull(tasks_dm.disabled_at) AND In(tasks_dm.status, 'TO_DO':Varchar, 'IN_PROGRESS':Varchar) } { output: [ tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_portfolio_id, tasks_dm.disabled_at ], stream key: [ tasks_dm.task_id ] }
    └── StreamTableScan { table: tasks_dm, columns: [task_id, status, priority, resource_portfolio_id, disabled_at] } { output: [ tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_portfolio_id, tasks_dm.disabled_at ], stream key: [ tasks_dm.task_id ] }
        ├── Upstream { output: [ task_id, status, priority, resource_portfolio_id, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ task_id, status, priority, resource_portfolio_id, disabled_at ], stream key: [] }

Fragment 17983 (Actor 158374,158375)
StreamLocalityProvider { locality_columns: [party_active_portfolio_involvements_mv.portfolio_id] } { output: [ party_active_portfolio_involvements_mv.party_id, party_active_portfolio_involvements_mv.portfolio_id, party_active_portfolio_involvements_mv.involvement_type ], stream key: [ party_active_portfolio_involvements_mv.portfolio_id, party_active_portfolio_involvements_mv.party_id, party_active_portfolio_involvements_mv.involvement_type ] }
└── MergeExecutor { output: [ party_active_portfolio_involvements_mv.party_id, party_active_portfolio_involvements_mv.portfolio_id, party_active_portfolio_involvements_mv.involvement_type ], stream key: [ party_active_portfolio_involvements_mv.party_id, party_active_portfolio_involvements_mv.portfolio_id, party_active_portfolio_involvements_mv.involvement_type ] }

Fragment 17984 (Actor 158376,158377)
StreamTableScan { table: party_active_portfolio_involvements_mv, columns: [party_id, portfolio_id, involvement_type] } { output: [ party_active_portfolio_involvements_mv.party_id, party_active_portfolio_involvements_mv.portfolio_id, party_active_portfolio_involvements_mv.involvement_type ], stream key: [ party_active_portfolio_involvements_mv.party_id, party_active_portfolio_involvements_mv.portfolio_id, party_active_portfolio_involvements_mv.involvement_type ] }
├── Upstream { output: [ party_id, portfolio_id, involvement_type ], stream key: [] }
└── BatchPlanNode { output: [ party_id, portfolio_id, involvement_type ], stream key: [] }