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

← cluster alinma_bff objects party_task_unique_mv explain
Overview Objects Graph History
materialized view · alinma_bff.party_task_unique_mv profiled over 5s
seconds (1–30)
Stateful hash join (4 state tables) — consider a temporal join for dimension lookupsDynamic filter — verify it pairs with a temporal condition to clean state
116 operators
Materialize · alinma_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_party_id = party_active_customer_parties_…
2 actors
HashJoin · Inner · tasks_dm.resource_party_id = party_active_customer_parties_… 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
Filter · party_active_customer_parties_mv
0% idle 2 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
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 = 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% 2/s 0 actors
Now
0% 2/s 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% 2/s 0 actors
Now
0% 2/s 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 · alinma_bff.party_task_unique_mv Materialize alinma_bff.party_task_u… 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 = party_active_customer_parties_… SyncLogStore Inner · tasks_dm.resour… — · 2 actors HashJoin · Inner · tasks_dm.resource_party_id = party_active_customer_parties_… HashJoin Inner · tasks_dm.resour… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Filter · party_active_customer_parties_mv Filter party_active_customer_p… idle · 2 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 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 = 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 2/s · 0 actors Now Now 2/s · 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 2/s · 0 actors Now Now 2/s · 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 61137 (Actor 738303,738302)
StreamMaterialize { columns: [party_id, task_id, status, priority, is_overdue], 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, tasks_dm.is_overdue ]
├── 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, tasks_dm.is_overdue] }
    ├── output: [ party_account_direct_mv_next.party_id, tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.is_overdue ]
    ├── 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
        │   ├── tasks_dm.is_overdue
        │   ├── 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
            │   ├── tasks_dm.is_overdue
            │   ├── 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
                │   ├── tasks_dm.is_overdue
                │   ├── 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 61138 (Actor 738305,738304)
StreamUnion { all: true }
├── output:
│   ┌── party_account_direct_mv_next.party_id
│   ├── tasks_dm.task_id
│   ├── tasks_dm.status
│   ├── tasks_dm.priority
│   ├── tasks_dm.is_overdue
│   ├── 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
│   │   ├── tasks_dm.is_overdue
│   │   ├── 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.is_overdue
│   │   ├── tasks_dm.resource_portfolio_id
│   │   ├── tasks_dm.task_id
│   │   ├── party_active_portfolio_involvements_mv.party_id
│   │   ├── party_active_portfolio_involvements_mv.customer_relationship_id
│   │   ├── party_active_portfolio_involvements_mv.involvement_type
│   │   ├── 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.customer_relationship_id
│       └── party_active_portfolio_involvements_mv.involvement_type
└── MergeExecutor
    ├── output:
    │   ┌── tasks_dm.resource_party_id
    │   ├── tasks_dm.task_id
    │   ├── tasks_dm.status
    │   ├── tasks_dm.priority
    │   ├── tasks_dm.is_overdue
    │   ├── tasks_dm.resource_party_id
    │   ├── tasks_dm.task_id
    │   ├── null:Varchar
    │   ├── null:Varchar
    │   ├── null:Varchar
    │   ├── null:Varchar
    │   ├── null:Int32
    │   ├── null:Int32
    │   └── 2:Int32
    └── stream key: [ tasks_dm.resource_party_id, tasks_dm.task_id ]

