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

← cluster search objects active_identifier_edges_mv explain
Overview Objects Graph History
materialized view · search.active_identifier_edges_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 lookups
221 operators
Materialize · search.active_identifier_edges_mv
0% idle 2 actors
Union
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · party_holder_edges_mv_next
2 actors
StreamScan · party_holder_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 · 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
SyncLogStore · Inner · account_to_portfolios_dm.portfolio_id = portfolios_dm.portf…
2 actors
HashJoin · Inner · account_to_portfolios_dm.portfolio_id = portfolios_dm.portf… 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 · 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
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
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
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
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
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
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
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · account_to_portfolios_dm.portfolio_id = portfolios_dm.portf…
2 actors
HashJoin · Inner · account_to_portfolios_dm.portfolio_id = portfolios_dm.portf… 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
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
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
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
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 · 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_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_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.active_identifier_edges_mv Materialize search.active_identifie… idle · 2 actors Union Union idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · party_holder_edges_mv_next Project party_holder_edges_mv_n… — · 2 actors StreamScan · party_holder_edges_mv_next StreamScan party_holder_edges_mv_n… 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 · 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 SyncLogStore · Inner · account_to_portfolios_dm.portfolio_id = portfolios_dm.portf… SyncLogStore Inner · account_to_port… — · 2 actors HashJoin · Inner · account_to_portfolios_dm.portfolio_id = portfolios_dm.portf… HashJoin Inner · account_to_port… 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 · 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 Project — · 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 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 Project — · 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 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 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 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 Project — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · account_to_portfolios_dm.portfolio_id = portfolios_dm.portf… SyncLogStore Inner · account_to_port… — · 2 actors HashJoin · Inner · account_to_portfolios_dm.portfolio_id = portfolios_dm.portf… HashJoin Inner · account_to_port… 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 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 Project — · 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 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 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 · 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_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_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 24818 (Actor 118500,118499)
StreamMaterialize { columns: [target_entity_id, target_entity_type, owner_entity_id, owner_entity_type, null:Int32(hidden), accounts_dm.account_id(hidden), null:Varchar(hidden), null:Date(hidden), $src(hidden)], stream_key: [accounts_dm.account_id, null:Varchar, null:Date, null:Int32, $src], pk_columns: [accounts_dm.account_id, null:Varchar, null:Date, null:Int32, $src], pk_conflict: NoCheck }
├── output: [ accounts_dm.account_id, 'account':Varchar, accounts_dm.account_id, 'account':Varchar, null:Int32, accounts_dm.account_id, null:Varchar, null:Date, $src ]
├── stream key: [ accounts_dm.account_id, null:Varchar, null:Date, null:Int32, $src ]
└── StreamUnion { all: true } { output: [ accounts_dm.account_id, 'account':Varchar, accounts_dm.account_id, 'account':Varchar, null:Int32, accounts_dm.account_id, null:Varchar, null:Date, $src ], stream key: [ accounts_dm.account_id, null:Varchar, null:Date, null:Int32, $src ] }
    ├── MergeExecutor { output: [ accounts_dm.account_id, 'account':Varchar, accounts_dm.account_id, 'account':Varchar, null:Int32, accounts_dm.account_id, null:Varchar, null:Date, 0:Int32 ], stream key: [ accounts_dm.account_id ] }
    ├── MergeExecutor { output: [ clients_dm.id, 'client':Varchar, clients_dm.id, 'client':Varchar, null:Int32, clients_dm.id, null:Varchar, null:Date, 1:Int32 ], stream key: [ clients_dm.id ] }
    ├── MergeExecutor { output: [ portfolios_dm.portfolio_id, 'portfolio':Varchar, portfolios_dm.portfolio_id, 'portfolio':Varchar, null:Int32, portfolios_dm.portfolio_id, null:Varchar, null:Date, 2:Int32 ], stream key: [ portfolios_dm.portfolio_id ] }
    ├── MergeExecutor
    │   ├── output: [ accounts_to_clients_dm.account_id, 'account':Varchar, accounts_to_clients_dm.client_id, 'client':Varchar, null:Int32, accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, 3:Int32 ]
    │   └── stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ]
    ├── MergeExecutor
    │   ├── output: [ account_to_portfolios_dm.account_id, 'account':Varchar, account_to_portfolios_dm.portfolio_id, 'portfolio':Varchar, null:Int32, account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date, 4:Int32 ]
    │   └── stream key: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ]
    ├── MergeExecutor
    │   ├── output: [ accounts_to_clients_dm.client_id, 'client':Varchar, accounts_to_clients_dm.account_id, 'account':Varchar, null:Int32, accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, 5:Int32 ]
    │   └── stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ]
    ├── MergeExecutor
    │   ├── output: [ clients_portfolios_dm.client_id, 'client':Varchar, clients_portfolios_dm.portfolio_id, 'portfolio':Varchar, null:Int32, clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, 6:Int32 ]
    │   └── stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ]
    ├── MergeExecutor
    │   ├── output: [ clients_portfolios_dm.portfolio_id, 'portfolio':Varchar, clients_portfolios_dm.client_id, 'client':Varchar, null:Int32, clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, 7:Int32 ]
    │   └── stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ]
    ├── MergeExecutor
    │   ├── output: [ account_to_portfolios_dm.portfolio_id, 'portfolio':Varchar, account_to_portfolios_dm.account_id, 'account':Varchar, null:Int32, account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date, 8:Int32 ]
    │   └── stream key: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ]
    └── MergeExecutor
        ├── output: [ party_holder_edges_mv_next.party_id, 'party':Varchar, party_holder_edges_mv_next.entity_id, party_holder_edges_mv_next.entity_type, party_holder_edges_mv_next.$src, party_holder_edges_mv_next.party_id, party_holder_edges_mv_next.entity_id, null:Date, 9:Int32 ]
        └── stream key: [ party_holder_edges_mv_next.party_id, party_holder_edges_mv_next.entity_id, party_holder_edges_mv_next.$src ]

