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

← cluster search objects items_mv explain
Overview Objects Graph History
materialized view · search.items_mv profiled over 5s
seconds (1–30)

Job is idle — throughput ~0; structure shown.

Stateful hash join (4 state tables) — consider a temporal join for dimension lookupsAggregation state — unbounded unless keyed or temporally filtered
441 operators
Materialize · search.items_mv
0% idle 2 actors
Project
2 actors
HashAgg Aggregation state — unbounded unless keyed or temporally filtered
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Union
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · party_identifier_edges_mv_next.owner_entity_id = party_refe…
2 actors
HashJoin · Inner · party_identifier_edges_mv_next.owner_entity_id = party_refe… Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · party_reference_identifier_terms_mv
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · party_identifier_edges_mv_next
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · active_identifier_edges_mv_next.owner_entity_id = olap_refe…
2 actors
HashJoin · Inner · active_identifier_edges_mv_next.owner_entity_id = olap_refe… Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · olap_reference_identifier_terms_mv
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Filter · active_identifier_edges_mv_next
0% idle 2 actors
StreamScan · active_identifier_edges_mv_next
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · clients_portfolios_dm.client_id = entity_to_teams_dm.entity…
2 actors
HashJoin · Inner · clients_portfolios_dm.client_id = entity_to_teams_dm.entity… Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · entity_to_teams_dm
2 actors
Filter · entity_to_teams_dm
0% idle 2 actors
StreamScan · entity_to_teams_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · clients_portfolios_dm
2 actors
Filter · clients_portfolios_dm
0% idle 2 actors
StreamScan · clients_portfolios_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · entity_to_teams_dm.entity_id = portfolios_dm.portfolio_id
2 actors
HashJoin · Inner · entity_to_teams_dm.entity_id = portfolios_dm.portfolio_id Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · portfolios_dm
2 actors
Filter · portfolios_dm
0% idle 2 actors
StreamScan · portfolios_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · entity_to_teams_dm
2 actors
Filter · entity_to_teams_dm
0% idle 2 actors
StreamScan · entity_to_teams_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · portfolios_dm.status_label_id = labels_dm.label_id
2 actors
HashJoin · Inner · portfolios_dm.status_label_id = labels_dm.label_id Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · labels_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · portfolios_dm
2 actors
Filter · portfolios_dm
0% idle 2 actors
StreamScan · portfolios_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · Inner · portfolios_dm.service_type_id = service_types_dm.service_ty…
2 actors
TemporalJoin · Inner · portfolios_dm.service_type_id = service_types_dm.service_ty…
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · service_types_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · portfolios_dm
2 actors
Filter · portfolios_dm
0% idle 2 actors
StreamScan · portfolios_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · accounts_to_clients_dm.client_id = entity_to_teams_dm.entit…
2 actors
HashJoin · Inner · accounts_to_clients_dm.client_id = entity_to_teams_dm.entit… Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · entity_to_teams_dm
2 actors
Filter · entity_to_teams_dm
0% idle 2 actors
StreamScan · entity_to_teams_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Filter · IsNull(accounts_to_clients_dm.disabled_at) AND IsNull(accou…
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · accounts_to_clients_dm.account_id = accounts_dm.account_id
2 actors
HashJoin · Inner · accounts_to_clients_dm.account_id = accounts_dm.account_id Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · accounts_dm
2 actors
Filter · accounts_dm
0% idle 2 actors
StreamScan · accounts_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Filter · accounts_to_clients_dm
0% idle 2 actors
StreamScan · accounts_to_clients_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · accounts_dm.status_label_id = labels_dm.label_id
2 actors
HashJoin · Inner · accounts_dm.status_label_id = labels_dm.label_id Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · labels_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · accounts_dm
2 actors
Filter · accounts_dm
0% idle 2 actors
StreamScan · accounts_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · IsNull(accounts_dm.disabled_at) AND Not(IsNull(product_type…
2 actors
Filter · IsNull(accounts_dm.disabled_at) AND Not(IsNull(product_type…
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
TemporalJoin · Inner · accounts_dm.product_type_id = product_types_dm.product_type…
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · product_types_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Filter · accounts_dm
0% idle 2 actors
StreamScan · accounts_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · IsNull(accounts_dm.disabled_at)
2 actors
Filter · IsNull(accounts_dm.disabled_at)
0% idle 2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · entity_to_teams_dm.entity_id = clients_dm.id
2 actors
HashJoin · Inner · entity_to_teams_dm.entity_id = clients_dm.id Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · clients_dm
2 actors
Filter · clients_dm
0% idle 2 actors
StreamScan · clients_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · entity_to_teams_dm
2 actors
Filter · entity_to_teams_dm
0% idle 2 actors
StreamScan · entity_to_teams_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · clients_dm.status_label_id = labels_dm.label_id
2 actors
HashJoin · Inner · clients_dm.status_label_id = labels_dm.label_id Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · labels_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · clients_dm
2 actors
Filter · clients_dm
0% idle 2 actors
StreamScan · clients_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · clients_dm
2 actors
Filter · clients_dm
0% idle 2 actors
StreamScan · clients_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · assets_dm
2 actors
Filter · assets_dm
0% idle 2 actors
StreamScan · assets_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · assets_dm
2 actors
ProjectSet · assets_dm
0% idle 2 actors
Project · assets_dm
2 actors
Filter · assets_dm
0% idle 2 actors
StreamScan · assets_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
2 actors
ProjectSet
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · clients_portfolios_dm.client_id = clients_contacts_dm.clien…
2 actors
HashJoin · Inner · clients_portfolios_dm.client_id = clients_contacts_dm.clien… Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · clients_contacts_dm
2 actors
Filter · clients_contacts_dm
0% idle 2 actors
StreamScan · clients_contacts_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · clients_portfolios_dm
2 actors
Filter · clients_portfolios_dm
0% idle 2 actors
StreamScan · clients_portfolios_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
2 actors
ProjectSet
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · clients_portfolios_dm.client_id = clients_dm.id
2 actors
HashJoin · Inner · clients_portfolios_dm.client_id = clients_dm.id Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · clients_dm
2 actors
Filter · clients_dm
0% idle 2 actors
StreamScan · clients_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · clients_portfolios_dm
2 actors
Filter · clients_portfolios_dm
0% idle 2 actors
StreamScan · clients_portfolios_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
2 actors
ProjectSet
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · account_to_portfolios_dm.account_id = accounts_dm.account_id
2 actors
HashJoin · Inner · account_to_portfolios_dm.account_id = accounts_dm.account_id Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · accounts_dm
2 actors
Filter · accounts_dm
0% idle 2 actors
StreamScan · accounts_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · account_to_portfolios_dm
2 actors
Filter · account_to_portfolios_dm
0% idle 2 actors
StreamScan · account_to_portfolios_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · portfolios_dm
2 actors
ProjectSet · portfolios_dm
0% idle 2 actors
Project · portfolios_dm
2 actors
Filter · portfolios_dm
0% idle 2 actors
StreamScan · portfolios_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · party_items_mv_next
2 actors
StreamScan · party_items_mv_next
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
2 actors
ProjectSet
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · clients_portfolios_dm.portfolio_id = portfolios_dm.portfoli…
2 actors
HashJoin · Inner · clients_portfolios_dm.portfolio_id = portfolios_dm.portfoli… Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · portfolios_dm
2 actors
Filter · portfolios_dm
0% idle 2 actors
StreamScan · portfolios_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · clients_portfolios_dm
2 actors
Filter · clients_portfolios_dm
0% idle 2 actors
StreamScan · clients_portfolios_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · IsNull(accounts_to_clients_dm.disabled_at)
2 actors
ProjectSet · IsNull(accounts_to_clients_dm.disabled_at)
0% idle 2 actors
Project · IsNull(accounts_to_clients_dm.disabled_at)
2 actors
Filter · IsNull(accounts_to_clients_dm.disabled_at)
0% idle 2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · clients_contacts_dm
2 actors
ProjectSet · clients_contacts_dm
0% idle 2 actors
Project · clients_contacts_dm
2 actors
Filter · clients_contacts_dm
0% idle 2 actors
StreamScan · clients_contacts_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
2 actors
ProjectSet
0% idle 2 actors
Project
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
2 actors
ProjectSet
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · accounts_to_clients_dm.client_id = clients_contacts_dm.clie…
2 actors
HashJoin · Inner · accounts_to_clients_dm.client_id = clients_contacts_dm.clie… Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · clients_contacts_dm
2 actors
Filter · clients_contacts_dm
0% idle 2 actors
StreamScan · clients_contacts_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · accounts_to_clients_dm
2 actors
Filter · accounts_to_clients_dm
0% idle 2 actors
StreamScan · accounts_to_clients_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
2 actors
ProjectSet
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · accounts_to_clients_dm.client_id = clients_dm.id
2 actors
HashJoin · Inner · accounts_to_clients_dm.client_id = clients_dm.id Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · clients_dm
2 actors
Filter · clients_dm
0% idle 2 actors
StreamScan · clients_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · accounts_to_clients_dm
2 actors
Filter · accounts_to_clients_dm
0% idle 2 actors
StreamScan · accounts_to_clients_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · accounts_dm
2 actors
ProjectSet · accounts_dm
0% idle 2 actors
Project · accounts_dm
2 actors
Filter · accounts_dm
0% idle 2 actors
StreamScan · accounts_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Heat = the operator's output-buffer backpressure over the sampling window. Click a node to fold its subtree.
Materialize · search.items_mv Materialize search.items_mv idle · 2 actors Project Project — · 2 actors HashAgg HashAgg idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Union Union idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · party_identifier_edges_mv_next.owner_entity_id = party_refe… SyncLogStore Inner · party_identifie… — · 2 actors HashJoin · Inner · party_identifier_edges_mv_next.owner_entity_id = party_refe… HashJoin Inner · party_identifie… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · party_reference_identifier_terms_mv StreamScan party_reference_identif… idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · party_identifier_edges_mv_next StreamScan party_identifier_edges_… idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · active_identifier_edges_mv_next.owner_entity_id = olap_refe… SyncLogStore Inner · active_identifi… — · 2 actors HashJoin · Inner · active_identifier_edges_mv_next.owner_entity_id = olap_refe… HashJoin Inner · active_identifi… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · olap_reference_identifier_terms_mv StreamScan olap_reference_identifi… idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Filter · active_identifier_edges_mv_next Filter active_identifier_edges… idle · 2 actors StreamScan · active_identifier_edges_mv_next StreamScan active_identifier_edges… idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · clients_portfolios_dm.client_id = entity_to_teams_dm.entity… SyncLogStore Inner · clients_portfol… — · 2 actors HashJoin · Inner · clients_portfolios_dm.client_id = entity_to_teams_dm.entity… HashJoin Inner · clients_portfol… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · entity_to_teams_dm Project entity_to_teams_dm — · 2 actors Filter · entity_to_teams_dm Filter entity_to_teams_dm idle · 2 actors StreamScan · entity_to_teams_dm StreamScan entity_to_teams_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · clients_portfolios_dm Project clients_portfolios_dm — · 2 actors Filter · clients_portfolios_dm Filter clients_portfolios_dm idle · 2 actors StreamScan · clients_portfolios_dm StreamScan clients_portfolios_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · entity_to_teams_dm.entity_id = portfolios_dm.portfolio_id SyncLogStore Inner · entity_to_teams… — · 2 actors HashJoin · Inner · entity_to_teams_dm.entity_id = portfolios_dm.portfolio_id HashJoin Inner · entity_to_teams… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · portfolios_dm Project portfolios_dm — · 2 actors Filter · portfolios_dm Filter portfolios_dm idle · 2 actors StreamScan · portfolios_dm StreamScan portfolios_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · entity_to_teams_dm Project entity_to_teams_dm — · 2 actors Filter · entity_to_teams_dm Filter entity_to_teams_dm idle · 2 actors StreamScan · entity_to_teams_dm StreamScan entity_to_teams_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · portfolios_dm.status_label_id = labels_dm.label_id SyncLogStore Inner · portfolios_dm.s… — · 2 actors HashJoin · Inner · portfolios_dm.status_label_id = labels_dm.label_id HashJoin Inner · portfolios_dm.s… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · labels_dm StreamScan labels_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · portfolios_dm Project portfolios_dm — · 2 actors Filter · portfolios_dm Filter portfolios_dm idle · 2 actors StreamScan · portfolios_dm StreamScan portfolios_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · Inner · portfolios_dm.service_type_id = service_types_dm.service_ty… Project Inner · portfolios_dm.s… — · 2 actors TemporalJoin · Inner · portfolios_dm.service_type_id = service_types_dm.service_ty… TemporalJoin Inner · portfolios_dm.s… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · service_types_dm StreamScan service_types_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · portfolios_dm Project portfolios_dm — · 2 actors Filter · portfolios_dm Filter portfolios_dm idle · 2 actors StreamScan · portfolios_dm StreamScan portfolios_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · accounts_to_clients_dm.client_id = entity_to_teams_dm.entit… SyncLogStore Inner · accounts_to_cli… — · 2 actors HashJoin · Inner · accounts_to_clients_dm.client_id = entity_to_teams_dm.entit… HashJoin Inner · accounts_to_cli… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · entity_to_teams_dm Project entity_to_teams_dm — · 2 actors Filter · entity_to_teams_dm Filter entity_to_teams_dm idle · 2 actors StreamScan · entity_to_teams_dm StreamScan entity_to_teams_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Filter · IsNull(accounts_to_clients_dm.disabled_at) AND IsNull(accou… Filter IsNull(accounts_to_clie… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore SyncLogStore — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · accounts_to_clients_dm.account_id = accounts_dm.account_id SyncLogStore Inner · accounts_to_cli… — · 2 actors HashJoin · Inner · accounts_to_clients_dm.account_id = accounts_dm.account_id HashJoin Inner · accounts_to_cli… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · accounts_dm Project accounts_dm — · 2 actors Filter · accounts_dm Filter accounts_dm idle · 2 actors StreamScan · accounts_dm StreamScan accounts_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Filter · accounts_to_clients_dm Filter accounts_to_clients_dm idle · 2 actors StreamScan · accounts_to_clients_dm StreamScan accounts_to_clients_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · accounts_dm.status_label_id = labels_dm.label_id SyncLogStore Inner · accounts_dm.sta… — · 2 actors HashJoin · Inner · accounts_dm.status_label_id = labels_dm.label_id HashJoin Inner · accounts_dm.sta… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · labels_dm StreamScan labels_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · accounts_dm Project accounts_dm — · 2 actors Filter · accounts_dm Filter accounts_dm idle · 2 actors StreamScan · accounts_dm StreamScan accounts_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · IsNull(accounts_dm.disabled_at) AND Not(IsNull(product_type… Project IsNull(accounts_dm.disa… — · 2 actors Filter · IsNull(accounts_dm.disabled_at) AND Not(IsNull(product_type… Filter IsNull(accounts_dm.disa… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors TemporalJoin · Inner · accounts_dm.product_type_id = product_types_dm.product_type… TemporalJoin Inner · accounts_dm.pro… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · product_types_dm StreamScan product_types_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Filter · accounts_dm Filter accounts_dm idle · 2 actors StreamScan · accounts_dm StreamScan accounts_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · IsNull(accounts_dm.disabled_at) Project IsNull(accounts_dm.disa… — · 2 actors Filter · IsNull(accounts_dm.disabled_at) Filter IsNull(accounts_dm.disa… idle · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · entity_to_teams_dm.entity_id = clients_dm.id SyncLogStore Inner · entity_to_teams… — · 2 actors HashJoin · Inner · entity_to_teams_dm.entity_id = clients_dm.id HashJoin Inner · entity_to_teams… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · clients_dm Project clients_dm — · 2 actors Filter · clients_dm Filter clients_dm idle · 2 actors StreamScan · clients_dm StreamScan clients_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · entity_to_teams_dm Project entity_to_teams_dm — · 2 actors Filter · entity_to_teams_dm Filter entity_to_teams_dm idle · 2 actors StreamScan · entity_to_teams_dm StreamScan entity_to_teams_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · clients_dm.status_label_id = labels_dm.label_id SyncLogStore Inner · clients_dm.stat… — · 2 actors HashJoin · Inner · clients_dm.status_label_id = labels_dm.label_id HashJoin Inner · clients_dm.stat… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · labels_dm StreamScan labels_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · clients_dm Project clients_dm — · 2 actors Filter · clients_dm Filter clients_dm idle · 2 actors StreamScan · clients_dm StreamScan clients_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · clients_dm Project clients_dm — · 2 actors Filter · clients_dm Filter clients_dm idle · 2 actors StreamScan · clients_dm StreamScan clients_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · assets_dm Project assets_dm — · 2 actors Filter · assets_dm Filter assets_dm idle · 2 actors StreamScan · assets_dm StreamScan assets_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · assets_dm Project assets_dm — · 2 actors ProjectSet · assets_dm ProjectSet assets_dm idle · 2 actors Project · assets_dm Project assets_dm — · 2 actors Filter · assets_dm Filter assets_dm idle · 2 actors StreamScan · assets_dm StreamScan assets_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 2 actors ProjectSet ProjectSet idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · clients_portfolios_dm.client_id = clients_contacts_dm.clien… SyncLogStore Inner · clients_portfol… — · 2 actors HashJoin · Inner · clients_portfolios_dm.client_id = clients_contacts_dm.clien… HashJoin Inner · clients_portfol… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · clients_contacts_dm Project clients_contacts_dm — · 2 actors Filter · clients_contacts_dm Filter clients_contacts_dm idle · 2 actors StreamScan · clients_contacts_dm StreamScan clients_contacts_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · clients_portfolios_dm Project clients_portfolios_dm — · 2 actors Filter · clients_portfolios_dm Filter clients_portfolios_dm idle · 2 actors StreamScan · clients_portfolios_dm StreamScan clients_portfolios_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 2 actors ProjectSet ProjectSet idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · clients_portfolios_dm.client_id = clients_dm.id SyncLogStore Inner · clients_portfol… — · 2 actors HashJoin · Inner · clients_portfolios_dm.client_id = clients_dm.id HashJoin Inner · clients_portfol… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · clients_dm Project clients_dm — · 2 actors Filter · clients_dm Filter clients_dm idle · 2 actors StreamScan · clients_dm StreamScan clients_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · clients_portfolios_dm Project clients_portfolios_dm — · 2 actors Filter · clients_portfolios_dm Filter clients_portfolios_dm idle · 2 actors StreamScan · clients_portfolios_dm StreamScan clients_portfolios_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 2 actors ProjectSet ProjectSet idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · account_to_portfolios_dm.account_id = accounts_dm.account_id SyncLogStore Inner · account_to_port… — · 2 actors HashJoin · Inner · account_to_portfolios_dm.account_id = accounts_dm.account_id HashJoin Inner · account_to_port… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · accounts_dm Project accounts_dm — · 2 actors Filter · accounts_dm Filter accounts_dm idle · 2 actors StreamScan · accounts_dm StreamScan accounts_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · account_to_portfolios_dm Project account_to_portfolios_dm — · 2 actors Filter · account_to_portfolios_dm Filter account_to_portfolios_dm idle · 2 actors StreamScan · account_to_portfolios_dm StreamScan account_to_portfolios_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · portfolios_dm Project portfolios_dm — · 2 actors ProjectSet · portfolios_dm ProjectSet portfolios_dm idle · 2 actors Project · portfolios_dm Project portfolios_dm — · 2 actors Filter · portfolios_dm Filter portfolios_dm idle · 2 actors StreamScan · portfolios_dm StreamScan portfolios_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · party_items_mv_next Project party_items_mv_next — · 2 actors StreamScan · party_items_mv_next StreamScan party_items_mv_next idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 2 actors ProjectSet ProjectSet idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · clients_portfolios_dm.portfolio_id = portfolios_dm.portfoli… SyncLogStore Inner · clients_portfol… — · 2 actors HashJoin · Inner · clients_portfolios_dm.portfolio_id = portfolios_dm.portfoli… HashJoin Inner · clients_portfol… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · portfolios_dm Project portfolios_dm — · 2 actors Filter · portfolios_dm Filter portfolios_dm idle · 2 actors StreamScan · portfolios_dm StreamScan portfolios_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · clients_portfolios_dm Project clients_portfolios_dm — · 2 actors Filter · clients_portfolios_dm Filter clients_portfolios_dm idle · 2 actors StreamScan · clients_portfolios_dm StreamScan clients_portfolios_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · IsNull(accounts_to_clients_dm.disabled_at) Project IsNull(accounts_to_clie… — · 2 actors ProjectSet · IsNull(accounts_to_clients_dm.disabled_at) ProjectSet IsNull(accounts_to_clie… idle · 2 actors Project · IsNull(accounts_to_clients_dm.disabled_at) Project IsNull(accounts_to_clie… — · 2 actors Filter · IsNull(accounts_to_clients_dm.disabled_at) Filter IsNull(accounts_to_clie… idle · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · clients_contacts_dm Project clients_contacts_dm — · 2 actors ProjectSet · clients_contacts_dm ProjectSet clients_contacts_dm idle · 2 actors Project · clients_contacts_dm Project clients_contacts_dm — · 2 actors Filter · clients_contacts_dm Filter clients_contacts_dm idle · 2 actors StreamScan · clients_contacts_dm StreamScan clients_contacts_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 2 actors ProjectSet ProjectSet idle · 2 actors Project Project — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 2 actors ProjectSet ProjectSet idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · accounts_to_clients_dm.client_id = clients_contacts_dm.clie… SyncLogStore Inner · accounts_to_cli… — · 2 actors HashJoin · Inner · accounts_to_clients_dm.client_id = clients_contacts_dm.clie… HashJoin Inner · accounts_to_cli… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · clients_contacts_dm Project clients_contacts_dm — · 2 actors Filter · clients_contacts_dm Filter clients_contacts_dm idle · 2 actors StreamScan · clients_contacts_dm StreamScan clients_contacts_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · accounts_to_clients_dm Project accounts_to_clients_dm — · 2 actors Filter · accounts_to_clients_dm Filter accounts_to_clients_dm idle · 2 actors StreamScan · accounts_to_clients_dm StreamScan accounts_to_clients_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 2 actors ProjectSet ProjectSet idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · accounts_to_clients_dm.client_id = clients_dm.id SyncLogStore Inner · accounts_to_cli… — · 2 actors HashJoin · Inner · accounts_to_clients_dm.client_id = clients_dm.id HashJoin Inner · accounts_to_cli… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · clients_dm Project clients_dm — · 2 actors Filter · clients_dm Filter clients_dm idle · 2 actors StreamScan · clients_dm StreamScan clients_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · accounts_to_clients_dm Project accounts_to_clients_dm — · 2 actors Filter · accounts_to_clients_dm Filter accounts_to_clients_dm idle · 2 actors StreamScan · accounts_to_clients_dm StreamScan accounts_to_clients_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · accounts_dm Project accounts_dm — · 2 actors ProjectSet · accounts_dm ProjectSet accounts_dm idle · 2 actors Project · accounts_dm Project accounts_dm — · 2 actors Filter · accounts_dm Filter accounts_dm idle · 2 actors StreamScan · accounts_dm StreamScan accounts_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors
Streaming operator plan from EXPLAIN ANALYZE. Node heat = backpressure. Drag to pan, scroll to zoom.
Fragments (DESCRIBE FRAGMENTS) — click to expand
Fragment 24887 (Actor 119453,119454)
StreamMaterialize { columns: [entity_id, entity_type, search_term_en, search_term_ar, filters], stream_key: [entity_id, entity_type], pk_columns: [entity_id, entity_type], pk_conflict: NoCheck }
├── output: [ accounts_dm.account_id, 'account':Varchar, array_agg(distinct Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar)))) filter(Not(IsNull(Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))))) AND (Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))) <> '':Varchar) AND (Trim(Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar)))) <> '':Varchar)), array_agg(distinct null:Varchar) filter(Not(IsNull(null:Varchar)) AND (null:Varchar <> '':Varchar) AND (Trim(null:Varchar) <> '':Varchar)), array_agg(distinct null:Varchar) filter(Not(IsNull(null:Varchar))) ]
├── stream key: [ accounts_dm.account_id, 'account':Varchar ]
└── StreamProject { exprs: [accounts_dm.account_id, 'account':Varchar, array_agg(distinct Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar)))) filter(Not(IsNull(Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))))) AND (Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))) <> '':Varchar) AND (Trim(Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar)))) <> '':Varchar)), array_agg(distinct null:Varchar) filter(Not(IsNull(null:Varchar)) AND (null:Varchar <> '':Varchar) AND (Trim(null:Varchar) <> '':Varchar)), array_agg(distinct null:Varchar) filter(Not(IsNull(null:Varchar)))] }
    ├── output: [ accounts_dm.account_id, 'account':Varchar, array_agg(distinct Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar)))) filter(Not(IsNull(Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))))) AND (Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))) <> '':Varchar) AND (Trim(Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar)))) <> '':Varchar)), array_agg(distinct null:Varchar) filter(Not(IsNull(null:Varchar)) AND (null:Varchar <> '':Varchar) AND (Trim(null:Varchar) <> '':Varchar)), array_agg(distinct null:Varchar) filter(Not(IsNull(null:Varchar))) ]
    ├── stream key: [ accounts_dm.account_id, 'account':Varchar ]
    └── StreamHashAgg { group_key: [accounts_dm.account_id, 'account':Varchar], aggs: [array_agg(distinct Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar)))) filter(Not(IsNull(Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))))) AND (Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))) <> '':Varchar) AND (Trim(Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar)))) <> '':Varchar)), array_agg(distinct null:Varchar) filter(Not(IsNull(null:Varchar)) AND (null:Varchar <> '':Varchar) AND (Trim(null:Varchar) <> '':Varchar)), array_agg(distinct null:Varchar) filter(Not(IsNull(null:Varchar))), count] }
        ├── output: [ accounts_dm.account_id, 'account':Varchar, array_agg(distinct Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar)))) filter(Not(IsNull(Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))))) AND (Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))) <> '':Varchar) AND (Trim(Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar)))) <> '':Varchar)), array_agg(distinct null:Varchar) filter(Not(IsNull(null:Varchar)) AND (null:Varchar <> '':Varchar) AND (Trim(null:Varchar) <> '':Varchar)), array_agg(distinct null:Varchar) filter(Not(IsNull(null:Varchar))), count ]
        ├── stream key: [ accounts_dm.account_id, 'account':Varchar ]
        └── MergeExecutor { output: [ accounts_dm.account_id, 'account':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), null:Varchar, null:Varchar, accounts_dm.account_id, null:Varchar, null:Varchar, null:Varchar, null:Date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, $src ], stream key: [ accounts_dm.account_id, _rw_projected_row_id, null:Varchar, null:Date, null:Varchar, null:Varchar, null:Int32, null:Int64, null:Int32, null:Int32, null:Date, $src ] }

