Job is idle — throughput ~0; structure shown.
Fragment 63279 (Actor 745789,745788)
StreamMaterialize { columns: [account_id, client_json], stream_key: [account_id], pk_columns: [account_id], pk_conflict: NoCheck }
├── output: [ party_account_owner_candidates_mv_next.account_id, $expr2 ]
├── stream key: [ party_account_owner_candidates_mv_next.account_id ]
└── StreamProject { exprs: [party_account_owner_candidates_mv_next.account_id, JsonbAccess(jsonb_agg($expr1 order_by(min(party_account_owner_candidates_mv_next.party_id) ASC)), 0:Int32) as $expr2] }
├── output: [ party_account_owner_candidates_mv_next.account_id, $expr2 ]
├── stream key: [ party_account_owner_candidates_mv_next.account_id ]
└── StreamHashAgg { group_key: [party_account_owner_candidates_mv_next.account_id], aggs: [jsonb_agg($expr1 order_by(min(party_account_owner_candidates_mv_next.party_id) ASC)), count] }
├── output: [ party_account_owner_candidates_mv_next.account_id, jsonb_agg($expr1 order_by(min(party_account_owner_candidates_mv_next.party_id) ASC)), count ]
├── stream key: [ party_account_owner_candidates_mv_next.account_id ]
└── StreamLocalityProvider { locality_columns: [party_account_owner_candidates_mv_next.account_id] }
├── output: [ party_account_owner_candidates_mv_next.account_id, min(party_account_owner_candidates_mv_next.party_id), $expr1, parties.id, reference_identifiers_next.id ]
├── stream key: [ party_account_owner_candidates_mv_next.account_id, parties.id, min(party_account_owner_candidates_mv_next.party_id), reference_identifiers_next.id ]
└── MergeExecutor
├── output: [ party_account_owner_candidates_mv_next.account_id, min(party_account_owner_candidates_mv_next.party_id), $expr1, parties.id, reference_identifiers_next.id ]
└── stream key: [ parties.id, min(party_account_owner_candidates_mv_next.party_id), party_account_owner_candidates_mv_next.account_id, reference_identifiers_next.id ]
Fragment 63280 (Actor 745828,745829)
StreamProject { exprs: [party_account_owner_candidates_mv_next.account_id, min(party_account_owner_candidates_mv_next.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_next.account_id, min(party_account_owner_candidates_mv_next.party_id), $expr1, parties.id, reference_identifiers_next.id ]
├── stream key: [ parties.id, min(party_account_owner_candidates_mv_next.party_id), party_account_owner_candidates_mv_next.account_id, reference_identifiers_next.id ]
└── MergeExecutor { output: [ party_account_owner_candidates_mv_next.account_id, min(party_account_owner_candidates_mv_next.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_next.party_id), party_account_owner_candidates_mv_next.account_id, reference_identifiers_next.id ] }
Fragment 63281 (Actor 745827,745826)
StreamSyncLogStore { output: [ party_account_owner_candidates_mv_next.account_id, min(party_account_owner_candidates_mv_next.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_next.party_id), party_account_owner_candidates_mv_next.account_id, reference_identifiers_next.id ] }
└── StreamHashJoin { type: LeftOuter, predicate: parties.id = reference_identifiers_next.entity_id } { output: [ party_account_owner_candidates_mv_next.account_id, min(party_account_owner_candidates_mv_next.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_next.party_id), party_account_owner_candidates_mv_next.account_id, reference_identifiers_next.id ] }
├── MergeExecutor { output: [ party_account_owner_candidates_mv_next.account_id, min(party_account_owner_candidates_mv_next.party_id), parties.id, parties.type, parties.display_name ], stream key: [ parties.id, min(party_account_owner_candidates_mv_next.party_id), party_account_owner_candidates_mv_next.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 63282 (Actor 745830,745831)
StreamLocalityProvider { locality_columns: [parties.id] } { output: [ party_account_owner_candidates_mv_next.account_id, min(party_account_owner_candidates_mv_next.party_id), parties.id, parties.type, parties.display_name ], stream key: [ parties.id, min(party_account_owner_candidates_mv_next.party_id), party_account_owner_candidates_mv_next.account_id ] }
└── MergeExecutor { output: [ party_account_owner_candidates_mv_next.account_id, min(party_account_owner_candidates_mv_next.party_id), parties.id, parties.type, parties.display_name ], stream key: [ min(party_account_owner_candidates_mv_next.party_id), party_account_owner_candidates_mv_next.account_id ] }
Fragment 63283 (Actor 745832,745833)
StreamSyncLogStore { output: [ party_account_owner_candidates_mv_next.account_id, min(party_account_owner_candidates_mv_next.party_id), parties.id, parties.type, parties.display_name ], stream key: [ min(party_account_owner_candidates_mv_next.party_id), party_account_owner_candidates_mv_next.account_id ] }
└── StreamHashJoin { type: Inner, predicate: min(party_account_owner_candidates_mv_next.party_id) = parties.id } { output: [ party_account_owner_candidates_mv_next.account_id, min(party_account_owner_candidates_mv_next.party_id), parties.id, parties.type, parties.display_name ], stream key: [ min(party_account_owner_candidates_mv_next.party_id), party_account_owner_candidates_mv_next.account_id ] }
├── MergeExecutor { output: [ party_account_owner_candidates_mv_next.account_id, min(party_account_owner_candidates_mv_next.party_id), count ], stream key: [ min(party_account_owner_candidates_mv_next.party_id), party_account_owner_candidates_mv_next.account_id ] }
└── MergeExecutor { output: [ parties.id, parties.type, parties.display_name ], stream key: [ parties.id ] }
Fragment 63284 (Actor 745834,745835)
StreamLocalityProvider { locality_columns: [min(party_account_owner_candidates_mv_next.party_id)] } { output: [ party_account_owner_candidates_mv_next.account_id, min(party_account_owner_candidates_mv_next.party_id), count ], stream key: [ min(party_account_owner_candidates_mv_next.party_id), party_account_owner_candidates_mv_next.account_id ] }
└── MergeExecutor { output: [ party_account_owner_candidates_mv_next.account_id, min(party_account_owner_candidates_mv_next.party_id), count ], stream key: [ party_account_owner_candidates_mv_next.account_id ] }
Fragment 63285 (Actor 745836,745837)
StreamFilter { predicate: (count = 1:Int32) } { output: [ party_account_owner_candidates_mv_next.account_id, min(party_account_owner_candidates_mv_next.party_id), count ], stream key: [ party_account_owner_candidates_mv_next.account_id ] }
└── StreamHashAgg { group_key: [party_account_owner_candidates_mv_next.account_id], aggs: [min(party_account_owner_candidates_mv_next.party_id), count] } { output: [ party_account_owner_candidates_mv_next.account_id, min(party_account_owner_candidates_mv_next.party_id), count ], stream key: [ party_account_owner_candidates_mv_next.account_id ] }
└── StreamLocalityProvider { locality_columns: [party_account_owner_candidates_mv_next.account_id] } { output: [ party_account_owner_candidates_mv_next.account_id, party_account_owner_candidates_mv_next.party_id ], stream key: [ party_account_owner_candidates_mv_next.account_id, party_account_owner_candidates_mv_next.party_id ] }
└── MergeExecutor { output: [ party_account_owner_candidates_mv_next.account_id, party_account_owner_candidates_mv_next.party_id ], stream key: [ party_account_owner_candidates_mv_next.account_id, party_account_owner_candidates_mv_next.party_id ] }
Fragment 63286 (Actor 745841,745840)
StreamProject { exprs: [party_account_owner_candidates_mv_next.account_id, party_account_owner_candidates_mv_next.party_id] } { output: [ party_account_owner_candidates_mv_next.account_id, party_account_owner_candidates_mv_next.party_id ], stream key: [ party_account_owner_candidates_mv_next.account_id, party_account_owner_candidates_mv_next.party_id ] }
└── StreamFilter { predicate: (party_account_owner_candidates_mv_next.owner_type = 'PRIMARY':Varchar) } { output: [ party_account_owner_candidates_mv_next.account_id, party_account_owner_candidates_mv_next.party_id, party_account_owner_candidates_mv_next.owner_type ], stream key: [ party_account_owner_candidates_mv_next.account_id, party_account_owner_candidates_mv_next.party_id ] }
└── StreamTableScan { table: party_account_owner_candidates_mv_next, columns: [account_id, party_id, owner_type] } { output: [ party_account_owner_candidates_mv_next.account_id, party_account_owner_candidates_mv_next.party_id, party_account_owner_candidates_mv_next.owner_type ], stream key: [ party_account_owner_candidates_mv_next.account_id, party_account_owner_candidates_mv_next.party_id ] }
├── Upstream { output: [ account_id, party_id, owner_type ], stream key: [] }
└── BatchPlanNode { output: [ account_id, party_id, owner_type ], stream key: [] }
Fragment 63287 (Actor 745727,745726)
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 63288 (Actor 745838,745839)
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 63289 (Actor 745849,745848)
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: [] }