Fragment 24819 (Actor 118551,118552)
StreamProject { exprs: [accounts_dm.account_id, 'account':Varchar, accounts_dm.account_id, 'account':Varchar, null:Int32, accounts_dm.account_id, null:Varchar, null:Date, 0:Int32] } { output: [ accounts_dm.account_id, 'account':Varchar, accounts_dm.account_id, 'account':Varchar, null:Int32, accounts_dm.account_id, null:Varchar, null:Date, 0:Int32 ], stream key: [ accounts_dm.account_id ] }
└── StreamFilter { predicate: IsNull(accounts_dm.disabled_at) } { output: [ accounts_dm.account_id, accounts_dm.disabled_at ], stream key: [ accounts_dm.account_id ] }
    └── StreamTableScan { table: accounts_dm, columns: [account_id, disabled_at] } { output: [ accounts_dm.account_id, accounts_dm.disabled_at ], stream key: [ accounts_dm.account_id ] }
        ├── Upstream { output: [ account_id, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ account_id, disabled_at ], stream key: [] }

Fragment 24820 (Actor 118554,118553)
StreamProject { exprs: [clients_dm.id, 'client':Varchar, clients_dm.id, 'client':Varchar, null:Int32, clients_dm.id, null:Varchar, null:Date, 1:Int32] } { output: [ clients_dm.id, 'client':Varchar, clients_dm.id, 'client':Varchar, null:Int32, clients_dm.id, null:Varchar, null:Date, 1:Int32 ], 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 24821 (Actor 118555,118556)
StreamProject { exprs: [portfolios_dm.portfolio_id, 'portfolio':Varchar, portfolios_dm.portfolio_id, 'portfolio':Varchar, null:Int32, portfolios_dm.portfolio_id, null:Varchar, null:Date, 2:Int32] }
├── output: [ portfolios_dm.portfolio_id, 'portfolio':Varchar, portfolios_dm.portfolio_id, 'portfolio':Varchar, null:Int32, portfolios_dm.portfolio_id, null:Varchar, null:Date, 2:Int32 ]
├── 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 24822 (Actor 118502,118501)
StreamProject { exprs: [accounts_to_clients_dm.account_id, 'account':Varchar, accounts_to_clients_dm.client_id, 'client':Varchar, null:Int32, accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, 3:Int32] }
├── output: [ accounts_to_clients_dm.account_id, 'account':Varchar, accounts_to_clients_dm.client_id, 'client':Varchar, null:Int32, accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, 3:Int32 ]
├── 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, 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 24823 (Actor 118504,118503)
StreamSyncLogStore { output: [ accounts_to_clients_dm.account_id, 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, 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, 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: [ clients_dm.id ], stream key: [ clients_dm.id ] }

Fragment 24824 (Actor 118505,118506)
StreamSyncLogStore { output: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, 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.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.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 ], stream key: [ accounts_dm.account_id ] }

Fragment 24825 (Actor 118558,118557)
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) 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.effective_start_date, accounts_to_clients_dm.disabled_at, accounts_to_clients_dm.effective_end_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, effective_start_date, disabled_at, effective_end_date] }
        ├── 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, accounts_to_clients_dm.effective_end_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, effective_start_date, disabled_at, effective_end_date ], stream key: [] }
        └── BatchPlanNode { output: [ account_id, client_id, effective_start_date, disabled_at, effective_end_date ], stream key: [] }

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