Fragment 24888 (Actor 119455,119456)
StreamUnion { all: true } { output: [ accounts_dm.account_id, 'account':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), null:Varchar, null:Varchar, accounts_dm.account_id, null:Varchar, null:Varchar, null:Varchar, null:Date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, $src ], stream key: [ accounts_dm.account_id, _rw_projected_row_id, null:Varchar, null:Date, null:Varchar, null:Varchar, null:Int32, null:Int64, null:Int32, null:Int32, null:Date, $src ] }
├── MergeExecutor { output: [ accounts_dm.account_id, 'account':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), null:Varchar, null:Varchar, accounts_dm.account_id, null:Varchar, null:Varchar, null:Varchar, null:Date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 0:Int32 ], stream key: [ accounts_dm.account_id, _rw_projected_row_id ] }
├── MergeExecutor
│   ├── output: [ accounts_to_clients_dm.account_id, 'account':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar), $3, Replace($3, '-':Varchar, '':Varchar), $4, Replace($4, '-':Varchar, '':Varchar))), Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar), $3, Replace($3, '-':Varchar, '':Varchar), $4, Replace($4, '-':Varchar, '':Varchar))), null:Varchar, accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, null:Varchar, null:Varchar, accounts_to_clients_dm.effective_start_date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 1:Int32 ]
│   └── stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, _rw_projected_row_id ]
├── MergeExecutor { output: [ accounts_to_clients_dm.account_id, 'account':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar))), null:Varchar, null:Varchar, accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, clients_contacts_dm.type, clients_contacts_dm.value, accounts_to_clients_dm.effective_start_date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 2:Int32 ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, clients_contacts_dm.type, clients_contacts_dm.value, _rw_projected_row_id ] }
├── MergeExecutor { output: [ clients_dm.id, 'client':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar), $3, Replace($3, '-':Varchar, '':Varchar), $4, Replace($4, '-':Varchar, '':Varchar))), Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar), $3, Replace($3, '-':Varchar, '':Varchar), $4, Replace($4, '-':Varchar, '':Varchar))), null:Varchar, clients_dm.id, null:Varchar, null:Varchar, null:Varchar, null:Date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 3:Int32 ], stream key: [ clients_dm.id, _rw_projected_row_id ] }
├── MergeExecutor { output: [ clients_contacts_dm.client_id, 'client':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar))), null:Varchar, null:Varchar, clients_contacts_dm.client_id, clients_contacts_dm.type, clients_contacts_dm.value, null:Varchar, null:Date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 4:Int32 ], stream key: [ clients_contacts_dm.client_id, clients_contacts_dm.type, clients_contacts_dm.value, _rw_projected_row_id ] }
├── MergeExecutor { output: [ accounts_to_clients_dm.client_id, 'client':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), null:Varchar, accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, null:Varchar, null:Varchar, accounts_to_clients_dm.effective_start_date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 5:Int32 ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, _rw_projected_row_id ] }
├── MergeExecutor { output: [ clients_portfolios_dm.client_id, 'client':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), null:Varchar, clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, null:Varchar, null:Varchar, clients_portfolios_dm.effective_start_date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 6:Int32 ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, _rw_projected_row_id ] }
├── MergeExecutor
│   ├── output: [ party_items_mv_next.entity_id, party_items_mv_next.entity_type, party_items_mv_next.val, party_items_mv_next.ar_val, party_items_mv_next.filter_val, party_items_mv_next.active_parties_mv.id, party_items_mv_next.null:Varchar, party_items_mv_next.null:Varchar#1, party_items_mv_next.null:Varchar#2, party_items_mv_next.null:Date, null:Date, party_items_mv_next._rw_projected_row_id, party_items_mv_next.null:Int64, party_items_mv_next.null:Int32, party_items_mv_next.null:Int32#1, party_items_mv_next.$src, 7:Int32 ]
│   └── stream key: [ party_items_mv_next.active_parties_mv.id, party_items_mv_next._rw_projected_row_id, party_items_mv_next.null:Varchar, party_items_mv_next.null:Int32, party_items_mv_next.null:Int64, party_items_mv_next.null:Date, party_items_mv_next.null:Int32#1, party_items_mv_next.null:Varchar#1, party_items_mv_next.null:Varchar#2, party_items_mv_next.$src ]
├── MergeExecutor { output: [ portfolios_dm.portfolio_id, 'portfolio':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), null:Varchar, portfolios_dm.portfolio_id, null:Varchar, null:Varchar, null:Varchar, null:Date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 8:Int32 ], stream key: [ portfolios_dm.portfolio_id, _rw_projected_row_id ] }
├── MergeExecutor { output: [ account_to_portfolios_dm.portfolio_id, 'portfolio':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), null:Varchar, account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, null:Varchar, null:Varchar, account_to_portfolios_dm.effective_start_date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 9:Int32 ], stream key: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date, _rw_projected_row_id ] }
├── MergeExecutor
│   ├── output: [ clients_portfolios_dm.portfolio_id, 'portfolio':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar), $3, Replace($3, '-':Varchar, '':Varchar), $4, Replace($4, '-':Varchar, '':Varchar))), Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar), $3, Replace($3, '-':Varchar, '':Varchar), $4, Replace($4, '-':Varchar, '':Varchar))), null:Varchar, clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, null:Varchar, null:Varchar, clients_portfolios_dm.effective_start_date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 10:Int32 ]
│   └── stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, _rw_projected_row_id ]
├── MergeExecutor { output: [ clients_portfolios_dm.portfolio_id, 'portfolio':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar))), null:Varchar, null:Varchar, clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_contacts_dm.type, clients_contacts_dm.value, clients_portfolios_dm.effective_start_date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 11:Int32 ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, clients_contacts_dm.type, clients_contacts_dm.value, _rw_projected_row_id ] }
├── MergeExecutor { output: [ assets_dm.id, 'asset':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $3, Replace($3, '-':Varchar, '':Varchar), $4, Replace($4, '-':Varchar, '':Varchar), $7, Replace($7, '-':Varchar, '':Varchar), $5, Replace($5, '-':Varchar, '':Varchar), $6, Replace($6, '-':Varchar, '':Varchar), $8, Replace($8, '-':Varchar, '':Varchar))), null:Varchar, $expr1, assets_dm.id, null:Varchar, null:Varchar, null:Varchar, null:Date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 12:Int32 ], stream key: [ assets_dm.id, _rw_projected_row_id ] }
├── MergeExecutor { output: [ assets_dm.id, 'asset':Varchar, null:Varchar, assets_dm.name_ar, $expr2, assets_dm.id, null:Varchar, null:Varchar, null:Varchar, null:Date, null:Date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 13:Int32 ], stream key: [ assets_dm.id ] }
├── MergeExecutor { output: [ clients_dm.id, 'client':Varchar, null:Varchar, null:Varchar, $expr3, clients_dm.id, null:Varchar, null:Varchar, null:Varchar, null:Date, null:Date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 14:Int32 ], stream key: [ clients_dm.id ] }
├── MergeExecutor { output: [ clients_dm.id, 'client':Varchar, null:Varchar, null:Varchar, $expr4, clients_dm.id, clients_dm.status_label_id, null:Varchar, null:Varchar, null:Date, null:Date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 15:Int32 ], stream key: [ clients_dm.id, clients_dm.status_label_id ] }
├── MergeExecutor { output: [ entity_to_teams_dm.entity_id, 'client':Varchar, null:Varchar, null:Varchar, $expr5, entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, null:Varchar, entity_to_teams_dm.effective_start_date, null:Date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 16:Int32 ], stream key: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }
├── MergeExecutor { output: [ accounts_dm.account_id, 'account':Varchar, null:Varchar, null:Varchar, $expr6, accounts_dm.account_id, accounts_dm.product_type_id, null:Varchar, null:Varchar, null:Date, null:Date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 17:Int32 ], stream key: [ accounts_dm.account_id, accounts_dm.product_type_id ] }
├── MergeExecutor { output: [ accounts_dm.account_id, 'account':Varchar, null:Varchar, null:Varchar, $expr7, accounts_dm.account_id, accounts_dm.product_type_id, null:Varchar, null:Varchar, null:Date, null:Date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 18:Int32 ], stream key: [ accounts_dm.account_id, accounts_dm.product_type_id ] }
├── MergeExecutor { output: [ accounts_dm.account_id, 'account':Varchar, null:Varchar, null:Varchar, $expr8, accounts_dm.account_id, accounts_dm.status_label_id, null:Varchar, null:Varchar, null:Date, null:Date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 19:Int32 ], stream key: [ accounts_dm.account_id, accounts_dm.status_label_id ] }
├── MergeExecutor { output: [ accounts_to_clients_dm.account_id, 'account':Varchar, null:Varchar, null:Varchar, $expr9, accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, entity_to_teams_dm.team_id, entity_to_teams_dm.entity_type, accounts_to_clients_dm.effective_start_date, entity_to_teams_dm.effective_start_date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 20:Int32 ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, entity_to_teams_dm.team_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }
├── MergeExecutor { output: [ portfolios_dm.portfolio_id, 'portfolio':Varchar, null:Varchar, null:Varchar, $expr10, portfolios_dm.portfolio_id, portfolios_dm.service_type_id, null:Varchar, null:Varchar, null:Date, null:Date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 21:Int32 ], stream key: [ portfolios_dm.portfolio_id, portfolios_dm.service_type_id ] }
├── MergeExecutor { output: [ portfolios_dm.portfolio_id, 'portfolio':Varchar, null:Varchar, null:Varchar, $expr11, portfolios_dm.portfolio_id, portfolios_dm.status_label_id, null:Varchar, null:Varchar, null:Date, null:Date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 22:Int32 ], stream key: [ portfolios_dm.portfolio_id, portfolios_dm.status_label_id ] }
├── MergeExecutor { output: [ entity_to_teams_dm.entity_id, 'portfolio':Varchar, null:Varchar, null:Varchar, $expr12, entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, null:Varchar, entity_to_teams_dm.effective_start_date, null:Date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 23:Int32 ], stream key: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }
├── MergeExecutor { output: [ clients_portfolios_dm.portfolio_id, 'portfolio':Varchar, null:Varchar, null:Varchar, $expr13, clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, entity_to_teams_dm.team_id, entity_to_teams_dm.entity_type, clients_portfolios_dm.effective_start_date, entity_to_teams_dm.effective_start_date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 24:Int32 ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, entity_to_teams_dm.team_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }
├── MergeExecutor
│   ├── output: [ active_identifier_edges_mv_next.target_entity_id, active_identifier_edges_mv_next.target_entity_type, olap_reference_identifier_terms_mv.val, olap_reference_identifier_terms_mv.ar_val, null:Varchar, active_identifier_edges_mv_next.accounts_dm.account_id, active_identifier_edges_mv_next.null:Varchar, active_identifier_edges_mv_next.owner_entity_id, active_identifier_edges_mv_next.owner_entity_type, active_identifier_edges_mv_next.null:Date, null:Date, olap_reference_identifier_terms_mv.reference_identifiers.id, olap_reference_identifier_terms_mv._rw_projected_row_id, active_identifier_edges_mv_next.null:Int32, active_identifier_edges_mv_next.$src, null:Int32, 25:Int32 ]
│   └── stream key: [ active_identifier_edges_mv_next.accounts_dm.account_id, active_identifier_edges_mv_next.null:Varchar, active_identifier_edges_mv_next.null:Date, active_identifier_edges_mv_next.null:Int32, active_identifier_edges_mv_next.$src, olap_reference_identifier_terms_mv.reference_identifiers.id, olap_reference_identifier_terms_mv._rw_projected_row_id, active_identifier_edges_mv_next.owner_entity_id, active_identifier_edges_mv_next.owner_entity_type ]
└── MergeExecutor
    ├── output: [ party_identifier_edges_mv_next.target_entity_id, party_identifier_edges_mv_next.target_entity_type, party_reference_identifier_terms_mv.val, party_reference_identifier_terms_mv.ar_val, null:Varchar, party_identifier_edges_mv_next.owner_entity_id, party_identifier_edges_mv_next.target_entity_id, party_identifier_edges_mv_next.owner_entity_type, null:Varchar, null:Date, null:Date, party_reference_identifier_terms_mv.reference_identifiers_next.id, party_reference_identifier_terms_mv._rw_projected_row_id, party_identifier_edges_mv_next.party_holder_edges_mv_next.$src, null:Int32, null:Int32, 26:Int32 ]
    └── stream key: [ party_identifier_edges_mv_next.owner_entity_id, party_identifier_edges_mv_next.target_entity_id, party_identifier_edges_mv_next.party_holder_edges_mv_next.$src, party_reference_identifier_terms_mv.reference_identifiers_next.id, party_reference_identifier_terms_mv._rw_projected_row_id, party_identifier_edges_mv_next.owner_entity_type ]

Fragment 24889 (Actor 119457,119458)
StreamProject { exprs: [accounts_dm.account_id, 'account':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), null:Varchar, null:Varchar, accounts_dm.account_id, null:Varchar, null:Varchar, null:Varchar, null:Date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 0:Int32] } { output: [ accounts_dm.account_id, 'account':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), null:Varchar, null:Varchar, accounts_dm.account_id, null:Varchar, null:Varchar, null:Varchar, null:Date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 0:Int32 ], stream key: [ accounts_dm.account_id, _rw_projected_row_id ] }
└── StreamProjectSet { select_list: [$0, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar)))] } { output: [ _rw_projected_row_id, accounts_dm.account_id, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))) ], stream key: [ accounts_dm.account_id, _rw_projected_row_id ] }
    └── StreamProject { exprs: [accounts_dm.account_id, accounts_dm.name, accounts_dm.number] } { output: [ accounts_dm.account_id, accounts_dm.name, accounts_dm.number ], stream key: [ accounts_dm.account_id ] }
        └── StreamFilter { predicate: IsNull(accounts_dm.disabled_at) } { output: [ accounts_dm.account_id, accounts_dm.name, accounts_dm.number, accounts_dm.disabled_at ], stream key: [ accounts_dm.account_id ] }
            └── StreamTableScan { table: accounts_dm, columns: [account_id, name, number, disabled_at] } { output: [ accounts_dm.account_id, accounts_dm.name, accounts_dm.number, accounts_dm.disabled_at ], stream key: [ accounts_dm.account_id ] }
                ├── Upstream { output: [ account_id, name, number, disabled_at ], stream key: [] }
                └── BatchPlanNode { output: [ account_id, name, number, disabled_at ], stream key: [] }