Fragment 61139 (Actor 738334,738335)
StreamProject { exprs: [party_account_direct_mv_next.party_id, tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.is_overdue, 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, tasks_dm.is_overdue, 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.is_overdue, 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 61140 (Actor 738336,738337)
StreamSyncLogStore
├── output: [ party_account_direct_mv_next.party_id, tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.is_overdue, 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.is_overdue, 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, tasks_dm.is_overdue, 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 61141 (Actor 738346,738347)
StreamLocalityProvider { locality_columns: [party_account_direct_mv_next.party_id] }
├── output: [ tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.is_overdue, 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, tasks_dm.is_overdue, 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 61142 (Actor 738403,738404)
StreamSyncLogStore
├── output: [ tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.is_overdue, 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, tasks_dm.is_overdue, 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, tasks_dm.is_overdue ], 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 61143 (Actor 738406,738405)
StreamLocalityProvider { locality_columns: [tasks_dm.resource_account_id] } { output: [ tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_account_id, tasks_dm.is_overdue ], 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, tasks_dm.is_overdue ], stream key: [ tasks_dm.task_id ] }

Fragment 61144 (Actor 738451,738450)
StreamProject { exprs: [tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_account_id, tasks_dm.is_overdue] } { output: [ tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_account_id, tasks_dm.is_overdue ], 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.is_overdue, tasks_dm.disabled_at ], stream key: [ tasks_dm.task_id ] }
    └── StreamTableScan { table: tasks_dm, columns: [task_id, status, priority, resource_account_id, is_overdue, disabled_at] } { output: [ tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_account_id, tasks_dm.is_overdue, tasks_dm.disabled_at ], stream key: [ tasks_dm.task_id ] }
        ├── Upstream { output: [ task_id, status, priority, resource_account_id, is_overdue, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ task_id, status, priority, resource_account_id, is_overdue, disabled_at ], stream key: [] }

Fragment 61145 (Actor 738408,738407)
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 61146 (Actor 738409,738410)
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 61147 (Actor 738418,738417)
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 61148 (Actor 738416,738415)
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 61149 (Actor 738411)
StreamNow { output: [ now ], stream key: [] }

Fragment 61150 (Actor 738412)
StreamNow { output: [ now ], stream key: [] }

Fragment 61151 (Actor 738413,738414)
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 61152 (Actor 738428,738427)
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 61153 (Actor 738432,738431)
StreamProject { exprs: [party_active_portfolio_involvements_mv.party_id, tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.is_overdue, tasks_dm.resource_portfolio_id, tasks_dm.task_id, party_active_portfolio_involvements_mv.party_id, party_active_portfolio_involvements_mv.customer_relationship_id, party_active_portfolio_involvements_mv.involvement_type, 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.is_overdue, tasks_dm.resource_portfolio_id, tasks_dm.task_id, party_active_portfolio_involvements_mv.party_id, party_active_portfolio_involvements_mv.customer_relationship_id, party_active_portfolio_involvements_mv.involvement_type, 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.customer_relationship_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.is_overdue, tasks_dm.resource_portfolio_id, party_active_portfolio_involvements_mv.portfolio_id, party_active_portfolio_involvements_mv.customer_relationship_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.customer_relationship_id, party_active_portfolio_involvements_mv.involvement_type ]

Fragment 61154 (Actor 738429,738430)
StreamSyncLogStore
├── output: [ party_active_portfolio_involvements_mv.party_id, tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.is_overdue, tasks_dm.resource_portfolio_id, party_active_portfolio_involvements_mv.portfolio_id, party_active_portfolio_involvements_mv.customer_relationship_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.customer_relationship_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.is_overdue, tasks_dm.resource_portfolio_id, party_active_portfolio_involvements_mv.portfolio_id, party_active_portfolio_involvements_mv.customer_relationship_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.customer_relationship_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, tasks_dm.is_overdue ], 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.customer_relationship_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.customer_relationship_id, party_active_portfolio_involvements_mv.involvement_type ]

Fragment 61155 (Actor 738434,738433)
StreamLocalityProvider { locality_columns: [tasks_dm.resource_portfolio_id] } { output: [ tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_portfolio_id, tasks_dm.is_overdue ], 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, tasks_dm.is_overdue ], stream key: [ tasks_dm.task_id ] }

Fragment 61156 (Actor 738439,738438)
StreamProject { exprs: [tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_portfolio_id, tasks_dm.is_overdue] } { output: [ tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_portfolio_id, tasks_dm.is_overdue ], 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.is_overdue, tasks_dm.disabled_at ], stream key: [ tasks_dm.task_id ] }
    └── StreamTableScan { table: tasks_dm, columns: [task_id, status, priority, resource_portfolio_id, is_overdue, disabled_at] } { output: [ tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_portfolio_id, tasks_dm.is_overdue, tasks_dm.disabled_at ], stream key: [ tasks_dm.task_id ] }
        ├── Upstream { output: [ task_id, status, priority, resource_portfolio_id, is_overdue, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ task_id, status, priority, resource_portfolio_id, is_overdue, disabled_at ], stream key: [] }

Fragment 61157 (Actor 738440,738441)
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.customer_relationship_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.customer_relationship_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.customer_relationship_id, party_active_portfolio_involvements_mv.involvement_type ]
    └── stream key: [ party_active_portfolio_involvements_mv.party_id, party_active_portfolio_involvements_mv.customer_relationship_id, party_active_portfolio_involvements_mv.portfolio_id, party_active_portfolio_involvements_mv.involvement_type ]