Fragment 24827 (Actor 118562,118561)
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 24828 (Actor 118508,118507)
StreamProject { exprs: [account_to_portfolios_dm.account_id, 'account':Varchar, account_to_portfolios_dm.portfolio_id, 'portfolio':Varchar, null:Int32, account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date, 4:Int32] }
├── output: [ account_to_portfolios_dm.account_id, 'account':Varchar, account_to_portfolios_dm.portfolio_id, 'portfolio':Varchar, null:Int32, account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date, 4:Int32 ]
├── 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, portfolios_dm.portfolio_id ], stream key: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ] }

Fragment 24829 (Actor 118509,118510)
StreamSyncLogStore { output: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date, portfolios_dm.portfolio_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.portfolio_id = portfolios_dm.portfolio_id }
    ├── output: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date, portfolios_dm.portfolio_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, 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: [ portfolios_dm.portfolio_id ], stream key: [ portfolios_dm.portfolio_id ] }

Fragment 24830 (Actor 118511,118512)
StreamSyncLogStore { output: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_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.account_id, account_to_portfolios_dm.portfolio_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 ], stream key: [ accounts_dm.account_id ] }

Fragment 24831 (Actor 118563,118564)
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) AND IsNull(account_to_portfolios_dm.effective_end_date) }
    ├── 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, account_to_portfolios_dm.effective_end_date ]
    ├── 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, effective_end_date] }
        ├── 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, account_to_portfolios_dm.effective_end_date ]
        ├── 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, effective_end_date ], stream key: [] }
        └── BatchPlanNode { output: [ account_id, portfolio_id, effective_start_date, disabled_at, effective_end_date ], stream key: [] }

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

Fragment 24833 (Actor 118568,118567)
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 24834 (Actor 118515,118516)
StreamProject { exprs: [accounts_to_clients_dm.client_id, 'client':Varchar, accounts_to_clients_dm.account_id, 'account':Varchar, null:Int32, accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, 5:Int32] }
├── output: [ accounts_to_clients_dm.client_id, 'client':Varchar, accounts_to_clients_dm.account_id, 'account':Varchar, null:Int32, accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, 5:Int32 ]
├── 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, 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 24835 (Actor 118513,118514)
StreamSyncLogStore { output: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, 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.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.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_dm.account_id ], stream key: [ accounts_dm.account_id ] }