Fragment 24890 (Actor 119461,119462)
StreamProject { exprs: [accounts_to_clients_dm.account_id, 'account':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar), $3, Replace($3, '-':Varchar, '':Varchar), $4, Replace($4, '-':Varchar, '':Varchar))), Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar), $3, Replace($3, '-':Varchar, '':Varchar), $4, Replace($4, '-':Varchar, '':Varchar))), null:Varchar, accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, null:Varchar, null:Varchar, accounts_to_clients_dm.effective_start_date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 1:Int32] }
├── output: [ accounts_to_clients_dm.account_id, 'account':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar), $3, Replace($3, '-':Varchar, '':Varchar), $4, Replace($4, '-':Varchar, '':Varchar))), Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar), $3, Replace($3, '-':Varchar, '':Varchar), $4, Replace($4, '-':Varchar, '':Varchar))), null:Varchar, accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, null:Varchar, null:Varchar, accounts_to_clients_dm.effective_start_date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 1:Int32 ]
├── stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, _rw_projected_row_id ]
└── StreamProjectSet { select_list: [$0, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar), $3, Replace($3, '-':Varchar, '':Varchar), $4, Replace($4, '-':Varchar, '':Varchar))), $5, $6] } { output: [ _rw_projected_row_id, accounts_to_clients_dm.account_id, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar), $3, Replace($3, '-':Varchar, '':Varchar), $4, Replace($4, '-':Varchar, '':Varchar))), accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, _rw_projected_row_id ] }
    └── MergeExecutor { output: [ accounts_to_clients_dm.account_id, clients_dm.display_name, clients_dm.local_display_name, clients_dm.preferred_name, clients_dm.customer_identification_file, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, clients_dm.id ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ] }

Fragment 24891 (Actor 119460,119459)
StreamSyncLogStore { output: [ accounts_to_clients_dm.account_id, clients_dm.display_name, clients_dm.local_display_name, clients_dm.preferred_name, clients_dm.customer_identification_file, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, clients_dm.id ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ] }
└── StreamHashJoin { type: Inner, predicate: accounts_to_clients_dm.client_id = clients_dm.id } { output: [ accounts_to_clients_dm.account_id, clients_dm.display_name, clients_dm.local_display_name, clients_dm.preferred_name, clients_dm.customer_identification_file, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, clients_dm.id ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ] }
    ├── MergeExecutor { output: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ] }
    └── MergeExecutor { output: [ clients_dm.id, clients_dm.display_name, clients_dm.local_display_name, clients_dm.preferred_name, clients_dm.customer_identification_file ], stream key: [ clients_dm.id ] }

Fragment 24892 (Actor 119463,119464)
StreamProject { exprs: [accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date] } { output: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ] }
└── StreamFilter { predicate: IsNull(accounts_to_clients_dm.disabled_at) } { output: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, accounts_to_clients_dm.disabled_at ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ] }
    └── StreamTableScan { table: accounts_to_clients_dm, columns: [account_id, client_id, effective_start_date, disabled_at] } { output: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, accounts_to_clients_dm.disabled_at ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ] }
        ├── Upstream { output: [ account_id, client_id, effective_start_date, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ account_id, client_id, effective_start_date, disabled_at ], stream key: [] }

Fragment 24893 (Actor 119466,119465)
StreamProject { exprs: [clients_dm.id, clients_dm.display_name, clients_dm.local_display_name, clients_dm.preferred_name, clients_dm.customer_identification_file] } { output: [ clients_dm.id, clients_dm.display_name, clients_dm.local_display_name, clients_dm.preferred_name, clients_dm.customer_identification_file ], stream key: [ clients_dm.id ] }
└── StreamFilter { predicate: IsNull(clients_dm.closing_date) } { output: [ clients_dm.id, clients_dm.display_name, clients_dm.local_display_name, clients_dm.preferred_name, clients_dm.customer_identification_file, clients_dm.closing_date ], stream key: [ clients_dm.id ] }
    └── StreamTableScan { table: clients_dm, columns: [id, display_name, local_display_name, preferred_name, customer_identification_file, closing_date] } { output: [ clients_dm.id, clients_dm.display_name, clients_dm.local_display_name, clients_dm.preferred_name, clients_dm.customer_identification_file, clients_dm.closing_date ], stream key: [ clients_dm.id ] }
        ├── Upstream { output: [ id, display_name, local_display_name, preferred_name, customer_identification_file, closing_date ], stream key: [] }
        └── BatchPlanNode { output: [ id, display_name, local_display_name, preferred_name, customer_identification_file, closing_date ], stream key: [] }

Fragment 24894 (Actor 119468,119467)
StreamProject { exprs: [accounts_to_clients_dm.account_id, 'account':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar))), null:Varchar, null:Varchar, accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, clients_contacts_dm.type, clients_contacts_dm.value, accounts_to_clients_dm.effective_start_date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 2:Int32] }
├── output: [ accounts_to_clients_dm.account_id, 'account':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar))), null:Varchar, null:Varchar, accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, clients_contacts_dm.type, clients_contacts_dm.value, accounts_to_clients_dm.effective_start_date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 2:Int32 ]
├── stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, clients_contacts_dm.type, clients_contacts_dm.value, _rw_projected_row_id ]
└── StreamProjectSet { select_list: [$0, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar))), $2, $3, $5, $1] } { output: [ _rw_projected_row_id, accounts_to_clients_dm.account_id, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar))), accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, clients_contacts_dm.type, clients_contacts_dm.value ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, clients_contacts_dm.type, clients_contacts_dm.value, _rw_projected_row_id ] }
    └── MergeExecutor { output: [ accounts_to_clients_dm.account_id, clients_contacts_dm.value, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, clients_contacts_dm.client_id, clients_contacts_dm.type ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, clients_contacts_dm.type, clients_contacts_dm.value ] }

Fragment 24895 (Actor 119469,119470)
StreamSyncLogStore { output: [ accounts_to_clients_dm.account_id, clients_contacts_dm.value, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, clients_contacts_dm.client_id, clients_contacts_dm.type ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, clients_contacts_dm.type, clients_contacts_dm.value ] }
└── StreamHashJoin { type: Inner, predicate: accounts_to_clients_dm.client_id = clients_contacts_dm.client_id } { output: [ accounts_to_clients_dm.account_id, clients_contacts_dm.value, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, clients_contacts_dm.client_id, clients_contacts_dm.type ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, clients_contacts_dm.type, clients_contacts_dm.value ] }
    ├── MergeExecutor { output: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ] }
    └── MergeExecutor { output: [ clients_contacts_dm.client_id, clients_contacts_dm.value, clients_contacts_dm.type ], stream key: [ clients_contacts_dm.client_id, clients_contacts_dm.type, clients_contacts_dm.value ] }

Fragment 24896 (Actor 119471,119472)
StreamProject { exprs: [accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date] } { output: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ] }
└── StreamFilter { predicate: IsNull(accounts_to_clients_dm.disabled_at) } { output: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, accounts_to_clients_dm.disabled_at ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ] }
    └── StreamTableScan { table: accounts_to_clients_dm, columns: [account_id, client_id, effective_start_date, disabled_at] } { output: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, accounts_to_clients_dm.disabled_at ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ] }
        ├── Upstream { output: [ account_id, client_id, effective_start_date, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ account_id, client_id, effective_start_date, disabled_at ], stream key: [] }

Fragment 24897 (Actor 119473,119474)
StreamProject { exprs: [clients_contacts_dm.client_id, clients_contacts_dm.value, clients_contacts_dm.type] } { output: [ clients_contacts_dm.client_id, clients_contacts_dm.value, clients_contacts_dm.type ], stream key: [ clients_contacts_dm.client_id, clients_contacts_dm.type, clients_contacts_dm.value ] }
└── StreamFilter { predicate: IsNull(clients_contacts_dm.disabled_at) } { output: [ clients_contacts_dm.client_id, clients_contacts_dm.value, clients_contacts_dm.type, clients_contacts_dm.disabled_at ], stream key: [ clients_contacts_dm.client_id, clients_contacts_dm.type, clients_contacts_dm.value ] }
    └── StreamTableScan { table: clients_contacts_dm, columns: [client_id, value, type, disabled_at] } { output: [ clients_contacts_dm.client_id, clients_contacts_dm.value, clients_contacts_dm.type, clients_contacts_dm.disabled_at ], stream key: [ clients_contacts_dm.client_id, clients_contacts_dm.type, clients_contacts_dm.value ] }
        ├── Upstream { output: [ client_id, value, type, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ client_id, value, type, disabled_at ], stream key: [] }

Fragment 24898 (Actor 119479,119480)
StreamProject { exprs: [clients_dm.id, 'client':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar), $3, Replace($3, '-':Varchar, '':Varchar), $4, Replace($4, '-':Varchar, '':Varchar))), Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar), $3, Replace($3, '-':Varchar, '':Varchar), $4, Replace($4, '-':Varchar, '':Varchar))), null:Varchar, clients_dm.id, null:Varchar, null:Varchar, null:Varchar, null:Date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 3:Int32] }
├── output: [ clients_dm.id, 'client':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar), $3, Replace($3, '-':Varchar, '':Varchar), $4, Replace($4, '-':Varchar, '':Varchar))), Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar), $3, Replace($3, '-':Varchar, '':Varchar), $4, Replace($4, '-':Varchar, '':Varchar))), null:Varchar, clients_dm.id, null:Varchar, null:Varchar, null:Varchar, null:Date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 3:Int32 ]
├── stream key: [ clients_dm.id, _rw_projected_row_id ]
└── StreamProjectSet { select_list: [$0, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar), $3, Replace($3, '-':Varchar, '':Varchar), $4, Replace($4, '-':Varchar, '':Varchar)))] } { output: [ _rw_projected_row_id, clients_dm.id, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar), $3, Replace($3, '-':Varchar, '':Varchar), $4, Replace($4, '-':Varchar, '':Varchar))) ], stream key: [ clients_dm.id, _rw_projected_row_id ] }
    └── StreamProject { exprs: [clients_dm.id, clients_dm.display_name, clients_dm.local_display_name, clients_dm.preferred_name, clients_dm.customer_identification_file] } { output: [ clients_dm.id, clients_dm.display_name, clients_dm.local_display_name, clients_dm.preferred_name, clients_dm.customer_identification_file ], stream key: [ clients_dm.id ] }
        └── MergeExecutor { output: [ clients_dm.id, clients_dm.display_name, clients_dm.local_display_name, clients_dm.preferred_name, clients_dm.customer_identification_file, clients_dm.type ], stream key: [ clients_dm.id ] }