Fragment 61158 (Actor 738454,738455)
StreamTableScan { table: party_active_portfolio_involvements_mv, columns: [party_id, portfolio_id, customer_relationship_id, involvement_type] }
├── output: [ party_active_portfolio_involvements_mv.party_id, party_active_portfolio_involvements_mv.portfolio_id, party_active_portfolio_involvements_mv.customer_relationship_id, party_active_portfolio_involvements_mv.involvement_type ]
├── stream key: [ party_active_portfolio_involvements_mv.party_id, party_active_portfolio_involvements_mv.customer_relationship_id, party_active_portfolio_involvements_mv.portfolio_id, party_active_portfolio_involvements_mv.involvement_type ]
├── Upstream { output: [ party_id, portfolio_id, customer_relationship_id, involvement_type ], stream key: [] }
└── BatchPlanNode { output: [ party_id, portfolio_id, customer_relationship_id, involvement_type ], stream key: [] }

Fragment 61159 (Actor 738444,738445)
StreamProject { exprs: [tasks_dm.resource_party_id, tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.is_overdue, tasks_dm.resource_party_id, tasks_dm.task_id, null:Varchar, null:Varchar, null:Varchar, null:Varchar, null:Int32, null:Int32, 2:Int32] }
├── output: [ tasks_dm.resource_party_id, tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.is_overdue, tasks_dm.resource_party_id, tasks_dm.task_id, null:Varchar, null:Varchar, null:Varchar, null:Varchar, null:Int32, null:Int32, 2:Int32 ]
├── stream key: [ tasks_dm.resource_party_id, tasks_dm.task_id ]
└── MergeExecutor { output: [ tasks_dm.resource_party_id, tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.is_overdue, party_active_customer_parties_mv.party_id ], stream key: [ tasks_dm.resource_party_id, tasks_dm.task_id ] }

Fragment 61160 (Actor 738446,738447)
StreamSyncLogStore { output: [ tasks_dm.resource_party_id, tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.is_overdue, party_active_customer_parties_mv.party_id ], stream key: [ tasks_dm.resource_party_id, tasks_dm.task_id ] }
└── StreamHashJoin { type: Inner, predicate: tasks_dm.resource_party_id = party_active_customer_parties_mv.party_id } { output: [ tasks_dm.resource_party_id, tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.is_overdue, party_active_customer_parties_mv.party_id ], stream key: [ tasks_dm.resource_party_id, tasks_dm.task_id ] }
    ├── MergeExecutor { output: [ tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_party_id, tasks_dm.is_overdue ], stream key: [ tasks_dm.resource_party_id, tasks_dm.task_id ] }
    └── MergeExecutor { output: [ party_active_customer_parties_mv.party_id ], stream key: [ party_active_customer_parties_mv.party_id ] }

Fragment 61161 (Actor 738449,738448)
StreamLocalityProvider { locality_columns: [tasks_dm.resource_party_id] } { output: [ tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_party_id, tasks_dm.is_overdue ], stream key: [ tasks_dm.resource_party_id, tasks_dm.task_id ] }
└── MergeExecutor { output: [ tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_party_id, tasks_dm.is_overdue ], stream key: [ tasks_dm.task_id ] }

Fragment 61162 (Actor 738459,738458)
StreamProject { exprs: [tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_party_id, tasks_dm.is_overdue] } { output: [ tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_party_id, tasks_dm.is_overdue ], stream key: [ tasks_dm.task_id ] }
└── StreamFilter { predicate: Not(IsNull(tasks_dm.resource_party_id)) AND 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_party_id, tasks_dm.is_overdue, tasks_dm.disabled_at ], stream key: [ tasks_dm.task_id ] }
    └── StreamTableScan { table: tasks_dm, columns: [task_id, status, priority, resource_party_id, is_overdue, disabled_at] } { output: [ tasks_dm.task_id, tasks_dm.status, tasks_dm.priority, tasks_dm.resource_party_id, tasks_dm.is_overdue, tasks_dm.disabled_at ], stream key: [ tasks_dm.task_id ] }
        ├── Upstream { output: [ task_id, status, priority, resource_party_id, is_overdue, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ task_id, status, priority, resource_party_id, is_overdue, disabled_at ], stream key: [] }

Fragment 61163 (Actor 738460,738461)
StreamFilter { predicate: Not(IsNull(party_active_customer_parties_mv.party_id)) } { output: [ party_active_customer_parties_mv.party_id ], stream key: [ party_active_customer_parties_mv.party_id ] }
└── 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: [] }