Fragment 24836 (Actor 118518,118517)
StreamSyncLogStore { output: [ accounts_to_clients_dm.account_id, 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, 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 ], stream key: [ clients_dm.id ] }

Fragment 24837 (Actor 118570,118569)
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) 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.effective_start_date, accounts_to_clients_dm.disabled_at, accounts_to_clients_dm.effective_end_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, effective_start_date, disabled_at, effective_end_date] }
        ├── 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, accounts_to_clients_dm.effective_end_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, effective_start_date, disabled_at, effective_end_date ], stream key: [] }
        └── BatchPlanNode { output: [ account_id, client_id, effective_start_date, disabled_at, effective_end_date ], stream key: [] }

Fragment 24838 (Actor 118571,118572)
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 24839 (Actor 118520,118519)
StreamProject { exprs: [accounts_dm.account_id] } { output: [ accounts_dm.account_id ], stream key: [ accounts_dm.account_id ] }
└── StreamFilter { predicate: IsNull(accounts_dm.disabled_at) } { output: [ accounts_dm.account_id, accounts_dm.disabled_at ], stream key: [ accounts_dm.account_id ] }
    └── StreamTableScan { table: accounts_dm, columns: [account_id, disabled_at] } { output: [ accounts_dm.account_id, accounts_dm.disabled_at ], stream key: [ accounts_dm.account_id ] }
        ├── Upstream { output: [ account_id, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ account_id, disabled_at ], stream key: [] }

Fragment 24840 (Actor 118523,118524)
StreamProject { exprs: [clients_portfolios_dm.client_id, 'client':Varchar, clients_portfolios_dm.portfolio_id, 'portfolio':Varchar, null:Int32, clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, 6:Int32] }
├── output: [ clients_portfolios_dm.client_id, 'client':Varchar, clients_portfolios_dm.portfolio_id, 'portfolio':Varchar, null:Int32, clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, 6:Int32 ]
├── 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, portfolios_dm.portfolio_id ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ] }

Fragment 24841 (Actor 118526,118525)
StreamSyncLogStore { output: [ clients_portfolios_dm.client_id, 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, 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, clients_dm.id ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ] }
    └── MergeExecutor { output: [ portfolios_dm.portfolio_id ], stream key: [ portfolios_dm.portfolio_id ] }

Fragment 24842 (Actor 118527,118528)
StreamSyncLogStore { output: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_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.client_id, clients_portfolios_dm.portfolio_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 ], stream key: [ clients_dm.id ] }

Fragment 24843 (Actor 118522,118521)
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 24844 (Actor 118574,118573)
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 24845 (Actor 118576,118575)
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 24846 (Actor 118531,118532)
StreamProject { exprs: [clients_portfolios_dm.portfolio_id, 'portfolio':Varchar, clients_portfolios_dm.client_id, 'client':Varchar, null:Int32, clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, 7:Int32] }
├── output: [ clients_portfolios_dm.portfolio_id, 'portfolio':Varchar, clients_portfolios_dm.client_id, 'client':Varchar, null:Int32, clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date, 7:Int32 ]
├── 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, clients_dm.id ], stream key: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_id, clients_portfolios_dm.effective_start_date ] }

Fragment 24847 (Actor 118530,118529)
StreamSyncLogStore { output: [ clients_portfolios_dm.client_id, clients_portfolios_dm.portfolio_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.client_id, clients_portfolios_dm.portfolio_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, 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_dm.id ], stream key: [ clients_dm.id ] }

Fragment 24848 (Actor 118534,118533)
StreamSyncLogStore { output: [ clients_portfolios_dm.client_id, 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, 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 ], stream key: [ portfolios_dm.portfolio_id ] }