Fragment 24899 (Actor 119476,119475)
StreamProject { exprs: [clients_dm.id, clients_dm.display_name, clients_dm.local_display_name, clients_dm.preferred_name, clients_dm.customer_identification_file, clients_dm.type] } { output: [ clients_dm.id, clients_dm.display_name, clients_dm.local_display_name, clients_dm.preferred_name, clients_dm.customer_identification_file, clients_dm.type ], stream key: [ clients_dm.id ] }
└── StreamFilter { predicate: IsNull(clients_dm.closing_date) } { output: [ clients_dm.id, clients_dm.display_name, clients_dm.local_display_name, clients_dm.preferred_name, clients_dm.customer_identification_file, clients_dm.type, clients_dm.closing_date ], stream key: [ clients_dm.id ] }
    └── StreamTableScan { table: clients_dm, columns: [id, display_name, local_display_name, preferred_name, customer_identification_file, type, closing_date] } { output: [ clients_dm.id, clients_dm.display_name, clients_dm.local_display_name, clients_dm.preferred_name, clients_dm.customer_identification_file, clients_dm.type, clients_dm.closing_date ], stream key: [ clients_dm.id ] }
        ├── Upstream { output: [ id, display_name, local_display_name, preferred_name, customer_identification_file, type, closing_date ], stream key: [] }
        └── BatchPlanNode { output: [ id, display_name, local_display_name, preferred_name, customer_identification_file, type, closing_date ], stream key: [] }

Fragment 24900 (Actor 119482,119481)
StreamProject { exprs: [clients_contacts_dm.client_id, 'client':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar))), null:Varchar, null:Varchar, clients_contacts_dm.client_id, clients_contacts_dm.type, clients_contacts_dm.value, null:Varchar, null:Date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 4:Int32] }
├── output: [ clients_contacts_dm.client_id, 'client':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar))), null:Varchar, null:Varchar, clients_contacts_dm.client_id, clients_contacts_dm.type, clients_contacts_dm.value, null:Varchar, null:Date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 4:Int32 ]
├── stream key: [ clients_contacts_dm.client_id, clients_contacts_dm.type, clients_contacts_dm.value, _rw_projected_row_id ]
└── StreamProjectSet { select_list: [$0, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar))), $2, $1] } { output: [ _rw_projected_row_id, clients_contacts_dm.client_id, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar))), clients_contacts_dm.type, clients_contacts_dm.value ], stream key: [ clients_contacts_dm.client_id, clients_contacts_dm.type, clients_contacts_dm.value, _rw_projected_row_id ] }
    └── StreamProject { exprs: [clients_contacts_dm.client_id, clients_contacts_dm.value, clients_contacts_dm.type] } { output: [ clients_contacts_dm.client_id, clients_contacts_dm.value, clients_contacts_dm.type ], stream key: [ clients_contacts_dm.client_id, clients_contacts_dm.type, clients_contacts_dm.value ] }
        └── StreamFilter { predicate: IsNull(clients_contacts_dm.disabled_at) } { output: [ clients_contacts_dm.client_id, clients_contacts_dm.value, clients_contacts_dm.type, clients_contacts_dm.disabled_at ], stream key: [ clients_contacts_dm.client_id, clients_contacts_dm.type, clients_contacts_dm.value ] }
            └── StreamTableScan { table: clients_contacts_dm, columns: [client_id, value, type, disabled_at] } { output: [ clients_contacts_dm.client_id, clients_contacts_dm.value, clients_contacts_dm.type, clients_contacts_dm.disabled_at ], stream key: [ clients_contacts_dm.client_id, clients_contacts_dm.type, clients_contacts_dm.value ] }
                ├── Upstream { output: [ client_id, value, type, disabled_at ], stream key: [] }
                └── BatchPlanNode { output: [ client_id, value, type, disabled_at ], stream key: [] }

Fragment 24901 (Actor 119490,119489)
StreamProject { exprs: [accounts_to_clients_dm.client_id, 'client':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), null:Varchar, accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, null:Varchar, null:Varchar, accounts_to_clients_dm.effective_start_date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 5:Int32] }
├── output: [ accounts_to_clients_dm.client_id, 'client':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), null:Varchar, accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, null:Varchar, null:Varchar, accounts_to_clients_dm.effective_start_date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 5:Int32 ]
├── stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, _rw_projected_row_id ]
└── StreamProjectSet { select_list: [$0, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), $3, $4] } { output: [ _rw_projected_row_id, accounts_to_clients_dm.client_id, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, _rw_projected_row_id ] }
    └── StreamProject { exprs: [accounts_to_clients_dm.client_id, accounts_dm.name, accounts_dm.number, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date] } { output: [ accounts_to_clients_dm.client_id, accounts_dm.name, accounts_dm.number, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ] }
        └── StreamFilter { predicate: IsNull(accounts_to_clients_dm.disabled_at) } { output: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.disabled_at, accounts_to_clients_dm.effective_end_date, accounts_dm.name, accounts_dm.number, accounts_to_clients_dm.effective_start_date, accounts_dm.account_id ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ] }
            └── MergeExecutor { output: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.disabled_at, accounts_to_clients_dm.effective_end_date, accounts_dm.name, accounts_dm.number, accounts_to_clients_dm.effective_start_date, accounts_dm.account_id ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ] }

Fragment 24902 (Actor 119484,119483)
StreamSyncLogStore { output: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.disabled_at, accounts_to_clients_dm.effective_end_date, accounts_dm.name, accounts_dm.number, accounts_to_clients_dm.effective_start_date, accounts_dm.account_id ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ] }
└── MergeExecutor { output: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.disabled_at, accounts_to_clients_dm.effective_end_date, accounts_dm.name, accounts_dm.number, accounts_to_clients_dm.effective_start_date, accounts_dm.account_id ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ] }

Fragment 24903 (Actor 119485,119486)
StreamSyncLogStore { output: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.disabled_at, accounts_to_clients_dm.effective_end_date, accounts_dm.name, accounts_dm.number, accounts_to_clients_dm.effective_start_date, accounts_dm.account_id ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ] }
└── StreamHashJoin { type: Inner, predicate: accounts_to_clients_dm.account_id = accounts_dm.account_id } { output: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.disabled_at, accounts_to_clients_dm.effective_end_date, accounts_dm.name, accounts_dm.number, accounts_to_clients_dm.effective_start_date, accounts_dm.account_id ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ] }
    ├── MergeExecutor { output: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.disabled_at, accounts_to_clients_dm.effective_end_date, accounts_to_clients_dm.effective_start_date ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ] }
    └── MergeExecutor { output: [ accounts_dm.account_id, accounts_dm.name, accounts_dm.number ], stream key: [ accounts_dm.account_id ] }

Fragment 24904 (Actor 119491,119492)
StreamFilter { predicate: IsNull(accounts_to_clients_dm.disabled_at) } { output: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.disabled_at, accounts_to_clients_dm.effective_end_date, accounts_to_clients_dm.effective_start_date ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ] }
└── StreamTableScan { table: accounts_to_clients_dm, columns: [account_id, client_id, disabled_at, effective_end_date, effective_start_date] } { output: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.disabled_at, accounts_to_clients_dm.effective_end_date, accounts_to_clients_dm.effective_start_date ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ] }
    ├── Upstream { output: [ account_id, client_id, disabled_at, effective_end_date, effective_start_date ], stream key: [] }
    └── BatchPlanNode { output: [ account_id, client_id, disabled_at, effective_end_date, effective_start_date ], stream key: [] }

Fragment 24905 (Actor 119493,119494)
StreamProject { exprs: [accounts_dm.account_id, accounts_dm.name, accounts_dm.number] } { output: [ accounts_dm.account_id, accounts_dm.name, accounts_dm.number ], stream key: [ accounts_dm.account_id ] }
└── StreamFilter { predicate: IsNull(accounts_dm.disabled_at) } { output: [ accounts_dm.account_id, accounts_dm.name, accounts_dm.number, accounts_dm.disabled_at ], stream key: [ accounts_dm.account_id ] }
    └── StreamTableScan { table: accounts_dm, columns: [account_id, name, number, disabled_at] } { output: [ accounts_dm.account_id, accounts_dm.name, accounts_dm.number, accounts_dm.disabled_at ], stream key: [ accounts_dm.account_id ] }
        ├── Upstream { output: [ account_id, name, number, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ account_id, name, number, disabled_at ], stream key: [] }

Fragment 24906 (Actor 119495,119496)
StreamProject { exprs: [clients_portfolios_dm.client_id, 'client':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), null:Varchar, clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, null:Varchar, null:Varchar, clients_portfolios_dm.effective_start_date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 6:Int32] }
├── output: [ clients_portfolios_dm.client_id, 'client':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), null:Varchar, clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, null:Varchar, null:Varchar, clients_portfolios_dm.effective_start_date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 6:Int32 ]
├── stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, _rw_projected_row_id ]
└── StreamProjectSet { select_list: [$0, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), $3, $4] } { output: [ _rw_projected_row_id, clients_portfolios_dm.client_id, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, _rw_projected_row_id ] }
    └── MergeExecutor { output: [ clients_portfolios_dm.client_id, portfolios_dm.name, portfolios_dm.number, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, portfolios_dm.portfolio_id ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ] }

Fragment 24907 (Actor 119498,119497)
StreamSyncLogStore { output: [ clients_portfolios_dm.client_id, portfolios_dm.name, portfolios_dm.number, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, portfolios_dm.portfolio_id ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ] }
└── StreamHashJoin { type: Inner, predicate: clients_portfolios_dm.portfolio_id = portfolios_dm.portfolio_id } { output: [ clients_portfolios_dm.client_id, portfolios_dm.name, portfolios_dm.number, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, portfolios_dm.portfolio_id ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ] }
    ├── MergeExecutor { output: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ] }
    └── MergeExecutor { output: [ portfolios_dm.portfolio_id, portfolios_dm.name, portfolios_dm.number ], stream key: [ portfolios_dm.portfolio_id ] }

Fragment 24908 (Actor 119500,119499)
StreamProject { exprs: [clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date] } { output: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ] }
└── StreamFilter { predicate: IsNull(clients_portfolios_dm.disabled_at) } { output: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, clients_portfolios_dm.disabled_at ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ] }
    └── StreamTableScan { table: clients_portfolios_dm, columns: [client_id, portfolio_id, effective_start_date, disabled_at] } { output: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, clients_portfolios_dm.disabled_at ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ] }
        ├── Upstream { output: [ client_id, portfolio_id, effective_start_date, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ client_id, portfolio_id, effective_start_date, disabled_at ], stream key: [] }

Fragment 24909 (Actor 119501,119502)
StreamProject { exprs: [portfolios_dm.portfolio_id, portfolios_dm.name, portfolios_dm.number] } { output: [ portfolios_dm.portfolio_id, portfolios_dm.name, portfolios_dm.number ], stream key: [ portfolios_dm.portfolio_id ] }
└── StreamFilter { predicate: IsNull(portfolios_dm.disabled_at) } { output: [ portfolios_dm.portfolio_id, portfolios_dm.name, portfolios_dm.number, portfolios_dm.disabled_at ], stream key: [ portfolios_dm.portfolio_id ] }
    └── StreamTableScan { table: portfolios_dm, columns: [portfolio_id, name, number, disabled_at] } { output: [ portfolios_dm.portfolio_id, portfolios_dm.name, portfolios_dm.number, portfolios_dm.disabled_at ], stream key: [ portfolios_dm.portfolio_id ] }
        ├── Upstream { output: [ portfolio_id, name, number, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ portfolio_id, name, number, disabled_at ], stream key: [] }

Fragment 24910 (Actor 119504,119503)
StreamProject { exprs: [party_items_mv_next.entity_id, party_items_mv_next.entity_type, party_items_mv_next.val, party_items_mv_next.ar_val, party_items_mv_next.filter_val, party_items_mv_next.active_parties_mv.id, party_items_mv_next.null:Varchar, party_items_mv_next.null:Varchar#1, party_items_mv_next.null:Varchar#2, party_items_mv_next.null:Date, null:Date, party_items_mv_next._rw_projected_row_id, party_items_mv_next.null:Int64, party_items_mv_next.null:Int32, party_items_mv_next.null:Int32#1, party_items_mv_next.$src, 7:Int32] }
├── output: [ party_items_mv_next.entity_id, party_items_mv_next.entity_type, party_items_mv_next.val, party_items_mv_next.ar_val, party_items_mv_next.filter_val, party_items_mv_next.active_parties_mv.id, party_items_mv_next.null:Varchar, party_items_mv_next.null:Varchar#1, party_items_mv_next.null:Varchar#2, party_items_mv_next.null:Date, null:Date, party_items_mv_next._rw_projected_row_id, party_items_mv_next.null:Int64, party_items_mv_next.null:Int32, party_items_mv_next.null:Int32#1, party_items_mv_next.$src, 7:Int32 ]
├── stream key: [ party_items_mv_next.active_parties_mv.id, party_items_mv_next._rw_projected_row_id, party_items_mv_next.null:Varchar, party_items_mv_next.null:Int32, party_items_mv_next.null:Int64, party_items_mv_next.null:Date, party_items_mv_next.null:Int32#1, party_items_mv_next.null:Varchar#1, party_items_mv_next.null:Varchar#2, party_items_mv_next.$src ]
└── StreamTableScan { table: party_items_mv_next, columns: [entity_id, entity_type, val, ar_val, filter_val, active_parties_mv.id, _rw_projected_row_id, null:Varchar, null:Int32, null:Int64, null:Date, null:Int32#1, null:Varchar#1, null:Varchar#2, $src] }
    ├── output: [ party_items_mv_next.entity_id, party_items_mv_next.entity_type, party_items_mv_next.val, party_items_mv_next.ar_val, party_items_mv_next.filter_val, party_items_mv_next.active_parties_mv.id, party_items_mv_next._rw_projected_row_id, party_items_mv_next.null:Varchar, party_items_mv_next.null:Int32, party_items_mv_next.null:Int64, party_items_mv_next.null:Date, party_items_mv_next.null:Int32#1, party_items_mv_next.null:Varchar#1, party_items_mv_next.null:Varchar#2, party_items_mv_next.$src ]
    ├── stream key: [ party_items_mv_next.active_parties_mv.id, party_items_mv_next._rw_projected_row_id, party_items_mv_next.null:Varchar, party_items_mv_next.null:Int32, party_items_mv_next.null:Int64, party_items_mv_next.null:Date, party_items_mv_next.null:Int32#1, party_items_mv_next.null:Varchar#1, party_items_mv_next.null:Varchar#2, party_items_mv_next.$src ]
    ├── Upstream { output: [ entity_id, entity_type, val, ar_val, filter_val, active_parties_mv.id, _rw_projected_row_id, null:Varchar, null:Int32, null:Int64, null:Date, null:Int32#1, null:Varchar#1, null:Varchar#2, $src ], stream key: [] }
    └── BatchPlanNode { output: [ entity_id, entity_type, val, ar_val, filter_val, active_parties_mv.id, _rw_projected_row_id, null:Varchar, null:Int32, null:Int64, null:Date, null:Int32#1, null:Varchar#1, null:Varchar#2, $src ], stream key: [] }

Fragment 24911 (Actor 119506,119505)
StreamProject { exprs: [portfolios_dm.portfolio_id, 'portfolio':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), null:Varchar, portfolios_dm.portfolio_id, null:Varchar, null:Varchar, null:Varchar, null:Date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 8:Int32] }
├── output: [ portfolios_dm.portfolio_id, 'portfolio':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), null:Varchar, portfolios_dm.portfolio_id, null:Varchar, null:Varchar, null:Varchar, null:Date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 8:Int32 ]
├── stream key: [ portfolios_dm.portfolio_id, _rw_projected_row_id ]
└── StreamProjectSet { select_list: [$0, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar)))] } { output: [ _rw_projected_row_id, portfolios_dm.portfolio_id, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))) ], stream key: [ portfolios_dm.portfolio_id, _rw_projected_row_id ] }
    └── StreamProject { exprs: [portfolios_dm.portfolio_id, portfolios_dm.name, portfolios_dm.number] } { output: [ portfolios_dm.portfolio_id, portfolios_dm.name, portfolios_dm.number ], stream key: [ portfolios_dm.portfolio_id ] }
        └── StreamFilter { predicate: IsNull(portfolios_dm.disabled_at) } { output: [ portfolios_dm.portfolio_id, portfolios_dm.name, portfolios_dm.number, portfolios_dm.disabled_at ], stream key: [ portfolios_dm.portfolio_id ] }
            └── StreamTableScan { table: portfolios_dm, columns: [portfolio_id, name, number, disabled_at] } { output: [ portfolios_dm.portfolio_id, portfolios_dm.name, portfolios_dm.number, portfolios_dm.disabled_at ], stream key: [ portfolios_dm.portfolio_id ] }
                ├── Upstream { output: [ portfolio_id, name, number, disabled_at ], stream key: [] }
                └── BatchPlanNode { output: [ portfolio_id, name, number, disabled_at ], stream key: [] }

Fragment 24912 (Actor 119507,119508)
StreamProject { exprs: [account_to_portfolios_dm.portfolio_id, 'portfolio':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), null:Varchar, account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, null:Varchar, null:Varchar, account_to_portfolios_dm.effective_start_date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 9:Int32] }
├── output: [ account_to_portfolios_dm.portfolio_id, 'portfolio':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), null:Varchar, account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, null:Varchar, null:Varchar, account_to_portfolios_dm.effective_start_date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 9:Int32 ]
├── stream key: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date, _rw_projected_row_id ]
└── StreamProjectSet { select_list: [$0, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), $3, $4] } { output: [ _rw_projected_row_id, account_to_portfolios_dm.portfolio_id, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar))), account_to_portfolios_dm.account_id, account_to_portfolios_dm.effective_start_date ], stream key: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date, _rw_projected_row_id ] }
    └── MergeExecutor { output: [ account_to_portfolios_dm.portfolio_id, accounts_dm.name, accounts_dm.number, account_to_portfolios_dm.account_id, account_to_portfolios_dm.effective_start_date, accounts_dm.account_id ], stream key: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ] }

