Fragment 63106 (Actor 745361,745360)
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 63107 (Actor 745363,745362)
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_next.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_next.party_id
│ │ ├── party_active_portfolio_involvements_mv_next.customer_relationship_id
│ │ ├── party_active_portfolio_involvements_mv_next.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_next.party_id
│ ├── party_active_portfolio_involvements_mv_next.customer_relationship_id
│ └── party_active_portfolio_involvements_mv_next.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 63108 (Actor 745458,745457)
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_next.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 63109 (Actor 745455,745456)
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_next.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_next.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_next.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_next.party_id ], stream key: [ party_active_customer_parties_mv_next.party_id ] }
Fragment 63110 (Actor 745460,745459)
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 63111 (Actor 745461,745462)
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 63112 (Actor 745463,745464)
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 63113 (Actor 745269,745268)
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 63114 (Actor 745466,745465)
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 63115 (Actor 745468,745467)
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 63116 (Actor 745474,745473)
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 63117 (Actor 745476,745475)
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 63118 (Actor 745469)
StreamNow { output: [ now ], stream key: [] }
Fragment 63119 (Actor 745470)
StreamNow { output: [ now ], stream key: [] }
Fragment 63120 (Actor 745471,745472)
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 63121 (Actor 745271,745270)
StreamTableScan { table: party_active_customer_parties_mv_next, columns: [party_id] } { output: [ party_active_customer_parties_mv_next.party_id ], stream key: [ party_active_customer_parties_mv_next.party_id ] }
├── Upstream { output: [ party_id ], stream key: [] }
└── BatchPlanNode { output: [ party_id ], stream key: [] }
Fragment 63122 (Actor 745481,745482)
StreamProject { exprs: [party_active_portfolio_involvements_mv_next.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_next.party_id, party_active_portfolio_involvements_mv_next.customer_relationship_id, party_active_portfolio_involvements_mv_next.involvement_type, null:Varchar, null:Int32, null:Int32, 1:Int32] }
├── output: [ party_active_portfolio_involvements_mv_next.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_next.party_id, party_active_portfolio_involvements_mv_next.customer_relationship_id, party_active_portfolio_involvements_mv_next.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_next.party_id, party_active_portfolio_involvements_mv_next.customer_relationship_id, party_active_portfolio_involvements_mv_next.involvement_type ]
└── MergeExecutor
├── output: [ party_active_portfolio_involvements_mv_next.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_next.portfolio_id, party_active_portfolio_involvements_mv_next.customer_relationship_id, party_active_portfolio_involvements_mv_next.involvement_type ]
└── stream key: [ tasks_dm.resource_portfolio_id, tasks_dm.task_id, party_active_portfolio_involvements_mv_next.party_id, party_active_portfolio_involvements_mv_next.customer_relationship_id, party_active_portfolio_involvements_mv_next.involvement_type ]
Fragment 63123 (Actor 745483,745484)
StreamSyncLogStore
├── output: [ party_active_portfolio_involvements_mv_next.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_next.portfolio_id, party_active_portfolio_involvements_mv_next.customer_relationship_id, party_active_portfolio_involvements_mv_next.involvement_type ]
├── stream key: [ tasks_dm.resource_portfolio_id, tasks_dm.task_id, party_active_portfolio_involvements_mv_next.party_id, party_active_portfolio_involvements_mv_next.customer_relationship_id, party_active_portfolio_involvements_mv_next.involvement_type ]
└── StreamHashJoin { type: Inner, predicate: tasks_dm.resource_portfolio_id = party_active_portfolio_involvements_mv_next.portfolio_id }
├── output: [ party_active_portfolio_involvements_mv_next.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_next.portfolio_id, party_active_portfolio_involvements_mv_next.customer_relationship_id, party_active_portfolio_involvements_mv_next.involvement_type ]
├── stream key: [ tasks_dm.resource_portfolio_id, tasks_dm.task_id, party_active_portfolio_involvements_mv_next.party_id, party_active_portfolio_involvements_mv_next.customer_relationship_id, party_active_portfolio_involvements_mv_next.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_next.party_id, party_active_portfolio_involvements_mv_next.portfolio_id, party_active_portfolio_involvements_mv_next.customer_relationship_id, party_active_portfolio_involvements_mv_next.involvement_type ]
└── stream key: [ party_active_portfolio_involvements_mv_next.portfolio_id, party_active_portfolio_involvements_mv_next.party_id, party_active_portfolio_involvements_mv_next.customer_relationship_id, party_active_portfolio_involvements_mv_next.involvement_type ]
Fragment 63124 (Actor 745490,745489)
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 63125 (Actor 745510,745509)
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 63126 (Actor 745492,745491)
StreamLocalityProvider { locality_columns: [party_active_portfolio_involvements_mv_next.portfolio_id] }
├── output: [ party_active_portfolio_involvements_mv_next.party_id, party_active_portfolio_involvements_mv_next.portfolio_id, party_active_portfolio_involvements_mv_next.customer_relationship_id, party_active_portfolio_involvements_mv_next.involvement_type ]
├── stream key: [ party_active_portfolio_involvements_mv_next.portfolio_id, party_active_portfolio_involvements_mv_next.party_id, party_active_portfolio_involvements_mv_next.customer_relationship_id, party_active_portfolio_involvements_mv_next.involvement_type ]
└── MergeExecutor
├── output: [ party_active_portfolio_involvements_mv_next.party_id, party_active_portfolio_involvements_mv_next.portfolio_id, party_active_portfolio_involvements_mv_next.customer_relationship_id, party_active_portfolio_involvements_mv_next.involvement_type ]
└── stream key: [ party_active_portfolio_involvements_mv_next.party_id, party_active_portfolio_involvements_mv_next.customer_relationship_id, party_active_portfolio_involvements_mv_next.portfolio_id, party_active_portfolio_involvements_mv_next.involvement_type ]
Fragment 63127 (Actor 745512,745511)
StreamTableScan { table: party_active_portfolio_involvements_mv_next, columns: [party_id, portfolio_id, customer_relationship_id, involvement_type] }
├── output: [ party_active_portfolio_involvements_mv_next.party_id, party_active_portfolio_involvements_mv_next.portfolio_id, party_active_portfolio_involvements_mv_next.customer_relationship_id, party_active_portfolio_involvements_mv_next.involvement_type ]
├── stream key: [ party_active_portfolio_involvements_mv_next.party_id, party_active_portfolio_involvements_mv_next.customer_relationship_id, party_active_portfolio_involvements_mv_next.portfolio_id, party_active_portfolio_involvements_mv_next.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 63128 (Actor 745495,745496)
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_next.party_id ], stream key: [ tasks_dm.resource_party_id, tasks_dm.task_id ] }
Fragment 63129 (Actor 745494,745493)
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_next.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_next.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_next.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_next.party_id ], stream key: [ party_active_customer_parties_mv_next.party_id ] }
Fragment 63130 (Actor 745497,745498)
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 63131 (Actor 745502,745501)
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 63132 (Actor 745513,745514)
StreamFilter { predicate: Not(IsNull(party_active_customer_parties_mv_next.party_id)) } { output: [ party_active_customer_parties_mv_next.party_id ], stream key: [ party_active_customer_parties_mv_next.party_id ] }
└── StreamTableScan { table: party_active_customer_parties_mv_next, columns: [party_id] } { output: [ party_active_customer_parties_mv_next.party_id ], stream key: [ party_active_customer_parties_mv_next.party_id ] }
├── Upstream { output: [ party_id ], stream key: [] }
└── BatchPlanNode { output: [ party_id ], stream key: [] }