Fragment 24849 (Actor 118535,118536)
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 24850 (Actor 118577,118578)
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 24851 (Actor 118543,118544)
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 24852 (Actor 118538,118537)
StreamProject { exprs: [account_to_portfolios_dm.portfolio_id, 'portfolio':Varchar, account_to_portfolios_dm.account_id, 'account':Varchar, null:Int32, account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date, 8:Int32] }
├── output: [ account_to_portfolios_dm.portfolio_id, 'portfolio':Varchar, account_to_portfolios_dm.account_id, 'account':Varchar, null:Int32, account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date, 8:Int32 ]
├── 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, 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 24853 (Actor 118539,118540)
StreamSyncLogStore { output: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_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.account_id, account_to_portfolios_dm.portfolio_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, portfolios_dm.portfolio_id ], 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 ], stream key: [ accounts_dm.account_id ] }

Fragment 24854 (Actor 118542,118541)
StreamSyncLogStore { output: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date, portfolios_dm.portfolio_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.portfolio_id = portfolios_dm.portfolio_id }
    ├── output: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date, portfolios_dm.portfolio_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: [ portfolios_dm.portfolio_id ], stream key: [ portfolios_dm.portfolio_id ] }

Fragment 24855 (Actor 118580,118579)
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) AND IsNull(account_to_portfolios_dm.effective_end_date) }
    ├── 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, account_to_portfolios_dm.effective_end_date ]
    ├── 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, effective_end_date] }
        ├── 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, account_to_portfolios_dm.effective_end_date ]
        ├── 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, effective_end_date ], stream key: [] }
        └── BatchPlanNode { output: [ account_id, portfolio_id, effective_start_date, disabled_at, effective_end_date ], stream key: [] }

Fragment 24856 (Actor 118546,118545)
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 24857 (Actor 118548,118547)
StreamProject { exprs: [accounts_dm.account_id] } { output: [ accounts_dm.account_id ], stream key: [ accounts_dm.account_id ] }
└── StreamFilter { predicate: IsNull(accounts_dm.disabled_at) } { output: [ accounts_dm.account_id, accounts_dm.disabled_at ], stream key: [ accounts_dm.account_id ] }
    └── StreamTableScan { table: accounts_dm, columns: [account_id, disabled_at] } { output: [ accounts_dm.account_id, accounts_dm.disabled_at ], stream key: [ accounts_dm.account_id ] }
        ├── Upstream { output: [ account_id, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ account_id, disabled_at ], stream key: [] }

Fragment 24858 (Actor 118550,118549)
StreamProject { exprs: [party_holder_edges_mv_next.party_id, 'party':Varchar, party_holder_edges_mv_next.entity_id, party_holder_edges_mv_next.entity_type, party_holder_edges_mv_next.$src, party_holder_edges_mv_next.party_id, party_holder_edges_mv_next.entity_id, null:Date, 9:Int32] }
├── output: [ party_holder_edges_mv_next.party_id, 'party':Varchar, party_holder_edges_mv_next.entity_id, party_holder_edges_mv_next.entity_type, party_holder_edges_mv_next.$src, party_holder_edges_mv_next.party_id, party_holder_edges_mv_next.entity_id, null:Date, 9:Int32 ]
├── stream key: [ party_holder_edges_mv_next.party_id, party_holder_edges_mv_next.entity_id, party_holder_edges_mv_next.$src ]
└── StreamTableScan { table: party_holder_edges_mv_next, columns: [party_id, entity_id, entity_type, $src] } { output: [ party_holder_edges_mv_next.party_id, party_holder_edges_mv_next.entity_id, party_holder_edges_mv_next.entity_type, party_holder_edges_mv_next.$src ], stream key: [ party_holder_edges_mv_next.party_id, party_holder_edges_mv_next.entity_id, party_holder_edges_mv_next.$src ] }
    ├── Upstream { output: [ party_id, entity_id, entity_type, $src ], stream key: [] }
    └── BatchPlanNode { output: [ party_id, entity_id, entity_type, $src ], stream key: [] }