Fragment 24913 (Actor 119510,119509)
StreamSyncLogStore { output: [ account_to_portfolios_dm.portfolio_id, accounts_dm.name, accounts_dm.number, account_to_portfolios_dm.account_id, account_to_portfolios_dm.effective_start_date, accounts_dm.account_id ], stream key: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ] }
└── StreamHashJoin { type: Inner, predicate: account_to_portfolios_dm.account_id = accounts_dm.account_id } { output: [ account_to_portfolios_dm.portfolio_id, accounts_dm.name, accounts_dm.number, account_to_portfolios_dm.account_id, account_to_portfolios_dm.effective_start_date, accounts_dm.account_id ], stream key: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ] }
    ├── MergeExecutor { output: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ], stream key: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ] }
    └── MergeExecutor { output: [ accounts_dm.account_id, accounts_dm.name, accounts_dm.number ], stream key: [ accounts_dm.account_id ] }

Fragment 24914 (Actor 119511,119512)
StreamProject { exprs: [account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date] } { output: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ], stream key: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ] }
└── StreamFilter { predicate: IsNull(account_to_portfolios_dm.disabled_at) } { output: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date, account_to_portfolios_dm.disabled_at ], stream key: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ] }
    └── StreamTableScan { table: account_to_portfolios_dm, columns: [account_id, portfolio_id, effective_start_date, disabled_at] } { output: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date, account_to_portfolios_dm.disabled_at ], stream key: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ] }
        ├── Upstream { output: [ account_id, portfolio_id, effective_start_date, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ account_id, portfolio_id, effective_start_date, disabled_at ], stream key: [] }

Fragment 24915 (Actor 119514,119513)
StreamProject { exprs: [accounts_dm.account_id, accounts_dm.name, accounts_dm.number] } { output: [ accounts_dm.account_id, accounts_dm.name, accounts_dm.number ], stream key: [ accounts_dm.account_id ] }
└── StreamFilter { predicate: IsNull(accounts_dm.disabled_at) } { output: [ accounts_dm.account_id, accounts_dm.name, accounts_dm.number, accounts_dm.disabled_at ], stream key: [ accounts_dm.account_id ] }
    └── StreamTableScan { table: accounts_dm, columns: [account_id, name, number, disabled_at] } { output: [ accounts_dm.account_id, accounts_dm.name, accounts_dm.number, accounts_dm.disabled_at ], stream key: [ accounts_dm.account_id ] }
        ├── Upstream { output: [ account_id, name, number, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ account_id, name, number, disabled_at ], stream key: [] }

Fragment 24916 (Actor 119518,119517)
StreamProject { exprs: [clients_portfolios_dm.portfolio_id, 'portfolio':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar), $3, Replace($3, '-':Varchar, '':Varchar), $4, Replace($4, '-':Varchar, '':Varchar))), Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar), $3, Replace($3, '-':Varchar, '':Varchar), $4, Replace($4, '-':Varchar, '':Varchar))), null:Varchar, clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, null:Varchar, null:Varchar, clients_portfolios_dm.effective_start_date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 10:Int32] }
├── output: [ clients_portfolios_dm.portfolio_id, 'portfolio':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar), $3, Replace($3, '-':Varchar, '':Varchar), $4, Replace($4, '-':Varchar, '':Varchar))), Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar), $3, Replace($3, '-':Varchar, '':Varchar), $4, Replace($4, '-':Varchar, '':Varchar))), null:Varchar, clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, null:Varchar, null:Varchar, clients_portfolios_dm.effective_start_date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 10:Int32 ]
├── stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, _rw_projected_row_id ]
└── StreamProjectSet { select_list: [$0, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar), $3, Replace($3, '-':Varchar, '':Varchar), $4, Replace($4, '-':Varchar, '':Varchar))), $5, $6] } { output: [ _rw_projected_row_id, clients_portfolios_dm.portfolio_id, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $2, Replace($2, '-':Varchar, '':Varchar), $3, Replace($3, '-':Varchar, '':Varchar), $4, Replace($4, '-':Varchar, '':Varchar))), clients_portfolios_dm.client_id, clients_portfolios_dm.effective_start_date ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, _rw_projected_row_id ] }
    └── MergeExecutor { output: [ clients_portfolios_dm.portfolio_id, clients_dm.display_name, clients_dm.local_display_name, clients_dm.preferred_name, clients_dm.customer_identification_file, clients_portfolios_dm.client_id, clients_portfolios_dm.effective_start_date, clients_dm.id ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ] }

Fragment 24917 (Actor 119515,119516)
StreamSyncLogStore { output: [ clients_portfolios_dm.portfolio_id, clients_dm.display_name, clients_dm.local_display_name, clients_dm.preferred_name, clients_dm.customer_identification_file, clients_portfolios_dm.client_id, clients_portfolios_dm.effective_start_date, clients_dm.id ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ] }
└── StreamHashJoin { type: Inner, predicate: clients_portfolios_dm.client_id = clients_dm.id } { output: [ clients_portfolios_dm.portfolio_id, clients_dm.display_name, clients_dm.local_display_name, clients_dm.preferred_name, clients_dm.customer_identification_file, clients_portfolios_dm.client_id, clients_portfolios_dm.effective_start_date, clients_dm.id ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ] }
    ├── MergeExecutor { output: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ] }
    └── MergeExecutor { output: [ clients_dm.id, clients_dm.display_name, clients_dm.local_display_name, clients_dm.preferred_name, clients_dm.customer_identification_file ], stream key: [ clients_dm.id ] }

Fragment 24918 (Actor 119519,119520)
StreamProject { exprs: [clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date] } { output: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ] }
└── StreamFilter { predicate: IsNull(clients_portfolios_dm.disabled_at) } { output: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, clients_portfolios_dm.disabled_at ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ] }
    └── StreamTableScan { table: clients_portfolios_dm, columns: [client_id, portfolio_id, effective_start_date, disabled_at] } { output: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, clients_portfolios_dm.disabled_at ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ] }
        ├── Upstream { output: [ client_id, portfolio_id, effective_start_date, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ client_id, portfolio_id, effective_start_date, disabled_at ], stream key: [] }

Fragment 24919 (Actor 119521,119522)
StreamProject { exprs: [clients_dm.id, clients_dm.display_name, clients_dm.local_display_name, clients_dm.preferred_name, clients_dm.customer_identification_file] } { output: [ clients_dm.id, clients_dm.display_name, clients_dm.local_display_name, clients_dm.preferred_name, clients_dm.customer_identification_file ], stream key: [ clients_dm.id ] }
└── StreamFilter { predicate: IsNull(clients_dm.closing_date) } { output: [ clients_dm.id, clients_dm.display_name, clients_dm.local_display_name, clients_dm.preferred_name, clients_dm.customer_identification_file, clients_dm.closing_date ], stream key: [ clients_dm.id ] }
    └── StreamTableScan { table: clients_dm, columns: [id, display_name, local_display_name, preferred_name, customer_identification_file, closing_date] } { output: [ clients_dm.id, clients_dm.display_name, clients_dm.local_display_name, clients_dm.preferred_name, clients_dm.customer_identification_file, clients_dm.closing_date ], stream key: [ clients_dm.id ] }
        ├── Upstream { output: [ id, display_name, local_display_name, preferred_name, customer_identification_file, closing_date ], stream key: [] }
        └── BatchPlanNode { output: [ id, display_name, local_display_name, preferred_name, customer_identification_file, closing_date ], stream key: [] }

Fragment 24920 (Actor 119526,119525)
StreamProject { exprs: [clients_portfolios_dm.portfolio_id, 'portfolio':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar))), null:Varchar, null:Varchar, clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_contacts_dm.type, clients_contacts_dm.value, clients_portfolios_dm.effective_start_date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 11:Int32] }
├── output: [ clients_portfolios_dm.portfolio_id, 'portfolio':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar))), null:Varchar, null:Varchar, clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_contacts_dm.type, clients_contacts_dm.value, clients_portfolios_dm.effective_start_date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 11:Int32 ]
├── stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, clients_contacts_dm.type, clients_contacts_dm.value, _rw_projected_row_id ]
└── StreamProjectSet { select_list: [$0, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar))), $2, $3, $5, $1] } { output: [ _rw_projected_row_id, clients_portfolios_dm.portfolio_id, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar))), clients_portfolios_dm.client_id, clients_portfolios_dm.effective_start_date, clients_contacts_dm.type, clients_contacts_dm.value ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, clients_contacts_dm.type, clients_contacts_dm.value, _rw_projected_row_id ] }
    └── MergeExecutor { output: [ clients_portfolios_dm.portfolio_id, clients_contacts_dm.value, clients_portfolios_dm.client_id, clients_portfolios_dm.effective_start_date, clients_contacts_dm.client_id, clients_contacts_dm.type ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, clients_contacts_dm.type, clients_contacts_dm.value ] }

Fragment 24921 (Actor 119523,119524)
StreamSyncLogStore { output: [ clients_portfolios_dm.portfolio_id, clients_contacts_dm.value, clients_portfolios_dm.client_id, clients_portfolios_dm.effective_start_date, clients_contacts_dm.client_id, clients_contacts_dm.type ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, clients_contacts_dm.type, clients_contacts_dm.value ] }
└── StreamHashJoin { type: Inner, predicate: clients_portfolios_dm.client_id = clients_contacts_dm.client_id } { output: [ clients_portfolios_dm.portfolio_id, clients_contacts_dm.value, clients_portfolios_dm.client_id, clients_portfolios_dm.effective_start_date, clients_contacts_dm.client_id, clients_contacts_dm.type ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, clients_contacts_dm.type, clients_contacts_dm.value ] }
    ├── MergeExecutor { output: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ] }
    └── MergeExecutor { output: [ clients_contacts_dm.client_id, clients_contacts_dm.value, clients_contacts_dm.type ], stream key: [ clients_contacts_dm.client_id, clients_contacts_dm.type, clients_contacts_dm.value ] }

Fragment 24922 (Actor 119527,119528)
StreamProject { exprs: [clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date] } { output: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ] }
└── StreamFilter { predicate: IsNull(clients_portfolios_dm.disabled_at) } { output: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, clients_portfolios_dm.disabled_at ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ] }
    └── StreamTableScan { table: clients_portfolios_dm, columns: [client_id, portfolio_id, effective_start_date, disabled_at] } { output: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, clients_portfolios_dm.disabled_at ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ] }
        ├── Upstream { output: [ client_id, portfolio_id, effective_start_date, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ client_id, portfolio_id, effective_start_date, disabled_at ], stream key: [] }

Fragment 24923 (Actor 119530,119529)
StreamProject { exprs: [clients_contacts_dm.client_id, clients_contacts_dm.value, clients_contacts_dm.type] } { output: [ clients_contacts_dm.client_id, clients_contacts_dm.value, clients_contacts_dm.type ], stream key: [ clients_contacts_dm.client_id, clients_contacts_dm.type, clients_contacts_dm.value ] }
└── StreamFilter { predicate: IsNull(clients_contacts_dm.disabled_at) } { output: [ clients_contacts_dm.client_id, clients_contacts_dm.value, clients_contacts_dm.type, clients_contacts_dm.disabled_at ], stream key: [ clients_contacts_dm.client_id, clients_contacts_dm.type, clients_contacts_dm.value ] }
    └── StreamTableScan { table: clients_contacts_dm, columns: [client_id, value, type, disabled_at] } { output: [ clients_contacts_dm.client_id, clients_contacts_dm.value, clients_contacts_dm.type, clients_contacts_dm.disabled_at ], stream key: [ clients_contacts_dm.client_id, clients_contacts_dm.type, clients_contacts_dm.value ] }
        ├── Upstream { output: [ client_id, value, type, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ client_id, value, type, disabled_at ], stream key: [] }

Fragment 24924 (Actor 119532,119531)
StreamProject { exprs: [assets_dm.id, 'asset':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $3, Replace($3, '-':Varchar, '':Varchar), $4, Replace($4, '-':Varchar, '':Varchar), $7, Replace($7, '-':Varchar, '':Varchar), $5, Replace($5, '-':Varchar, '':Varchar), $6, Replace($6, '-':Varchar, '':Varchar), $8, Replace($8, '-':Varchar, '':Varchar))), null:Varchar, Lower(assets_dm.type) as $expr1, assets_dm.id, null:Varchar, null:Varchar, null:Varchar, null:Date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 12:Int32] }
├── output: [ assets_dm.id, 'asset':Varchar, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $3, Replace($3, '-':Varchar, '':Varchar), $4, Replace($4, '-':Varchar, '':Varchar), $7, Replace($7, '-':Varchar, '':Varchar), $5, Replace($5, '-':Varchar, '':Varchar), $6, Replace($6, '-':Varchar, '':Varchar), $8, Replace($8, '-':Varchar, '':Varchar))), null:Varchar, $expr1, assets_dm.id, null:Varchar, null:Varchar, null:Varchar, null:Date, null:Date, _rw_projected_row_id, null:Int64, null:Int32, null:Int32, null:Int32, 12:Int32 ]
├── stream key: [ assets_dm.id, _rw_projected_row_id ]
└── StreamProjectSet { select_list: [$0, $2, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $3, Replace($3, '-':Varchar, '':Varchar), $4, Replace($4, '-':Varchar, '':Varchar), $7, Replace($7, '-':Varchar, '':Varchar), $5, Replace($5, '-':Varchar, '':Varchar), $6, Replace($6, '-':Varchar, '':Varchar), $8, Replace($8, '-':Varchar, '':Varchar)))] }
    ├── output: [ _rw_projected_row_id, assets_dm.id, assets_dm.type, Unnest(Array($1, Replace($1, '-':Varchar, '':Varchar), $3, Replace($3, '-':Varchar, '':Varchar), $4, Replace($4, '-':Varchar, '':Varchar), $7, Replace($7, '-':Varchar, '':Varchar), $5, Replace($5, '-':Varchar, '':Varchar), $6, Replace($6, '-':Varchar, '':Varchar), $8, Replace($8, '-':Varchar, '':Varchar))) ]
    ├── stream key: [ assets_dm.id, _rw_projected_row_id ]
    └── StreamProject { exprs: [assets_dm.id, assets_dm.name_en, assets_dm.type, assets_dm.ticker, assets_dm.isin, assets_dm.cusip, assets_dm.sedol, assets_dm.ric, assets_dm.figi] } { output: [ assets_dm.id, assets_dm.name_en, assets_dm.type, assets_dm.ticker, assets_dm.isin, assets_dm.cusip, assets_dm.sedol, assets_dm.ric, assets_dm.figi ], stream key: [ assets_dm.id ] }
        └── StreamFilter { predicate: IsNull(assets_dm.disabled_at) } { output: [ assets_dm.id, assets_dm.name_en, assets_dm.type, assets_dm.ticker, assets_dm.isin, assets_dm.cusip, assets_dm.sedol, assets_dm.ric, assets_dm.figi, assets_dm.disabled_at ], stream key: [ assets_dm.id ] }
            └── StreamTableScan { table: assets_dm, columns: [id, name_en, type, ticker, isin, cusip, sedol, ric, figi, disabled_at] } { output: [ assets_dm.id, assets_dm.name_en, assets_dm.type, assets_dm.ticker, assets_dm.isin, assets_dm.cusip, assets_dm.sedol, assets_dm.ric, assets_dm.figi, assets_dm.disabled_at ], stream key: [ assets_dm.id ] }
                ├── Upstream { output: [ id, name_en, type, ticker, isin, cusip, sedol, ric, figi, disabled_at ], stream key: [] }
                └── BatchPlanNode { output: [ id, name_en, type, ticker, isin, cusip, sedol, ric, figi, disabled_at ], stream key: [] }

