Job is idle — throughput ~0; structure shown.
Fragment 54202 (Actor 738618,738619)
StreamMaterialize { columns: [account_id, client_json], stream_key: [account_id], pk_columns: [account_id], pk_conflict: NoCheck }
├── output: [ party_account_owner_candidates_mv.account_id, $expr2 ]
├── stream key: [ party_account_owner_candidates_mv.account_id ]
└── StreamProject { exprs: [party_account_owner_candidates_mv.account_id, JsonbAccess(jsonb_agg($expr1 order_by(min(party_account_owner_candidates_mv.party_id) ASC)), 0:Int32) as $expr2] }
├── output: [ party_account_owner_candidates_mv.account_id, $expr2 ]
├── stream key: [ party_account_owner_candidates_mv.account_id ]
└── StreamHashAgg { group_key: [party_account_owner_candidates_mv.account_id], aggs: [jsonb_agg($expr1 order_by(min(party_account_owner_candidates_mv.party_id) ASC)), count] }
├── output: [ party_account_owner_candidates_mv.account_id, jsonb_agg($expr1 order_by(min(party_account_owner_candidates_mv.party_id) ASC)), count ]
├── stream key: [ party_account_owner_candidates_mv.account_id ]
└── StreamLocalityProvider { locality_columns: [party_account_owner_candidates_mv.account_id] }
├── output: [ party_account_owner_candidates_mv.account_id, min(party_account_owner_candidates_mv.party_id), $expr1, parties.id, reference_identifiers_next.id ]
├── stream key: [ party_account_owner_candidates_mv.account_id, parties.id, min(party_account_owner_candidates_mv.party_id), reference_identifiers_next.id ]
└── MergeExecutor
├── output: [ party_account_owner_candidates_mv.account_id, min(party_account_owner_candidates_mv.party_id), $expr1, parties.id, reference_identifiers_next.id ]
└── stream key: [ parties.id, min(party_account_owner_candidates_mv.party_id), party_account_owner_candidates_mv.account_id, reference_identifiers_next.id ]
Fragment 54203 (Actor 738620,738621)
StreamProject { exprs: [party_account_owner_candidates_mv.account_id, min(party_account_owner_candidates_mv.party_id), JsonbBuildObject('id':Varchar, parties.id, 'name':Varchar, JsonbBuildObject('en':Varchar, Coalesce(parties.display_name, '':Varchar), 'ar':Varchar, '':Varchar), 'client_type':Varchar, Case((parties.type = 'INDIVIDUAL':Varchar), 1:Int32, (parties.type = 'ORGANIZATION':Varchar), 2:Int32, 0:Int32), 'organization':Varchar, 1:Int32, 'external_ids':Varchar, JsonbBuildObject('cif':Varchar, Coalesce(reference_identifiers_next.value, '':Varchar))) as $expr1, parties.id, reference_identifiers_next.id] }
├── output: [ party_account_owner_candidates_mv.account_id, min(party_account_owner_candidates_mv.party_id), $expr1, parties.id, reference_identifiers_next.id ]
├── stream key: [ parties.id, min(party_account_owner_candidates_mv.party_id), party_account_owner_candidates_mv.account_id, reference_identifiers_next.id ]
└── MergeExecutor { output: [ party_account_owner_candidates_mv.account_id, min(party_account_owner_candidates_mv.party_id), parties.id, parties.type, parties.display_name, reference_identifiers_next.value, reference_identifiers_next.entity_id, reference_identifiers_next.id ], stream key: [ parties.id, min(party_account_owner_candidates_mv.party_id), party_account_owner_candidates_mv.account_id, reference_identifiers_next.id ] }
Fragment 54204 (Actor 738622,738623)
StreamSyncLogStore { output: [ party_account_owner_candidates_mv.account_id, min(party_account_owner_candidates_mv.party_id), parties.id, parties.type, parties.display_name, reference_identifiers_next.value, reference_identifiers_next.entity_id, reference_identifiers_next.id ], stream key: [ parties.id, min(party_account_owner_candidates_mv.party_id), party_account_owner_candidates_mv.account_id, reference_identifiers_next.id ] }
└── StreamHashJoin { type: LeftOuter, predicate: parties.id = reference_identifiers_next.entity_id } { output: [ party_account_owner_candidates_mv.account_id, min(party_account_owner_candidates_mv.party_id), parties.id, parties.type, parties.display_name, reference_identifiers_next.value, reference_identifiers_next.entity_id, reference_identifiers_next.id ], stream key: [ parties.id, min(party_account_owner_candidates_mv.party_id), party_account_owner_candidates_mv.account_id, reference_identifiers_next.id ] }
├── MergeExecutor { output: [ party_account_owner_candidates_mv.account_id, min(party_account_owner_candidates_mv.party_id), parties.id, parties.type, parties.display_name ], stream key: [ parties.id, min(party_account_owner_candidates_mv.party_id), party_account_owner_candidates_mv.account_id ] }
└── MergeExecutor { output: [ reference_identifiers_next.entity_id, reference_identifiers_next.value, reference_identifiers_next.id ], stream key: [ reference_identifiers_next.entity_id, reference_identifiers_next.id ] }
Fragment 54205 (Actor 738629,738628)
StreamLocalityProvider { locality_columns: [parties.id] } { output: [ party_account_owner_candidates_mv.account_id, min(party_account_owner_candidates_mv.party_id), parties.id, parties.type, parties.display_name ], stream key: [ parties.id, min(party_account_owner_candidates_mv.party_id), party_account_owner_candidates_mv.account_id ] }
└── MergeExecutor { output: [ party_account_owner_candidates_mv.account_id, min(party_account_owner_candidates_mv.party_id), parties.id, parties.type, parties.display_name ], stream key: [ min(party_account_owner_candidates_mv.party_id), party_account_owner_candidates_mv.account_id ] }
Fragment 54206 (Actor 738633,738632)
StreamSyncLogStore { output: [ party_account_owner_candidates_mv.account_id, min(party_account_owner_candidates_mv.party_id), parties.id, parties.type, parties.display_name ], stream key: [ min(party_account_owner_candidates_mv.party_id), party_account_owner_candidates_mv.account_id ] }
└── StreamHashJoin { type: Inner, predicate: min(party_account_owner_candidates_mv.party_id) = parties.id } { output: [ party_account_owner_candidates_mv.account_id, min(party_account_owner_candidates_mv.party_id), parties.id, parties.type, parties.display_name ], stream key: [ min(party_account_owner_candidates_mv.party_id), party_account_owner_candidates_mv.account_id ] }
├── MergeExecutor { output: [ party_account_owner_candidates_mv.account_id, min(party_account_owner_candidates_mv.party_id), count ], stream key: [ min(party_account_owner_candidates_mv.party_id), party_account_owner_candidates_mv.account_id ] }
└── MergeExecutor { output: [ parties.id, parties.type, parties.display_name ], stream key: [ parties.id ] }
Fragment 54207 (Actor 738634,738635)
StreamLocalityProvider { locality_columns: [min(party_account_owner_candidates_mv.party_id)] } { output: [ party_account_owner_candidates_mv.account_id, min(party_account_owner_candidates_mv.party_id), count ], stream key: [ min(party_account_owner_candidates_mv.party_id), party_account_owner_candidates_mv.account_id ] }
└── MergeExecutor { output: [ party_account_owner_candidates_mv.account_id, min(party_account_owner_candidates_mv.party_id), count ], stream key: [ party_account_owner_candidates_mv.account_id ] }
Fragment 54208 (Actor 738637,738636)
StreamFilter { predicate: (count = 1:Int32) } { output: [ party_account_owner_candidates_mv.account_id, min(party_account_owner_candidates_mv.party_id), count ], stream key: [ party_account_owner_candidates_mv.account_id ] }
└── StreamHashAgg { group_key: [party_account_owner_candidates_mv.account_id], aggs: [min(party_account_owner_candidates_mv.party_id), count] } { output: [ party_account_owner_candidates_mv.account_id, min(party_account_owner_candidates_mv.party_id), count ], stream key: [ party_account_owner_candidates_mv.account_id ] }
└── StreamLocalityProvider { locality_columns: [party_account_owner_candidates_mv.account_id] } { output: [ party_account_owner_candidates_mv.account_id, party_account_owner_candidates_mv.party_id ], stream key: [ party_account_owner_candidates_mv.account_id, party_account_owner_candidates_mv.party_id ] }
└── MergeExecutor { output: [ party_account_owner_candidates_mv.account_id, party_account_owner_candidates_mv.party_id ], stream key: [ party_account_owner_candidates_mv.account_id, party_account_owner_candidates_mv.party_id ] }
Fragment 54209 (Actor 738662,738663)
StreamProject { exprs: [party_account_owner_candidates_mv.account_id, party_account_owner_candidates_mv.party_id] } { output: [ party_account_owner_candidates_mv.account_id, party_account_owner_candidates_mv.party_id ], stream key: [ party_account_owner_candidates_mv.account_id, party_account_owner_candidates_mv.party_id ] }
└── StreamFilter { predicate: (party_account_owner_candidates_mv.owner_type = 'PRIMARY':Varchar) } { output: [ party_account_owner_candidates_mv.account_id, party_account_owner_candidates_mv.party_id, party_account_owner_candidates_mv.owner_type ], stream key: [ party_account_owner_candidates_mv.account_id, party_account_owner_candidates_mv.party_id ] }
└── StreamTableScan { table: party_account_owner_candidates_mv, columns: [account_id, party_id, owner_type] } { output: [ party_account_owner_candidates_mv.account_id, party_account_owner_candidates_mv.party_id, party_account_owner_candidates_mv.owner_type ], stream key: [ party_account_owner_candidates_mv.account_id, party_account_owner_candidates_mv.party_id ] }
├── Upstream { output: [ account_id, party_id, owner_type ], stream key: [] }
└── BatchPlanNode { output: [ account_id, party_id, owner_type ], stream key: [] }
Fragment 54210 (Actor 738665,738664)
StreamProject { exprs: [parties.id, parties.type, parties.display_name] } { output: [ parties.id, parties.type, parties.display_name ], stream key: [ parties.id ] }
└── StreamFilter { predicate: IsNull(parties.disabled_at) } { output: [ parties.id, parties.type, parties.display_name, parties.disabled_at ], stream key: [ parties.id ] }
└── StreamTableScan { table: parties, columns: [id, type, display_name, disabled_at] } { output: [ parties.id, parties.type, parties.display_name, parties.disabled_at ], stream key: [ parties.id ] }
├── Upstream { output: [ id, type, display_name, disabled_at ], stream key: [] }
└── BatchPlanNode { output: [ id, type, display_name, disabled_at ], stream key: [] }
Fragment 54211 (Actor 738643,738642)
StreamLocalityProvider { locality_columns: [reference_identifiers_next.entity_id] } { output: [ reference_identifiers_next.entity_id, reference_identifiers_next.value, reference_identifiers_next.id ], stream key: [ reference_identifiers_next.entity_id, reference_identifiers_next.id ] }
└── MergeExecutor { output: [ reference_identifiers_next.entity_id, reference_identifiers_next.value, reference_identifiers_next.id ], stream key: [ reference_identifiers_next.id ] }
Fragment 54212 (Actor 738646,738647)
StreamProject { exprs: [reference_identifiers_next.entity_id, reference_identifiers_next.value, reference_identifiers_next.id] } { output: [ reference_identifiers_next.entity_id, reference_identifiers_next.value, reference_identifiers_next.id ], stream key: [ reference_identifiers_next.id ] }
└── StreamFilter { predicate: (reference_identifiers_next.entity_type = 'party':Varchar) AND (reference_identifiers_next.key = 'CustomerIdentificationFileId':Varchar) AND IsNull(reference_identifiers_next.disabled_at) } { output: [ reference_identifiers_next.entity_id, reference_identifiers_next.value, reference_identifiers_next.id, reference_identifiers_next.entity_type, reference_identifiers_next.key, reference_identifiers_next.disabled_at ], stream key: [ reference_identifiers_next.id ] }
└── StreamTableScan { table: reference_identifiers_next, columns: [entity_id, value, id, entity_type, key, disabled_at] } { output: [ reference_identifiers_next.entity_id, reference_identifiers_next.value, reference_identifiers_next.id, reference_identifiers_next.entity_type, reference_identifiers_next.key, reference_identifiers_next.disabled_at ], stream key: [ reference_identifiers_next.id ] }
├── Upstream { output: [ entity_id, value, id, entity_type, key, disabled_at ], stream key: [] }
└── BatchPlanNode { output: [ entity_id, value, id, entity_type, key, disabled_at ], stream key: [] }