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: [] }