Fragment 24925 (Actor 119534,119533)
StreamProject { exprs: [assets_dm.id, 'asset':Varchar, null:Varchar, assets_dm.name_ar, Lower(assets_dm.type) as $expr2, assets_dm.id, null:Varchar, null:Varchar, null:Varchar, null:Date, null:Date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 13:Int32] } { output: [ assets_dm.id, 'asset':Varchar, null:Varchar, assets_dm.name_ar, $expr2, assets_dm.id, null:Varchar, null:Varchar, null:Varchar, null:Date, null:Date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 13:Int32 ], stream key: [ assets_dm.id ] }
└── StreamFilter { predicate: IsNull(assets_dm.disabled_at) } { output: [ assets_dm.id, assets_dm.name_ar, assets_dm.type, assets_dm.disabled_at ], stream key: [ assets_dm.id ] }
    └── StreamTableScan { table: assets_dm, columns: [id, name_ar, type, disabled_at] } { output: [ assets_dm.id, assets_dm.name_ar, assets_dm.type, assets_dm.disabled_at ], stream key: [ assets_dm.id ] }
        ├── Upstream { output: [ id, name_ar, type, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ id, name_ar, type, disabled_at ], stream key: [] }

Fragment 24926 (Actor 119478,119477)
StreamProject { exprs: [clients_dm.id, 'client':Varchar, null:Varchar, null:Varchar, ConcatOp('client_type:':Varchar, Lower(clients_dm.type)) as $expr3, clients_dm.id, null:Varchar, null:Varchar, null:Varchar, null:Date, null:Date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 14:Int32] } { output: [ clients_dm.id, 'client':Varchar, null:Varchar, null:Varchar, $expr3, clients_dm.id, null:Varchar, null:Varchar, null:Varchar, null:Date, null:Date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 14:Int32 ], stream key: [ clients_dm.id ] }
└── MergeExecutor { output: [ clients_dm.id, clients_dm.display_name, clients_dm.local_display_name, clients_dm.preferred_name, clients_dm.customer_identification_file, clients_dm.type ], stream key: [ clients_dm.id ] }

Fragment 24927 (Actor 119535,119536)
StreamProject { exprs: [clients_dm.id, 'client':Varchar, null:Varchar, null:Varchar, ConcatOp('status:':Varchar, Lower(labels_dm.name_en)) as $expr4, clients_dm.id, clients_dm.status_label_id, null:Varchar, null:Varchar, null:Date, null:Date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 15:Int32] } { output: [ clients_dm.id, 'client':Varchar, null:Varchar, null:Varchar, $expr4, clients_dm.id, clients_dm.status_label_id, null:Varchar, null:Varchar, null:Date, null:Date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 15:Int32 ], stream key: [ clients_dm.id, clients_dm.status_label_id ] }
└── MergeExecutor { output: [ clients_dm.id, labels_dm.name_en, clients_dm.status_label_id, labels_dm.label_id ], stream key: [ clients_dm.id, clients_dm.status_label_id ] }

Fragment 24928 (Actor 119537,119538)
StreamSyncLogStore { output: [ clients_dm.id, labels_dm.name_en, clients_dm.status_label_id, labels_dm.label_id ], stream key: [ clients_dm.id, clients_dm.status_label_id ] }
└── StreamHashJoin { type: Inner, predicate: clients_dm.status_label_id = labels_dm.label_id } { output: [ clients_dm.id, labels_dm.name_en, clients_dm.status_label_id, labels_dm.label_id ], stream key: [ clients_dm.id, clients_dm.status_label_id ] }
    ├── MergeExecutor { output: [ clients_dm.id, clients_dm.status_label_id ], stream key: [ clients_dm.id ] }
    └── MergeExecutor { output: [ labels_dm.label_id, labels_dm.name_en ], stream key: [ labels_dm.label_id ] }

Fragment 24929 (Actor 119540,119539)
StreamProject { exprs: [clients_dm.id, clients_dm.status_label_id] } { output: [ clients_dm.id, clients_dm.status_label_id ], stream key: [ clients_dm.id ] }
└── StreamFilter { predicate: IsNull(clients_dm.closing_date) } { output: [ clients_dm.id, clients_dm.status_label_id, clients_dm.closing_date ], stream key: [ clients_dm.id ] }
    └── StreamTableScan { table: clients_dm, columns: [id, status_label_id, closing_date] } { output: [ clients_dm.id, clients_dm.status_label_id, clients_dm.closing_date ], stream key: [ clients_dm.id ] }
        ├── Upstream { output: [ id, status_label_id, closing_date ], stream key: [] }
        └── BatchPlanNode { output: [ id, status_label_id, closing_date ], stream key: [] }

Fragment 24930 (Actor 119541,119542)
StreamTableScan { table: labels_dm, columns: [label_id, name_en] } { output: [ labels_dm.label_id, labels_dm.name_en ], stream key: [ labels_dm.label_id ] }
├── Upstream { output: [ label_id, name_en ], stream key: [] }
└── BatchPlanNode { output: [ label_id, name_en ], stream key: [] }

Fragment 24931 (Actor 119546,119545)
StreamProject { exprs: [entity_to_teams_dm.entity_id, 'client':Varchar, null:Varchar, null:Varchar, ConcatOp('team:':Varchar, entity_to_teams_dm.team_id) as $expr5, entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, null:Varchar, entity_to_teams_dm.effective_start_date, null:Date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 16:Int32] }
├── output: [ entity_to_teams_dm.entity_id, 'client':Varchar, null:Varchar, null:Varchar, $expr5, entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, null:Varchar, entity_to_teams_dm.effective_start_date, null:Date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 16:Int32 ]
├── stream key: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ]
└── MergeExecutor { output: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date, clients_dm.id ], stream key: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }

Fragment 24932 (Actor 119543,119544)
StreamSyncLogStore { output: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date, clients_dm.id ], stream key: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }
└── StreamHashJoin { type: Inner, predicate: entity_to_teams_dm.entity_id = clients_dm.id } { output: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date, clients_dm.id ], stream key: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }
    ├── MergeExecutor { output: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ], stream key: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }
    └── MergeExecutor { output: [ clients_dm.id ], stream key: [ clients_dm.id ] }

Fragment 24933 (Actor 119548,119547)
StreamProject { exprs: [entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date] } { output: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ], stream key: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }
└── StreamFilter { predicate: (entity_to_teams_dm.entity_type = 'CLIENT':Varchar) AND IsNull(entity_to_teams_dm.disabled_at) AND IsNull(entity_to_teams_dm.effective_end_date) } { output: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date, entity_to_teams_dm.effective_end_date, entity_to_teams_dm.disabled_at ], stream key: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }
    └── StreamTableScan { table: entity_to_teams_dm, columns: [team_id, entity_id, entity_type, effective_start_date, effective_end_date, disabled_at] } { output: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date, entity_to_teams_dm.effective_end_date, entity_to_teams_dm.disabled_at ], stream key: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }
        ├── Upstream { output: [ team_id, entity_id, entity_type, effective_start_date, effective_end_date, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ team_id, entity_id, entity_type, effective_start_date, effective_end_date, disabled_at ], stream key: [] }

Fragment 24934 (Actor 119549,119550)
StreamProject { exprs: [clients_dm.id] } { output: [ clients_dm.id ], stream key: [ clients_dm.id ] }
└── StreamFilter { predicate: IsNull(clients_dm.closing_date) } { output: [ clients_dm.id, clients_dm.closing_date ], stream key: [ clients_dm.id ] }
    └── StreamTableScan { table: clients_dm, columns: [id, closing_date] } { output: [ clients_dm.id, clients_dm.closing_date ], stream key: [ clients_dm.id ] }
        ├── Upstream { output: [ id, closing_date ], stream key: [] }
        └── BatchPlanNode { output: [ id, closing_date ], stream key: [] }

Fragment 24935 (Actor 119007,119008)
StreamProject { exprs: [accounts_dm.account_id, 'account':Varchar, null:Varchar, null:Varchar, ConcatOp('account_type:':Varchar, Lower(product_types_dm.type)) as $expr6, accounts_dm.account_id, accounts_dm.product_type_id, null:Varchar, null:Varchar, null:Date, null:Date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 17:Int32] } { output: [ accounts_dm.account_id, 'account':Varchar, null:Varchar, null:Varchar, $expr6, accounts_dm.account_id, accounts_dm.product_type_id, null:Varchar, null:Varchar, null:Date, null:Date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 17:Int32 ], stream key: [ accounts_dm.account_id, accounts_dm.product_type_id ] }
└── StreamFilter { predicate: IsNull(accounts_dm.disabled_at) } { output: [ accounts_dm.account_id, accounts_dm.disabled_at, product_types_dm.type, product_types_dm.name_en, accounts_dm.product_type_id, product_types_dm.product_type_id ], stream key: [ accounts_dm.account_id, accounts_dm.product_type_id ] }
    └── MergeExecutor { output: [ accounts_dm.account_id, accounts_dm.disabled_at, product_types_dm.type, product_types_dm.name_en, accounts_dm.product_type_id, product_types_dm.product_type_id ], stream key: [ accounts_dm.account_id, accounts_dm.product_type_id ] }

Fragment 24936 (Actor 119003,119004)
StreamTemporalJoin { type: Inner, append_only: false, predicate: accounts_dm.product_type_id = product_types_dm.product_type_id, nested_loop: false } { output: [ accounts_dm.account_id, accounts_dm.disabled_at, product_types_dm.type, product_types_dm.name_en, accounts_dm.product_type_id, product_types_dm.product_type_id ], stream key: [ accounts_dm.account_id, accounts_dm.product_type_id ] }
├── MergeExecutor { output: [ accounts_dm.account_id, accounts_dm.product_type_id, accounts_dm.disabled_at ], stream key: [ accounts_dm.account_id ] }
└── MergeExecutor { output: [ product_types_dm.product_type_id, product_types_dm.type, product_types_dm.name_en ], stream key: [ product_types_dm.product_type_id ] }

Fragment 24937 (Actor 119551,119552)
StreamFilter { predicate: IsNull(accounts_dm.disabled_at) } { output: [ accounts_dm.account_id, accounts_dm.product_type_id, accounts_dm.disabled_at ], stream key: [ accounts_dm.account_id ] }
└── StreamTableScan { table: accounts_dm, columns: [account_id, product_type_id, disabled_at] } { output: [ accounts_dm.account_id, accounts_dm.product_type_id, accounts_dm.disabled_at ], stream key: [ accounts_dm.account_id ] }
    ├── Upstream { output: [ account_id, product_type_id, disabled_at ], stream key: [] }
    └── BatchPlanNode { output: [ account_id, product_type_id, disabled_at ], stream key: [] }

Fragment 24938 (Actor 119006,119005)
StreamTableScan { table: product_types_dm, columns: [product_type_id, type, name_en] } { output: [ product_types_dm.product_type_id, product_types_dm.type, product_types_dm.name_en ], stream key: [ product_types_dm.product_type_id ] }
├── Upstream { output: [ product_type_id, type, name_en ], stream key: [] }
└── BatchPlanNode { output: [ product_type_id, type, name_en ], stream key: [] }

Fragment 24939 (Actor 119001,119002)
StreamProject { exprs: [accounts_dm.account_id, 'account':Varchar, null:Varchar, null:Varchar, ConcatOp('product:':Varchar, Lower(Trim(product_types_dm.name_en))) as $expr7, accounts_dm.account_id, accounts_dm.product_type_id, null:Varchar, null:Varchar, null:Date, null:Date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 18:Int32] } { output: [ accounts_dm.account_id, 'account':Varchar, null:Varchar, null:Varchar, $expr7, accounts_dm.account_id, accounts_dm.product_type_id, null:Varchar, null:Varchar, null:Date, null:Date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 18:Int32 ], stream key: [ accounts_dm.account_id, accounts_dm.product_type_id ] }
└── StreamFilter { predicate: IsNull(accounts_dm.disabled_at) AND Not(IsNull(product_types_dm.name_en)) AND (Trim(product_types_dm.name_en) <> '':Varchar) } { output: [ accounts_dm.account_id, accounts_dm.disabled_at, product_types_dm.type, product_types_dm.name_en, accounts_dm.product_type_id, product_types_dm.product_type_id ], stream key: [ accounts_dm.account_id, accounts_dm.product_type_id ] }
    └── MergeExecutor { output: [ accounts_dm.account_id, accounts_dm.disabled_at, product_types_dm.type, product_types_dm.name_en, accounts_dm.product_type_id, product_types_dm.product_type_id ], stream key: [ accounts_dm.account_id, accounts_dm.product_type_id ] }

Fragment 24940 (Actor 119556,119555)
StreamProject { exprs: [accounts_dm.account_id, 'account':Varchar, null:Varchar, null:Varchar, ConcatOp('status:':Varchar, Lower(labels_dm.name_en)) as $expr8, accounts_dm.account_id, accounts_dm.status_label_id, null:Varchar, null:Varchar, null:Date, null:Date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 19:Int32] } { output: [ accounts_dm.account_id, 'account':Varchar, null:Varchar, null:Varchar, $expr8, accounts_dm.account_id, accounts_dm.status_label_id, null:Varchar, null:Varchar, null:Date, null:Date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 19:Int32 ], stream key: [ accounts_dm.account_id, accounts_dm.status_label_id ] }
└── MergeExecutor { output: [ accounts_dm.account_id, labels_dm.name_en, accounts_dm.status_label_id, labels_dm.label_id ], stream key: [ accounts_dm.account_id, accounts_dm.status_label_id ] }

Fragment 24941 (Actor 119553,119554)
StreamSyncLogStore { output: [ accounts_dm.account_id, labels_dm.name_en, accounts_dm.status_label_id, labels_dm.label_id ], stream key: [ accounts_dm.account_id, accounts_dm.status_label_id ] }
└── StreamHashJoin { type: Inner, predicate: accounts_dm.status_label_id = labels_dm.label_id } { output: [ accounts_dm.account_id, labels_dm.name_en, accounts_dm.status_label_id, labels_dm.label_id ], stream key: [ accounts_dm.account_id, accounts_dm.status_label_id ] }
    ├── MergeExecutor { output: [ accounts_dm.account_id, accounts_dm.status_label_id ], stream key: [ accounts_dm.account_id ] }
    └── MergeExecutor { output: [ labels_dm.label_id, labels_dm.name_en ], stream key: [ labels_dm.label_id ] }

Fragment 24942 (Actor 119558,119557)
StreamProject { exprs: [accounts_dm.account_id, accounts_dm.status_label_id] } { output: [ accounts_dm.account_id, accounts_dm.status_label_id ], stream key: [ accounts_dm.account_id ] }
└── StreamFilter { predicate: IsNull(accounts_dm.disabled_at) } { output: [ accounts_dm.account_id, accounts_dm.status_label_id, accounts_dm.disabled_at ], stream key: [ accounts_dm.account_id ] }
    └── StreamTableScan { table: accounts_dm, columns: [account_id, status_label_id, disabled_at] } { output: [ accounts_dm.account_id, accounts_dm.status_label_id, accounts_dm.disabled_at ], stream key: [ accounts_dm.account_id ] }
        ├── Upstream { output: [ account_id, status_label_id, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ account_id, status_label_id, disabled_at ], stream key: [] }

Fragment 24943 (Actor 119560,119559)
StreamTableScan { table: labels_dm, columns: [label_id, name_en] } { output: [ labels_dm.label_id, labels_dm.name_en ], stream key: [ labels_dm.label_id ] }
├── Upstream { output: [ label_id, name_en ], stream key: [] }
└── BatchPlanNode { output: [ label_id, name_en ], stream key: [] }

Fragment 24944 (Actor 119563,119564)
StreamProject { exprs: [accounts_to_clients_dm.account_id, 'account':Varchar, null:Varchar, null:Varchar, ConcatOp('team:':Varchar, entity_to_teams_dm.team_id) as $expr9, accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, entity_to_teams_dm.team_id, entity_to_teams_dm.entity_type, accounts_to_clients_dm.effective_start_date, entity_to_teams_dm.effective_start_date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 20:Int32] }
├── output: [ accounts_to_clients_dm.account_id, 'account':Varchar, null:Varchar, null:Varchar, $expr9, accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, entity_to_teams_dm.team_id, entity_to_teams_dm.entity_type, accounts_to_clients_dm.effective_start_date, entity_to_teams_dm.effective_start_date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 20:Int32 ]
├── stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, entity_to_teams_dm.team_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ]
└── MergeExecutor { output: [ accounts_to_clients_dm.account_id, entity_to_teams_dm.team_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, entity_to_teams_dm.team_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }

Fragment 24945 (Actor 119561,119562)
StreamSyncLogStore { output: [ accounts_to_clients_dm.account_id, entity_to_teams_dm.team_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, entity_to_teams_dm.team_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }
└── StreamHashJoin { type: Inner, predicate: accounts_to_clients_dm.client_id = entity_to_teams_dm.entity_id } { output: [ accounts_to_clients_dm.account_id, entity_to_teams_dm.team_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, entity_to_teams_dm.team_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }
    ├── MergeExecutor { output: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.disabled_at, accounts_to_clients_dm.effective_end_date, accounts_dm.name, accounts_dm.number, accounts_to_clients_dm.effective_start_date, accounts_dm.account_id ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ] }
    └── MergeExecutor { output: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ], stream key: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }

Fragment 24946 (Actor 119488,119487)
StreamFilter { predicate: IsNull(accounts_to_clients_dm.disabled_at) AND IsNull(accounts_to_clients_dm.effective_end_date) } { output: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.disabled_at, accounts_to_clients_dm.effective_end_date, accounts_dm.name, accounts_dm.number, accounts_to_clients_dm.effective_start_date, accounts_dm.account_id ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ] }
└── MergeExecutor { output: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.disabled_at, accounts_to_clients_dm.effective_end_date, accounts_dm.name, accounts_dm.number, accounts_to_clients_dm.effective_start_date, accounts_dm.account_id ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ] }

Fragment 24947 (Actor 119566,119565)
StreamProject { exprs: [entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date] } { output: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ], stream key: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }
└── StreamFilter { predicate: (entity_to_teams_dm.entity_type = 'CLIENT':Varchar) AND IsNull(entity_to_teams_dm.disabled_at) AND IsNull(entity_to_teams_dm.effective_end_date) } { output: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date, entity_to_teams_dm.effective_end_date, entity_to_teams_dm.disabled_at ], stream key: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }
    └── StreamTableScan { table: entity_to_teams_dm, columns: [team_id, entity_id, entity_type, effective_start_date, effective_end_date, disabled_at] } { output: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date, entity_to_teams_dm.effective_end_date, entity_to_teams_dm.disabled_at ], stream key: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }
        ├── Upstream { output: [ team_id, entity_id, entity_type, effective_start_date, effective_end_date, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ team_id, entity_id, entity_type, effective_start_date, effective_end_date, disabled_at ], stream key: [] }

Fragment 24948 (Actor 118987,118988)
StreamProject { exprs: [portfolios_dm.portfolio_id, 'portfolio':Varchar, null:Varchar, null:Varchar, ConcatOp('portfolio_type:':Varchar, Lower(service_types_dm.type)) as $expr10, portfolios_dm.portfolio_id, portfolios_dm.service_type_id, null:Varchar, null:Varchar, null:Date, null:Date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 21:Int32] } { output: [ portfolios_dm.portfolio_id, 'portfolio':Varchar, null:Varchar, null:Varchar, $expr10, portfolios_dm.portfolio_id, portfolios_dm.service_type_id, null:Varchar, null:Varchar, null:Date, null:Date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 21:Int32 ], stream key: [ portfolios_dm.portfolio_id, portfolios_dm.service_type_id ] }
└── StreamTemporalJoin { type: Inner, append_only: false, predicate: portfolios_dm.service_type_id = service_types_dm.service_type_id, nested_loop: false } { output: [ portfolios_dm.portfolio_id, service_types_dm.type, portfolios_dm.service_type_id, service_types_dm.service_type_id ], stream key: [ portfolios_dm.portfolio_id, portfolios_dm.service_type_id ] }
    ├── MergeExecutor { output: [ portfolios_dm.portfolio_id, portfolios_dm.service_type_id ], stream key: [ portfolios_dm.portfolio_id ] }
    └── MergeExecutor { output: [ service_types_dm.service_type_id, service_types_dm.type ], stream key: [ service_types_dm.service_type_id ] }

Fragment 24949 (Actor 119568,119567)
StreamProject { exprs: [portfolios_dm.portfolio_id, portfolios_dm.service_type_id] } { output: [ portfolios_dm.portfolio_id, portfolios_dm.service_type_id ], stream key: [ portfolios_dm.portfolio_id ] }
└── StreamFilter { predicate: IsNull(portfolios_dm.disabled_at) } { output: [ portfolios_dm.portfolio_id, portfolios_dm.service_type_id, portfolios_dm.disabled_at ], stream key: [ portfolios_dm.portfolio_id ] }
    └── StreamTableScan { table: portfolios_dm, columns: [portfolio_id, service_type_id, disabled_at] } { output: [ portfolios_dm.portfolio_id, portfolios_dm.service_type_id, portfolios_dm.disabled_at ], stream key: [ portfolios_dm.portfolio_id ] }
        ├── Upstream { output: [ portfolio_id, service_type_id, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ portfolio_id, service_type_id, disabled_at ], stream key: [] }

Fragment 24950 (Actor 118985,118986)
StreamTableScan { table: service_types_dm, columns: [service_type_id, type] } { output: [ service_types_dm.service_type_id, service_types_dm.type ], stream key: [ service_types_dm.service_type_id ] }
├── Upstream { output: [ service_type_id, type ], stream key: [] }
└── BatchPlanNode { output: [ service_type_id, type ], stream key: [] }

Fragment 24951 (Actor 119570,119569)
StreamProject { exprs: [portfolios_dm.portfolio_id, 'portfolio':Varchar, null:Varchar, null:Varchar, ConcatOp('status:':Varchar, Lower(labels_dm.name_en)) as $expr11, portfolios_dm.portfolio_id, portfolios_dm.status_label_id, null:Varchar, null:Varchar, null:Date, null:Date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 22:Int32] } { output: [ portfolios_dm.portfolio_id, 'portfolio':Varchar, null:Varchar, null:Varchar, $expr11, portfolios_dm.portfolio_id, portfolios_dm.status_label_id, null:Varchar, null:Varchar, null:Date, null:Date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 22:Int32 ], stream key: [ portfolios_dm.portfolio_id, portfolios_dm.status_label_id ] }
└── MergeExecutor { output: [ portfolios_dm.portfolio_id, labels_dm.name_en, portfolios_dm.status_label_id, labels_dm.label_id ], stream key: [ portfolios_dm.portfolio_id, portfolios_dm.status_label_id ] }

Fragment 24952 (Actor 119571,119572)
StreamSyncLogStore { output: [ portfolios_dm.portfolio_id, labels_dm.name_en, portfolios_dm.status_label_id, labels_dm.label_id ], stream key: [ portfolios_dm.portfolio_id, portfolios_dm.status_label_id ] }
└── StreamHashJoin { type: Inner, predicate: portfolios_dm.status_label_id = labels_dm.label_id } { output: [ portfolios_dm.portfolio_id, labels_dm.name_en, portfolios_dm.status_label_id, labels_dm.label_id ], stream key: [ portfolios_dm.portfolio_id, portfolios_dm.status_label_id ] }
    ├── MergeExecutor { output: [ portfolios_dm.portfolio_id, portfolios_dm.status_label_id ], stream key: [ portfolios_dm.portfolio_id ] }
    └── MergeExecutor { output: [ labels_dm.label_id, labels_dm.name_en ], stream key: [ labels_dm.label_id ] }

Fragment 24953 (Actor 119573,119574)
StreamProject { exprs: [portfolios_dm.portfolio_id, portfolios_dm.status_label_id] } { output: [ portfolios_dm.portfolio_id, portfolios_dm.status_label_id ], stream key: [ portfolios_dm.portfolio_id ] }
└── StreamFilter { predicate: IsNull(portfolios_dm.disabled_at) } { output: [ portfolios_dm.portfolio_id, portfolios_dm.status_label_id, portfolios_dm.disabled_at ], stream key: [ portfolios_dm.portfolio_id ] }
    └── StreamTableScan { table: portfolios_dm, columns: [portfolio_id, status_label_id, disabled_at] } { output: [ portfolios_dm.portfolio_id, portfolios_dm.status_label_id, portfolios_dm.disabled_at ], stream key: [ portfolios_dm.portfolio_id ] }
        ├── Upstream { output: [ portfolio_id, status_label_id, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ portfolio_id, status_label_id, disabled_at ], stream key: [] }

Fragment 24954 (Actor 119575,119576)
StreamTableScan { table: labels_dm, columns: [label_id, name_en] } { output: [ labels_dm.label_id, labels_dm.name_en ], stream key: [ labels_dm.label_id ] }
├── Upstream { output: [ label_id, name_en ], stream key: [] }
└── BatchPlanNode { output: [ label_id, name_en ], stream key: [] }

Fragment 24955 (Actor 119579,119580)
StreamProject { exprs: [entity_to_teams_dm.entity_id, 'portfolio':Varchar, null:Varchar, null:Varchar, ConcatOp('team:':Varchar, entity_to_teams_dm.team_id) as $expr12, entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, null:Varchar, entity_to_teams_dm.effective_start_date, null:Date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 23:Int32] }
├── output: [ entity_to_teams_dm.entity_id, 'portfolio':Varchar, null:Varchar, null:Varchar, $expr12, entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, null:Varchar, entity_to_teams_dm.effective_start_date, null:Date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 23:Int32 ]
├── stream key: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ]
└── MergeExecutor { output: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date, portfolios_dm.portfolio_id ], stream key: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }

Fragment 24956 (Actor 119578,119577)
StreamSyncLogStore { output: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date, portfolios_dm.portfolio_id ], stream key: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }
└── StreamHashJoin { type: Inner, predicate: entity_to_teams_dm.entity_id = portfolios_dm.portfolio_id } { output: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date, portfolios_dm.portfolio_id ], stream key: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }
    ├── MergeExecutor { output: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ], stream key: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }
    └── MergeExecutor { output: [ portfolios_dm.portfolio_id ], stream key: [ portfolios_dm.portfolio_id ] }

Fragment 24957 (Actor 119581,119582)
StreamProject { exprs: [entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date] } { output: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ], stream key: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }
└── StreamFilter { predicate: (entity_to_teams_dm.entity_type = 'PORTFOLIO':Varchar) AND IsNull(entity_to_teams_dm.disabled_at) AND IsNull(entity_to_teams_dm.effective_end_date) } { output: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date, entity_to_teams_dm.effective_end_date, entity_to_teams_dm.disabled_at ], stream key: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }
    └── StreamTableScan { table: entity_to_teams_dm, columns: [team_id, entity_id, entity_type, effective_start_date, effective_end_date, disabled_at] } { output: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date, entity_to_teams_dm.effective_end_date, entity_to_teams_dm.disabled_at ], stream key: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }
        ├── Upstream { output: [ team_id, entity_id, entity_type, effective_start_date, effective_end_date, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ team_id, entity_id, entity_type, effective_start_date, effective_end_date, disabled_at ], stream key: [] }

Fragment 24958 (Actor 119583,119584)
StreamProject { exprs: [portfolios_dm.portfolio_id] } { output: [ portfolios_dm.portfolio_id ], stream key: [ portfolios_dm.portfolio_id ] }
└── StreamFilter { predicate: IsNull(portfolios_dm.disabled_at) } { output: [ portfolios_dm.portfolio_id, portfolios_dm.disabled_at ], stream key: [ portfolios_dm.portfolio_id ] }
    └── StreamTableScan { table: portfolios_dm, columns: [portfolio_id, disabled_at] } { output: [ portfolios_dm.portfolio_id, portfolios_dm.disabled_at ], stream key: [ portfolios_dm.portfolio_id ] }
        ├── Upstream { output: [ portfolio_id, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ portfolio_id, disabled_at ], stream key: [] }

Fragment 24959 (Actor 119585,119586)
StreamProject { exprs: [clients_portfolios_dm.portfolio_id, 'portfolio':Varchar, null:Varchar, null:Varchar, ConcatOp('team:':Varchar, entity_to_teams_dm.team_id) as $expr13, clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, entity_to_teams_dm.team_id, entity_to_teams_dm.entity_type, clients_portfolios_dm.effective_start_date, entity_to_teams_dm.effective_start_date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 24:Int32] }
├── output: [ clients_portfolios_dm.portfolio_id, 'portfolio':Varchar, null:Varchar, null:Varchar, $expr13, clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, entity_to_teams_dm.team_id, entity_to_teams_dm.entity_type, clients_portfolios_dm.effective_start_date, entity_to_teams_dm.effective_start_date, null:Int64, null:Int64, null:Int32, null:Int32, null:Int32, 24:Int32 ]
├── stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, entity_to_teams_dm.team_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ]
└── MergeExecutor { output: [ clients_portfolios_dm.portfolio_id, entity_to_teams_dm.team_id, clients_portfolios_dm.client_id, clients_portfolios_dm.effective_start_date, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, entity_to_teams_dm.team_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }

Fragment 24960 (Actor 119587,119588)
StreamSyncLogStore { output: [ clients_portfolios_dm.portfolio_id, entity_to_teams_dm.team_id, clients_portfolios_dm.client_id, clients_portfolios_dm.effective_start_date, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, entity_to_teams_dm.team_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }
└── StreamHashJoin { type: Inner, predicate: clients_portfolios_dm.client_id = entity_to_teams_dm.entity_id } { output: [ clients_portfolios_dm.portfolio_id, entity_to_teams_dm.team_id, clients_portfolios_dm.client_id, clients_portfolios_dm.effective_start_date, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, entity_to_teams_dm.team_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }
    ├── MergeExecutor { output: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ] }
    └── MergeExecutor { output: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ], stream key: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }

Fragment 24961 (Actor 119589,119590)
StreamProject { exprs: [clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date] } { output: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ] }
└── StreamFilter { predicate: IsNull(clients_portfolios_dm.disabled_at) AND IsNull(clients_portfolios_dm.effective_end_date) } { output: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, clients_portfolios_dm.disabled_at, clients_portfolios_dm.effective_end_date ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ] }
    └── StreamTableScan { table: clients_portfolios_dm, columns: [client_id, portfolio_id, effective_start_date, disabled_at, effective_end_date] } { output: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, clients_portfolios_dm.disabled_at, clients_portfolios_dm.effective_end_date ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ] }
        ├── Upstream { output: [ client_id, portfolio_id, effective_start_date, disabled_at, effective_end_date ], stream key: [] }
        └── BatchPlanNode { output: [ client_id, portfolio_id, effective_start_date, disabled_at, effective_end_date ], stream key: [] }

Fragment 24962 (Actor 119592,119591)
StreamProject { exprs: [entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date] } { output: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ], stream key: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }
└── StreamFilter { predicate: (entity_to_teams_dm.entity_type = 'CLIENT':Varchar) AND IsNull(entity_to_teams_dm.disabled_at) AND IsNull(entity_to_teams_dm.effective_end_date) } { output: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date, entity_to_teams_dm.effective_end_date, entity_to_teams_dm.disabled_at ], stream key: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }
    └── StreamTableScan { table: entity_to_teams_dm, columns: [team_id, entity_id, entity_type, effective_start_date, effective_end_date, disabled_at] } { output: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date, entity_to_teams_dm.effective_end_date, entity_to_teams_dm.disabled_at ], stream key: [ entity_to_teams_dm.team_id, entity_to_teams_dm.entity_id, entity_to_teams_dm.entity_type, entity_to_teams_dm.effective_start_date ] }
        ├── Upstream { output: [ team_id, entity_id, entity_type, effective_start_date, effective_end_date, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ team_id, entity_id, entity_type, effective_start_date, effective_end_date, disabled_at ], stream key: [] }

Fragment 24963 (Actor 119596,119595)
StreamProject { exprs: [active_identifier_edges_mv_next.target_entity_id, active_identifier_edges_mv_next.target_entity_type, olap_reference_identifier_terms_mv.val, olap_reference_identifier_terms_mv.ar_val, null:Varchar, active_identifier_edges_mv_next.accounts_dm.account_id, active_identifier_edges_mv_next.null:Varchar, active_identifier_edges_mv_next.owner_entity_id, active_identifier_edges_mv_next.owner_entity_type, active_identifier_edges_mv_next.null:Date, null:Date, olap_reference_identifier_terms_mv.reference_identifiers.id, olap_reference_identifier_terms_mv._rw_projected_row_id, active_identifier_edges_mv_next.null:Int32, active_identifier_edges_mv_next.$src, null:Int32, 25:Int32] }
├── output: [ active_identifier_edges_mv_next.target_entity_id, active_identifier_edges_mv_next.target_entity_type, olap_reference_identifier_terms_mv.val, olap_reference_identifier_terms_mv.ar_val, null:Varchar, active_identifier_edges_mv_next.accounts_dm.account_id, active_identifier_edges_mv_next.null:Varchar, active_identifier_edges_mv_next.owner_entity_id, active_identifier_edges_mv_next.owner_entity_type, active_identifier_edges_mv_next.null:Date, null:Date, olap_reference_identifier_terms_mv.reference_identifiers.id, olap_reference_identifier_terms_mv._rw_projected_row_id, active_identifier_edges_mv_next.null:Int32, active_identifier_edges_mv_next.$src, null:Int32, 25:Int32 ]
├── stream key: [ active_identifier_edges_mv_next.accounts_dm.account_id, active_identifier_edges_mv_next.null:Varchar, active_identifier_edges_mv_next.null:Date, active_identifier_edges_mv_next.null:Int32, active_identifier_edges_mv_next.$src, olap_reference_identifier_terms_mv.reference_identifiers.id, olap_reference_identifier_terms_mv._rw_projected_row_id, active_identifier_edges_mv_next.owner_entity_id, active_identifier_edges_mv_next.owner_entity_type ]
└── MergeExecutor
    ├── output: [ active_identifier_edges_mv_next.target_entity_id, active_identifier_edges_mv_next.target_entity_type, olap_reference_identifier_terms_mv.val, olap_reference_identifier_terms_mv.ar_val, active_identifier_edges_mv_next.accounts_dm.account_id, active_identifier_edges_mv_next.null:Varchar, active_identifier_edges_mv_next.null:Date, active_identifier_edges_mv_next.null:Int32, active_identifier_edges_mv_next.$src, active_identifier_edges_mv_next.owner_entity_id, active_identifier_edges_mv_next.owner_entity_type, olap_reference_identifier_terms_mv.reference_identifiers.id, olap_reference_identifier_terms_mv._rw_projected_row_id ]
    └── stream key: [ active_identifier_edges_mv_next.accounts_dm.account_id, active_identifier_edges_mv_next.null:Varchar, active_identifier_edges_mv_next.null:Date, active_identifier_edges_mv_next.null:Int32, active_identifier_edges_mv_next.$src, olap_reference_identifier_terms_mv.reference_identifiers.id, olap_reference_identifier_terms_mv._rw_projected_row_id, active_identifier_edges_mv_next.owner_entity_id, active_identifier_edges_mv_next.owner_entity_type ]

Fragment 24964 (Actor 119594,119593)
StreamSyncLogStore
├── output: [ active_identifier_edges_mv_next.target_entity_id, active_identifier_edges_mv_next.target_entity_type, olap_reference_identifier_terms_mv.val, olap_reference_identifier_terms_mv.ar_val, active_identifier_edges_mv_next.accounts_dm.account_id, active_identifier_edges_mv_next.null:Varchar, active_identifier_edges_mv_next.null:Date, active_identifier_edges_mv_next.null:Int32, active_identifier_edges_mv_next.$src, active_identifier_edges_mv_next.owner_entity_id, active_identifier_edges_mv_next.owner_entity_type, olap_reference_identifier_terms_mv.reference_identifiers.id, olap_reference_identifier_terms_mv._rw_projected_row_id ]
├── stream key: [ active_identifier_edges_mv_next.accounts_dm.account_id, active_identifier_edges_mv_next.null:Varchar, active_identifier_edges_mv_next.null:Date, active_identifier_edges_mv_next.null:Int32, active_identifier_edges_mv_next.$src, olap_reference_identifier_terms_mv.reference_identifiers.id, olap_reference_identifier_terms_mv._rw_projected_row_id, active_identifier_edges_mv_next.owner_entity_id, active_identifier_edges_mv_next.owner_entity_type ]
└── StreamHashJoin { type: Inner, predicate: active_identifier_edges_mv_next.owner_entity_id = olap_reference_identifier_terms_mv.owner_entity_id AND active_identifier_edges_mv_next.owner_entity_type = olap_reference_identifier_terms_mv.owner_entity_type }
    ├── output: [ active_identifier_edges_mv_next.target_entity_id, active_identifier_edges_mv_next.target_entity_type, olap_reference_identifier_terms_mv.val, olap_reference_identifier_terms_mv.ar_val, active_identifier_edges_mv_next.accounts_dm.account_id, active_identifier_edges_mv_next.null:Varchar, active_identifier_edges_mv_next.null:Date, active_identifier_edges_mv_next.null:Int32, active_identifier_edges_mv_next.$src, active_identifier_edges_mv_next.owner_entity_id, active_identifier_edges_mv_next.owner_entity_type, olap_reference_identifier_terms_mv.reference_identifiers.id, olap_reference_identifier_terms_mv._rw_projected_row_id ]
    ├── stream key: [ active_identifier_edges_mv_next.accounts_dm.account_id, active_identifier_edges_mv_next.null:Varchar, active_identifier_edges_mv_next.null:Date, active_identifier_edges_mv_next.null:Int32, active_identifier_edges_mv_next.$src, olap_reference_identifier_terms_mv.reference_identifiers.id, olap_reference_identifier_terms_mv._rw_projected_row_id, active_identifier_edges_mv_next.owner_entity_id, active_identifier_edges_mv_next.owner_entity_type ]
    ├── MergeExecutor { output: [ active_identifier_edges_mv_next.target_entity_id, active_identifier_edges_mv_next.target_entity_type, active_identifier_edges_mv_next.owner_entity_id, active_identifier_edges_mv_next.owner_entity_type, active_identifier_edges_mv_next.accounts_dm.account_id, active_identifier_edges_mv_next.null:Varchar, active_identifier_edges_mv_next.null:Date, active_identifier_edges_mv_next.null:Int32, active_identifier_edges_mv_next.$src ], stream key: [ active_identifier_edges_mv_next.accounts_dm.account_id, active_identifier_edges_mv_next.null:Varchar, active_identifier_edges_mv_next.null:Date, active_identifier_edges_mv_next.null:Int32, active_identifier_edges_mv_next.$src ] }
    └── MergeExecutor { output: [ olap_reference_identifier_terms_mv.owner_entity_id, olap_reference_identifier_terms_mv.owner_entity_type, olap_reference_identifier_terms_mv.val, olap_reference_identifier_terms_mv.ar_val, olap_reference_identifier_terms_mv.reference_identifiers.id, olap_reference_identifier_terms_mv._rw_projected_row_id ], stream key: [ olap_reference_identifier_terms_mv.reference_identifiers.id, olap_reference_identifier_terms_mv._rw_projected_row_id ] }

Fragment 24965 (Actor 119597,119598)
StreamFilter { predicate: In(active_identifier_edges_mv_next.target_entity_type, 'account':Varchar, 'client':Varchar, 'portfolio':Varchar) }
├── output: [ active_identifier_edges_mv_next.target_entity_id, active_identifier_edges_mv_next.target_entity_type, active_identifier_edges_mv_next.owner_entity_id, active_identifier_edges_mv_next.owner_entity_type, active_identifier_edges_mv_next.accounts_dm.account_id, active_identifier_edges_mv_next.null:Varchar, active_identifier_edges_mv_next.null:Date, active_identifier_edges_mv_next.null:Int32, active_identifier_edges_mv_next.$src ]
├── stream key: [ active_identifier_edges_mv_next.accounts_dm.account_id, active_identifier_edges_mv_next.null:Varchar, active_identifier_edges_mv_next.null:Date, active_identifier_edges_mv_next.null:Int32, active_identifier_edges_mv_next.$src ]
└── StreamTableScan { table: active_identifier_edges_mv_next, columns: [target_entity_id, target_entity_type, owner_entity_id, owner_entity_type, accounts_dm.account_id, null:Varchar, null:Date, null:Int32, $src] }
    ├── output: [ active_identifier_edges_mv_next.target_entity_id, active_identifier_edges_mv_next.target_entity_type, active_identifier_edges_mv_next.owner_entity_id, active_identifier_edges_mv_next.owner_entity_type, active_identifier_edges_mv_next.accounts_dm.account_id, active_identifier_edges_mv_next.null:Varchar, active_identifier_edges_mv_next.null:Date, active_identifier_edges_mv_next.null:Int32, active_identifier_edges_mv_next.$src ]
    ├── stream key: [ active_identifier_edges_mv_next.accounts_dm.account_id, active_identifier_edges_mv_next.null:Varchar, active_identifier_edges_mv_next.null:Date, active_identifier_edges_mv_next.null:Int32, active_identifier_edges_mv_next.$src ]
    ├── Upstream { output: [ target_entity_id, target_entity_type, owner_entity_id, owner_entity_type, accounts_dm.account_id, null:Varchar, null:Date, null:Int32, $src ], stream key: [] }
    └── BatchPlanNode { output: [ target_entity_id, target_entity_type, owner_entity_id, owner_entity_type, accounts_dm.account_id, null:Varchar, null:Date, null:Int32, $src ], stream key: [] }

Fragment 24966 (Actor 119599,119600)
StreamTableScan { table: olap_reference_identifier_terms_mv, columns: [owner_entity_id, owner_entity_type, val, ar_val, reference_identifiers.id, _rw_projected_row_id] } { output: [ olap_reference_identifier_terms_mv.owner_entity_id, olap_reference_identifier_terms_mv.owner_entity_type, olap_reference_identifier_terms_mv.val, olap_reference_identifier_terms_mv.ar_val, olap_reference_identifier_terms_mv.reference_identifiers.id, olap_reference_identifier_terms_mv._rw_projected_row_id ], stream key: [ olap_reference_identifier_terms_mv.reference_identifiers.id, olap_reference_identifier_terms_mv._rw_projected_row_id ] }
├── Upstream { output: [ owner_entity_id, owner_entity_type, val, ar_val, reference_identifiers.id, _rw_projected_row_id ], stream key: [] }
└── BatchPlanNode { output: [ owner_entity_id, owner_entity_type, val, ar_val, reference_identifiers.id, _rw_projected_row_id ], stream key: [] }

Fragment 24967 (Actor 119603,119604)
StreamProject { exprs: [party_identifier_edges_mv_next.target_entity_id, party_identifier_edges_mv_next.target_entity_type, party_reference_identifier_terms_mv.val, party_reference_identifier_terms_mv.ar_val, null:Varchar, party_identifier_edges_mv_next.owner_entity_id, party_identifier_edges_mv_next.target_entity_id, party_identifier_edges_mv_next.owner_entity_type, null:Varchar, null:Date, null:Date, party_reference_identifier_terms_mv.reference_identifiers_next.id, party_reference_identifier_terms_mv._rw_projected_row_id, party_identifier_edges_mv_next.party_holder_edges_mv_next.$src, null:Int32, null:Int32, 26:Int32] }
├── output: [ party_identifier_edges_mv_next.target_entity_id, party_identifier_edges_mv_next.target_entity_type, party_reference_identifier_terms_mv.val, party_reference_identifier_terms_mv.ar_val, null:Varchar, party_identifier_edges_mv_next.owner_entity_id, party_identifier_edges_mv_next.target_entity_id, party_identifier_edges_mv_next.owner_entity_type, null:Varchar, null:Date, null:Date, party_reference_identifier_terms_mv.reference_identifiers_next.id, party_reference_identifier_terms_mv._rw_projected_row_id, party_identifier_edges_mv_next.party_holder_edges_mv_next.$src, null:Int32, null:Int32, 26:Int32 ]
├── stream key: [ party_identifier_edges_mv_next.owner_entity_id, party_identifier_edges_mv_next.target_entity_id, party_identifier_edges_mv_next.party_holder_edges_mv_next.$src, party_reference_identifier_terms_mv.reference_identifiers_next.id, party_reference_identifier_terms_mv._rw_projected_row_id, party_identifier_edges_mv_next.owner_entity_type ]
└── MergeExecutor
    ├── output: [ party_identifier_edges_mv_next.target_entity_id, party_identifier_edges_mv_next.target_entity_type, party_reference_identifier_terms_mv.val, party_reference_identifier_terms_mv.ar_val, party_identifier_edges_mv_next.owner_entity_id, party_identifier_edges_mv_next.party_holder_edges_mv_next.$src, party_identifier_edges_mv_next.owner_entity_type, party_reference_identifier_terms_mv.reference_identifiers_next.id, party_reference_identifier_terms_mv._rw_projected_row_id ]
    └── stream key: [ party_identifier_edges_mv_next.owner_entity_id, party_identifier_edges_mv_next.target_entity_id, party_identifier_edges_mv_next.party_holder_edges_mv_next.$src, party_reference_identifier_terms_mv.reference_identifiers_next.id, party_reference_identifier_terms_mv._rw_projected_row_id, party_identifier_edges_mv_next.owner_entity_type ]

Fragment 24968 (Actor 119601,119602)
StreamSyncLogStore
├── output: [ party_identifier_edges_mv_next.target_entity_id, party_identifier_edges_mv_next.target_entity_type, party_reference_identifier_terms_mv.val, party_reference_identifier_terms_mv.ar_val, party_identifier_edges_mv_next.owner_entity_id, party_identifier_edges_mv_next.party_holder_edges_mv_next.$src, party_identifier_edges_mv_next.owner_entity_type, party_reference_identifier_terms_mv.reference_identifiers_next.id, party_reference_identifier_terms_mv._rw_projected_row_id ]
├── stream key: [ party_identifier_edges_mv_next.owner_entity_id, party_identifier_edges_mv_next.target_entity_id, party_identifier_edges_mv_next.party_holder_edges_mv_next.$src, party_reference_identifier_terms_mv.reference_identifiers_next.id, party_reference_identifier_terms_mv._rw_projected_row_id, party_identifier_edges_mv_next.owner_entity_type ]
└── StreamHashJoin { type: Inner, predicate: party_identifier_edges_mv_next.owner_entity_id = party_reference_identifier_terms_mv.owner_entity_id AND party_identifier_edges_mv_next.owner_entity_type = party_reference_identifier_terms_mv.owner_entity_type }
    ├── output: [ party_identifier_edges_mv_next.target_entity_id, party_identifier_edges_mv_next.target_entity_type, party_reference_identifier_terms_mv.val, party_reference_identifier_terms_mv.ar_val, party_identifier_edges_mv_next.owner_entity_id, party_identifier_edges_mv_next.party_holder_edges_mv_next.$src, party_identifier_edges_mv_next.owner_entity_type, party_reference_identifier_terms_mv.reference_identifiers_next.id, party_reference_identifier_terms_mv._rw_projected_row_id ]
    ├── stream key: [ party_identifier_edges_mv_next.owner_entity_id, party_identifier_edges_mv_next.target_entity_id, party_identifier_edges_mv_next.party_holder_edges_mv_next.$src, party_reference_identifier_terms_mv.reference_identifiers_next.id, party_reference_identifier_terms_mv._rw_projected_row_id, party_identifier_edges_mv_next.owner_entity_type ]
    ├── MergeExecutor { output: [ party_identifier_edges_mv_next.target_entity_id, party_identifier_edges_mv_next.target_entity_type, party_identifier_edges_mv_next.owner_entity_id, party_identifier_edges_mv_next.owner_entity_type, party_identifier_edges_mv_next.party_holder_edges_mv_next.$src ], stream key: [ party_identifier_edges_mv_next.owner_entity_id, party_identifier_edges_mv_next.target_entity_id, party_identifier_edges_mv_next.party_holder_edges_mv_next.$src ] }
    └── MergeExecutor { output: [ party_reference_identifier_terms_mv.owner_entity_id, party_reference_identifier_terms_mv.owner_entity_type, party_reference_identifier_terms_mv.val, party_reference_identifier_terms_mv.ar_val, party_reference_identifier_terms_mv.reference_identifiers_next.id, party_reference_identifier_terms_mv._rw_projected_row_id ], stream key: [ party_reference_identifier_terms_mv.reference_identifiers_next.id, party_reference_identifier_terms_mv._rw_projected_row_id ] }

Fragment 24969 (Actor 119605,119606)
StreamTableScan { table: party_identifier_edges_mv_next, columns: [target_entity_id, target_entity_type, owner_entity_id, owner_entity_type, party_holder_edges_mv_next.$src] } { output: [ party_identifier_edges_mv_next.target_entity_id, party_identifier_edges_mv_next.target_entity_type, party_identifier_edges_mv_next.owner_entity_id, party_identifier_edges_mv_next.owner_entity_type, party_identifier_edges_mv_next.party_holder_edges_mv_next.$src ], stream key: [ party_identifier_edges_mv_next.owner_entity_id, party_identifier_edges_mv_next.target_entity_id, party_identifier_edges_mv_next.party_holder_edges_mv_next.$src ] }
├── Upstream { output: [ target_entity_id, target_entity_type, owner_entity_id, owner_entity_type, party_holder_edges_mv_next.$src ], stream key: [] }
└── BatchPlanNode { output: [ target_entity_id, target_entity_type, owner_entity_id, owner_entity_type, party_holder_edges_mv_next.$src ], stream key: [] }

Fragment 24970 (Actor 119607,119608)
StreamTableScan { table: party_reference_identifier_terms_mv, columns: [owner_entity_id, owner_entity_type, val, ar_val, reference_identifiers_next.id, _rw_projected_row_id] } { output: [ party_reference_identifier_terms_mv.owner_entity_id, party_reference_identifier_terms_mv.owner_entity_type, party_reference_identifier_terms_mv.val, party_reference_identifier_terms_mv.ar_val, party_reference_identifier_terms_mv.reference_identifiers_next.id, party_reference_identifier_terms_mv._rw_projected_row_id ], stream key: [ party_reference_identifier_terms_mv.reference_identifiers_next.id, party_reference_identifier_terms_mv._rw_projected_row_id ] }
├── Upstream { output: [ owner_entity_id, owner_entity_type, val, ar_val, reference_identifiers_next.id, _rw_projected_row_id ], stream key: [] }
└── BatchPlanNode { output: [ owner_entity_id, owner_entity_type, val, ar_val, reference_identifiers_next.id, _rw_projected_row_id ], stream key: [] }