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

← cluster alinma_bff objects party_current_account_membership_mv explain
Overview Objects Graph History
materialized view · alinma_bff.party_current_account_membership_mv profiled over 5s
seconds (1–30)
Stateful hash join (4 state tables) — consider a temporal join for dimension lookupsAggregation state — unbounded unless keyed or temporally filteredDynamic filter — verify it pairs with a temporal condition to clean state
205 operators
Materialize · alinma_bff.party_current_account_membership_mv
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · party_active_customer_relationships_mv.party_id = party_act…
2 actors
HashJoin · Inner · party_active_customer_relationships_mv.party_id = party_act… 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
LocalityProvider
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · LeftOuter · party_active_customer_relationships_mv.customer_relationshi…
2 actors
HashJoin · LeftOuter · party_active_customer_relationships_mv.customer_relationshi… 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
LocalityProvider
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · lifecycle_profiles
2 actors
Filter · lifecycle_profiles
0% idle 2 actors
StreamScan · lifecycle_profiles
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · party_active_customer_relationships_mv
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
LocalityProvider
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 · Not(IsTrue(open_accounts_mv.is_restricted))
2 actors
Filter · Not(IsTrue(open_accounts_mv.is_restricted))
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
2 actors
HashAgg Aggregation state — unbounded unless keyed or temporally filtered
0% idle 2 actors
LocalityProvider
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
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · party_involvements_dm.entity_id = account_to_portfolios_dm.…
2 actors
HashJoin · Inner · party_involvements_dm.entity_id = account_to_portfolios_dm.… 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
LocalityProvider
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · open_accounts_mv.account_id = account_to_portfolios_dm.acco…
2 actors
HashJoin · Inner · open_accounts_mv.account_id = account_to_portfolios_dm.acco… 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
LocalityProvider
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · account_to_portfolios_dm
2 actors
DynamicFilter · account_to_portfolios_dm Dynamic filter — verify it pairs with a temporal condition to clean state
0% idle 2 actors
Merge
2 actors
Exchange
0% 2/s 0 actors
Now
0% 2/s 1 actor
Project · account_to_portfolios_dm
2 actors
DynamicFilter · account_to_portfolios_dm Dynamic filter — verify it pairs with a temporal condition to clean state
3% idle 2 actors
Merge
2 actors
Exchange
0% 2/s 0 actors
Now
0% 2/s 1 actor
Project · account_to_portfolios_dm
2 actors
Filter · account_to_portfolios_dm
3% idle 2 actors
StreamScan · account_to_portfolios_dm
3% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · open_accounts_mv
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
LocalityProvider
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · party_active_customer_relationships_mv.customer_relationshi…
2 actors
HashJoin · Inner · party_active_customer_relationships_mv.customer_relationshi… 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
LocalityProvider
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · party_involvements_dm
2 actors
DynamicFilter · party_involvements_dm Dynamic filter — verify it pairs with a temporal condition to clean state
0% idle 2 actors
Merge
2 actors
Exchange
0% 2/s 0 actors
Now
0% 2/s 1 actor
Project · party_involvements_dm
2 actors
DynamicFilter · party_involvements_dm Dynamic filter — verify it pairs with a temporal condition to clean state
1% idle 2 actors
Merge
2 actors
Exchange
0% 2/s 0 actors
Now
0% 2/s 1 actor
Project · party_involvements_dm
2 actors
Filter · party_involvements_dm
4% idle 2 actors
StreamScan · party_involvements_dm
4% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
LocalityProvider
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · party_active_customer_relationships_mv
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
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · party_involvements_dm.entity_id = open_accounts_mv.account_…
2 actors
HashJoin · Inner · party_involvements_dm.entity_id = open_accounts_mv.account_… 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 · open_accounts_mv
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
LocalityProvider
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · party_active_customer_relationships_mv.customer_relationshi…
2 actors
HashJoin · Inner · party_active_customer_relationships_mv.customer_relationshi… 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
LocalityProvider
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · party_involvements_dm
2 actors
DynamicFilter · party_involvements_dm Dynamic filter — verify it pairs with a temporal condition to clean state
0% idle 2 actors
Merge
2 actors
Exchange
0% 2/s 0 actors
Now
0% 2/s 1 actor
Project · party_involvements_dm
2 actors
DynamicFilter · party_involvements_dm Dynamic filter — verify it pairs with a temporal condition to clean state
2% idle 2 actors
Merge
2 actors
Exchange
0% 2/s 0 actors
Now
0% 2/s 1 actor
Project · party_involvements_dm
2 actors
Filter · party_involvements_dm
5% idle 2 actors
StreamScan · party_involvements_dm
5% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
LocalityProvider
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · party_active_customer_relationships_mv
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · (open_accounts_mv.is_restricted = true:Boolean)
2 actors
Filter · (open_accounts_mv.is_restricted = true:Boolean)
0% idle 2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
2 actors
Merge
2 actors
Heat = the operator's output-buffer backpressure over the sampling window. Click a node to fold its subtree.
Materialize · alinma_bff.party_current_account_membership_mv Materialize alinma_bff.party_curren… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · party_active_customer_relationships_mv.party_id = party_act… SyncLogStore Inner · party_active_cu… — · 2 actors HashJoin · Inner · party_active_customer_relationships_mv.party_id = party_act… HashJoin Inner · party_active_cu… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors LocalityProvider LocalityProvider idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · LeftOuter · party_active_customer_relationships_mv.customer_relationshi… SyncLogStore LeftOuter · party_activ… — · 2 actors HashJoin · LeftOuter · party_active_customer_relationships_mv.customer_relationshi… HashJoin LeftOuter · party_activ… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors LocalityProvider LocalityProvider idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · lifecycle_profiles Project lifecycle_profiles — · 2 actors Filter · lifecycle_profiles Filter lifecycle_profiles idle · 2 actors StreamScan · lifecycle_profiles StreamScan lifecycle_profiles idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · party_active_customer_relationships_mv StreamScan party_active_customer_r… idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors LocalityProvider LocalityProvider 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 · Not(IsTrue(open_accounts_mv.is_restricted)) Project Not(IsTrue(open_account… — · 2 actors Filter · Not(IsTrue(open_accounts_mv.is_restricted)) Filter Not(IsTrue(open_account… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 2 actors HashAgg HashAgg idle · 2 actors LocalityProvider LocalityProvider 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 SyncLogStore — · 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 · party_involvements_dm.entity_id = account_to_portfolios_dm.… SyncLogStore Inner · party_involveme… — · 2 actors HashJoin · Inner · party_involvements_dm.entity_id = account_to_portfolios_dm.… HashJoin Inner · party_involveme… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors LocalityProvider LocalityProvider 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 SyncLogStore — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · open_accounts_mv.account_id = account_to_portfolios_dm.acco… SyncLogStore Inner · open_accounts_m… — · 2 actors HashJoin · Inner · open_accounts_mv.account_id = account_to_portfolios_dm.acco… HashJoin Inner · open_accounts_m… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors LocalityProvider LocalityProvider idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · account_to_portfolios_dm Project account_to_portfolios_dm — · 2 actors DynamicFilter · account_to_portfolios_dm DynamicFilter account_to_portfolios_dm idle · 2 actors Merge Merge — · 2 actors Exchange Exchange 2/s · 0 actors Now Now 2/s · 1 actor Project · account_to_portfolios_dm Project account_to_portfolios_dm — · 2 actors DynamicFilter · account_to_portfolios_dm DynamicFilter account_to_portfolios_dm idle · 2 actors Merge Merge — · 2 actors Exchange Exchange 2/s · 0 actors Now Now 2/s · 1 actor 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 StreamScan · open_accounts_mv StreamScan open_accounts_mv idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors LocalityProvider LocalityProvider 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 SyncLogStore — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · party_active_customer_relationships_mv.customer_relationshi… SyncLogStore Inner · party_active_cu… — · 2 actors HashJoin · Inner · party_active_customer_relationships_mv.customer_relationshi… HashJoin Inner · party_active_cu… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors LocalityProvider LocalityProvider idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · party_involvements_dm Project party_involvements_dm — · 2 actors DynamicFilter · party_involvements_dm DynamicFilter party_involvements_dm idle · 2 actors Merge Merge — · 2 actors Exchange Exchange 2/s · 0 actors Now Now 2/s · 1 actor Project · party_involvements_dm Project party_involvements_dm — · 2 actors DynamicFilter · party_involvements_dm DynamicFilter party_involvements_dm idle · 2 actors Merge Merge — · 2 actors Exchange Exchange 2/s · 0 actors Now Now 2/s · 1 actor Project · party_involvements_dm Project party_involvements_dm — · 2 actors Filter · party_involvements_dm Filter party_involvements_dm idle · 2 actors StreamScan · party_involvements_dm StreamScan party_involvements_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors LocalityProvider LocalityProvider idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · party_active_customer_relationships_mv StreamScan party_active_customer_r… 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 SyncLogStore — · 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 · party_involvements_dm.entity_id = open_accounts_mv.account_… SyncLogStore Inner · party_involveme… — · 2 actors HashJoin · Inner · party_involvements_dm.entity_id = open_accounts_mv.account_… HashJoin Inner · party_involveme… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · open_accounts_mv StreamScan open_accounts_mv idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors LocalityProvider LocalityProvider 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 SyncLogStore — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · party_active_customer_relationships_mv.customer_relationshi… SyncLogStore Inner · party_active_cu… — · 2 actors HashJoin · Inner · party_active_customer_relationships_mv.customer_relationshi… HashJoin Inner · party_active_cu… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors LocalityProvider LocalityProvider idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · party_involvements_dm Project party_involvements_dm — · 2 actors DynamicFilter · party_involvements_dm DynamicFilter party_involvements_dm idle · 2 actors Merge Merge — · 2 actors Exchange Exchange 2/s · 0 actors Now Now 2/s · 1 actor Project · party_involvements_dm Project party_involvements_dm — · 2 actors DynamicFilter · party_involvements_dm DynamicFilter party_involvements_dm idle · 2 actors Merge Merge — · 2 actors Exchange Exchange 2/s · 0 actors Now Now 2/s · 1 actor Project · party_involvements_dm Project party_involvements_dm — · 2 actors Filter · party_involvements_dm Filter party_involvements_dm idle · 2 actors StreamScan · party_involvements_dm StreamScan party_involvements_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors LocalityProvider LocalityProvider idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · party_active_customer_relationships_mv StreamScan party_active_customer_r… idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · (open_accounts_mv.is_restricted = true:Boolean) Project (open_accounts_mv.is_re… — · 2 actors Filter · (open_accounts_mv.is_restricted = true:Boolean) Filter (open_accounts_mv.is_re… idle · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 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 52177 (Actor 737109,737110)
StreamMaterialize { columns: [party_id, customer_relationship_id, account_id, account_group_type, party_currency, party_active_customer_relationships_mv.party_id(hidden), party_active_customer_relationships_mv.customer_relationship_id(hidden), party_involvements_dm.entity_id(hidden), open_accounts_mv.is_restricted(hidden), $src(hidden), party_active_customer_relationships_mv.party_id#1(hidden), party_active_customer_relationships_mv.customer_relationship_id#1(hidden), lifecycle_profiles.id(hidden)], stream_key: [party_id, customer_relationship_id, party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted, $src, lifecycle_profiles.id], pk_columns: [party_id, customer_relationship_id, party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted, $src, lifecycle_profiles.id], pk_conflict: NoCheck }
├── output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, 'all':Varchar, lifecycle_profiles.base_currency_code, party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted, $src, party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, lifecycle_profiles.id ]
├── stream key: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted, $src, lifecycle_profiles.id ]
└── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, 'all':Varchar, lifecycle_profiles.base_currency_code, party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted, $src, party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, lifecycle_profiles.id ], stream key: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted, $src, lifecycle_profiles.id ] }

Fragment 52178 (Actor 737112,737111)
StreamSyncLogStore { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, 'all':Varchar, lifecycle_profiles.base_currency_code, party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted, $src, party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, lifecycle_profiles.id ], stream key: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted, $src, lifecycle_profiles.id ] }
└── StreamHashJoin { type: Inner, predicate: party_active_customer_relationships_mv.party_id = party_active_customer_relationships_mv.party_id AND party_active_customer_relationships_mv.customer_relationship_id = party_active_customer_relationships_mv.customer_relationship_id }
    ├── output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, 'all':Varchar, lifecycle_profiles.base_currency_code, party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted, $src, party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, lifecycle_profiles.id ]
    ├── stream key: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted, $src, lifecycle_profiles.id ]
    ├── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, 'all':Varchar, open_accounts_mv.is_restricted, party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, $src ], stream key: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted, $src ] }
    └── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, lifecycle_profiles.base_currency_code, lifecycle_profiles.customer_relationship_id, lifecycle_profiles.id ], stream key: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, lifecycle_profiles.id ] }

Fragment 52179 (Actor 737113,737114)
StreamLocalityProvider { locality_columns: [party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id] } { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, 'all':Varchar, open_accounts_mv.is_restricted, party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, $src ], stream key: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted, $src ] }
└── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, 'all':Varchar, open_accounts_mv.is_restricted, party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, $src ], stream key: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted, $src ] }

Fragment 52180 (Actor 737365,737364)
StreamUnion { all: true } { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, 'all':Varchar, open_accounts_mv.is_restricted, party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, $src ], stream key: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted, $src ] }
├── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, 'all':Varchar, open_accounts_mv.is_restricted, party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, 0:Int32 ], stream key: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted ] }
├── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, 'restricted':Varchar, open_accounts_mv.is_restricted, party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, 1:Int32 ], stream key: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted ] }
└── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, 'un_restricted':Varchar, open_accounts_mv.is_restricted, party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, 2:Int32 ], stream key: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted ] }

Fragment 52181 (Actor 737370,737371)
StreamProject { exprs: [party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, 'all':Varchar, open_accounts_mv.is_restricted, party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, 0:Int32] } { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, 'all':Varchar, open_accounts_mv.is_restricted, party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, 0:Int32 ], stream key: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted ] }
└── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted ], stream key: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted ] }

Fragment 52182 (Actor 737367,737366)
StreamProject { exprs: [party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted] } { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted ], stream key: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted ] }
└── StreamHashAgg { group_key: [party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted], aggs: [count] } { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted, count ], stream key: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted ] }
    └── StreamLocalityProvider { locality_columns: [party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted] } { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted, party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id, null:Varchar, null:Date, $src ], stream key: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted, party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id, null:Varchar, null:Date, $src ] }
        └── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted, party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id, null:Varchar, null:Date, $src ], stream key: [ party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id, null:Varchar, null:Date, $src ] }

Fragment 52183 (Actor 737376,737377)
StreamUnion { all: true } { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted, party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id, null:Varchar, null:Date, $src ], stream key: [ party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id, null:Varchar, null:Date, $src ] }
├── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted, party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id, null:Varchar, null:Date, 0:Int32 ], stream key: [ party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id ] }
└── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, account_to_portfolios_dm.account_id, open_accounts_mv.is_restricted, party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date, 1:Int32 ], stream key: [ party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date ] }

Fragment 52184 (Actor 737387,737386)
StreamProject { exprs: [party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted, party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id, null:Varchar, null:Date, 0:Int32] } { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted, party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id, null:Varchar, null:Date, 0:Int32 ], stream key: [ party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id ] }
└── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted, party_involvements_dm.id, open_accounts_mv.account_id ], stream key: [ party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id ] }

Fragment 52185 (Actor 737385,737384)
StreamSyncLogStore { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted, party_involvements_dm.id, open_accounts_mv.account_id ], stream key: [ party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id ] }
└── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted, party_involvements_dm.id, open_accounts_mv.account_id ], stream key: [ party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id ] }

Fragment 52186 (Actor 737382,737383)
StreamSyncLogStore { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted, party_involvements_dm.id, open_accounts_mv.account_id ], stream key: [ party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id ] }
└── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted, party_involvements_dm.id, open_accounts_mv.account_id ], stream key: [ party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id ] }

Fragment 52187 (Actor 737381,737380)
StreamSyncLogStore { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted, party_involvements_dm.id, open_accounts_mv.account_id ], stream key: [ party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id ] }
└── StreamHashJoin { type: Inner, predicate: party_involvements_dm.entity_id = open_accounts_mv.account_id } { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted, party_involvements_dm.id, open_accounts_mv.account_id ], stream key: [ party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id ] }
    ├── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.party_id, party_involvements_dm.id ], stream key: [ party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id ] }
    └── MergeExecutor { output: [ open_accounts_mv.account_id, open_accounts_mv.is_restricted ], stream key: [ open_accounts_mv.account_id ] }

Fragment 52188 (Actor 737388,737389)
StreamLocalityProvider { locality_columns: [party_involvements_dm.entity_id] } { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.party_id, party_involvements_dm.id ], stream key: [ party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id ] }
└── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.party_id, party_involvements_dm.id ], stream key: [ party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id ] }

Fragment 52189 (Actor 737391,737390)
StreamSyncLogStore { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.party_id, party_involvements_dm.id ], stream key: [ party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id ] }
└── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.party_id, party_involvements_dm.id ], stream key: [ party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id ] }

Fragment 52190 (Actor 737394,737395)
StreamSyncLogStore { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.party_id, party_involvements_dm.id ], stream key: [ party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id ] }
└── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.party_id, party_involvements_dm.id ], stream key: [ party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id ] }

Fragment 52191 (Actor 737393,737392)
StreamSyncLogStore { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.party_id, party_involvements_dm.id ], stream key: [ party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id ] }
└── StreamHashJoin { type: Inner, predicate: party_active_customer_relationships_mv.customer_relationship_id = party_involvements_dm.customer_relationship_id AND party_active_customer_relationships_mv.party_id = party_involvements_dm.party_id } { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.party_id, party_involvements_dm.id ], stream key: [ party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id ] }
    ├── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id ], stream key: [ party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id ] }
    └── MergeExecutor { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.id ], stream key: [ party_involvements_dm.customer_relationship_id, party_involvements_dm.party_id, party_involvements_dm.id ] }

Fragment 52192 (Actor 737396,737397)
StreamLocalityProvider { locality_columns: [party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id] } { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id ], stream key: [ party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id ] }
└── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id ], stream key: [ party_active_customer_relationships_mv.customer_relationship_id ] }

Fragment 52193 (Actor 737351,737350)
StreamTableScan { table: party_active_customer_relationships_mv, columns: [party_id, customer_relationship_id] } { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id ], stream key: [ party_active_customer_relationships_mv.customer_relationship_id ] }
├── Upstream { output: [ party_id, customer_relationship_id ], stream key: [] }
└── BatchPlanNode { output: [ party_id, customer_relationship_id ], stream key: [] }

Fragment 52194 (Actor 737398,737399)
StreamLocalityProvider { locality_columns: [party_involvements_dm.customer_relationship_id, party_involvements_dm.party_id] } { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.id ], stream key: [ party_involvements_dm.customer_relationship_id, party_involvements_dm.party_id, party_involvements_dm.id ] }
└── MergeExecutor { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.id ], stream key: [ party_involvements_dm.id ] }

Fragment 52195 (Actor 737402,737403)
StreamProject { exprs: [party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.id] } { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.id ], stream key: [ party_involvements_dm.id ] }
└── StreamDynamicFilter { predicate: ($expr2 > now), output_watermarks: [[$expr2]], output: [party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, $expr2, party_involvements_dm.id], cleaned_by_watermark: true } { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, $expr2, party_involvements_dm.id ], stream key: [ party_involvements_dm.id ] }
    ├── StreamProject { exprs: [party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, AtTimeZone(Coalesce(party_involvements_dm.effective_to, '9999-12-31':Date)::Timestamp, 'UTC':Varchar) as $expr2, party_involvements_dm.id] } { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, $expr2, party_involvements_dm.id ], stream key: [ party_involvements_dm.id ] }
    │   └── StreamDynamicFilter { predicate: ($expr1 <= now), output: [party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.effective_to, $expr1, party_involvements_dm.id], cleaned_by_watermark: true } { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.effective_to, $expr1, party_involvements_dm.id ], stream key: [ party_involvements_dm.id ] }
    │       ├── StreamProject { exprs: [party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.effective_to, AtTimeZone(party_involvements_dm.effective_from::Timestamp, 'UTC':Varchar) as $expr1, party_involvements_dm.id] } { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.effective_to, $expr1, party_involvements_dm.id ], stream key: [ party_involvements_dm.id ] }
    │       │   └── StreamFilter { predicate: (party_involvements_dm.entity_type = 'ACCOUNT':Varchar) AND In(party_involvements_dm.involvement_type, 'ACCOUNT_HOLDER':Varchar, 'JOINT_ACCOUNT_HOLDER':Varchar) AND (party_involvements_dm.status = 'ACTIVE':Varchar) AND IsNull(party_involvements_dm.disabled_at) } { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.effective_from, party_involvements_dm.effective_to, party_involvements_dm.id, party_involvements_dm.involvement_type, party_involvements_dm.entity_type, party_involvements_dm.status, party_involvements_dm.disabled_at ], stream key: [ party_involvements_dm.id ] }
    │       │       └── StreamTableScan { table: party_involvements_dm, columns: [party_id, customer_relationship_id, entity_id, effective_from, effective_to, id, involvement_type, entity_type, status, disabled_at] } { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.effective_from, party_involvements_dm.effective_to, party_involvements_dm.id, party_involvements_dm.involvement_type, party_involvements_dm.entity_type, party_involvements_dm.status, party_involvements_dm.disabled_at ], stream key: [ party_involvements_dm.id ] }
    │       │           ├── Upstream { output: [ party_id, customer_relationship_id, entity_id, effective_from, effective_to, id, involvement_type, entity_type, status, disabled_at ], stream key: [] }
    │       │           └── BatchPlanNode { output: [ party_id, customer_relationship_id, entity_id, effective_from, effective_to, id, involvement_type, entity_type, status, disabled_at ], stream key: [] }
    │       └── MergeExecutor { output: [ now ], stream key: [] }
    └── MergeExecutor { output: [ now ], stream key: [] }

Fragment 52196 (Actor 737400)
StreamNow { output: [ now ], stream key: [] }

Fragment 52197 (Actor 737401)
StreamNow { output: [ now ], stream key: [] }

Fragment 52198 (Actor 737404,737405)
StreamTableScan { table: open_accounts_mv, columns: [account_id, is_restricted] } { output: [ open_accounts_mv.account_id, open_accounts_mv.is_restricted ], stream key: [ open_accounts_mv.account_id ] }
├── Upstream { output: [ account_id, is_restricted ], stream key: [] }
└── BatchPlanNode { output: [ account_id, is_restricted ], stream key: [] }

Fragment 52199 (Actor 737411,737410)
StreamProject { exprs: [party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, account_to_portfolios_dm.account_id, open_accounts_mv.is_restricted, party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date, 1:Int32] }
├── output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, account_to_portfolios_dm.account_id, open_accounts_mv.is_restricted, party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date, 1:Int32 ]
├── stream key: [ party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date ]
└── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, account_to_portfolios_dm.account_id, open_accounts_mv.is_restricted, party_involvements_dm.entity_id, party_involvements_dm.id, account_to_portfolios_dm.portfolio_id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date ], stream key: [ party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date ] }

Fragment 52200 (Actor 737407,737406)
StreamSyncLogStore { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, account_to_portfolios_dm.account_id, open_accounts_mv.is_restricted, party_involvements_dm.entity_id, party_involvements_dm.id, account_to_portfolios_dm.portfolio_id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date ], stream key: [ party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date ] }
└── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, account_to_portfolios_dm.account_id, open_accounts_mv.is_restricted, party_involvements_dm.entity_id, party_involvements_dm.id, account_to_portfolios_dm.portfolio_id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date ], stream key: [ party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date ] }

Fragment 52201 (Actor 737409,737408)
StreamSyncLogStore { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, account_to_portfolios_dm.account_id, open_accounts_mv.is_restricted, party_involvements_dm.entity_id, party_involvements_dm.id, account_to_portfolios_dm.portfolio_id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date ], stream key: [ party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date ] }
└── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, account_to_portfolios_dm.account_id, open_accounts_mv.is_restricted, party_involvements_dm.entity_id, party_involvements_dm.id, account_to_portfolios_dm.portfolio_id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date ], stream key: [ party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date ] }

Fragment 52202 (Actor 737412,737413)
StreamSyncLogStore { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, account_to_portfolios_dm.account_id, open_accounts_mv.is_restricted, party_involvements_dm.entity_id, party_involvements_dm.id, account_to_portfolios_dm.portfolio_id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date ], stream key: [ party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date ] }
└── StreamHashJoin { type: Inner, predicate: party_involvements_dm.entity_id = account_to_portfolios_dm.portfolio_id } { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, account_to_portfolios_dm.account_id, open_accounts_mv.is_restricted, party_involvements_dm.entity_id, party_involvements_dm.id, account_to_portfolios_dm.portfolio_id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date ], stream key: [ party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date ] }
    ├── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.party_id, party_involvements_dm.id ], stream key: [ party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id ] }
    └── MergeExecutor { output: [ open_accounts_mv.is_restricted, account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date ], stream key: [ account_to_portfolios_dm.portfolio_id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date ] }

Fragment 52203 (Actor 737415,737414)
StreamLocalityProvider { locality_columns: [party_involvements_dm.entity_id] } { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.party_id, party_involvements_dm.id ], stream key: [ party_involvements_dm.entity_id, party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id ] }
└── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.party_id, party_involvements_dm.id ], stream key: [ party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id ] }

Fragment 52204 (Actor 737421,737420)
StreamSyncLogStore { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.party_id, party_involvements_dm.id ], stream key: [ party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id ] }
└── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.party_id, party_involvements_dm.id ], stream key: [ party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id ] }

Fragment 52205 (Actor 737419,737418)
StreamSyncLogStore { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.party_id, party_involvements_dm.id ], stream key: [ party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id ] }
└── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.party_id, party_involvements_dm.id ], stream key: [ party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id ] }

Fragment 52206 (Actor 737416,737417)
StreamSyncLogStore { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.party_id, party_involvements_dm.id ], stream key: [ party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id ] }
└── StreamHashJoin { type: Inner, predicate: party_active_customer_relationships_mv.customer_relationship_id = party_involvements_dm.customer_relationship_id AND party_active_customer_relationships_mv.party_id = party_involvements_dm.party_id } { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.party_id, party_involvements_dm.id ], stream key: [ party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id, party_involvements_dm.id ] }
    ├── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id ], stream key: [ party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id ] }
    └── MergeExecutor { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.id ], stream key: [ party_involvements_dm.customer_relationship_id, party_involvements_dm.party_id, party_involvements_dm.id ] }

Fragment 52207 (Actor 737422,737423)
StreamLocalityProvider { locality_columns: [party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id] } { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id ], stream key: [ party_active_customer_relationships_mv.customer_relationship_id, party_active_customer_relationships_mv.party_id ] }
└── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id ], stream key: [ party_active_customer_relationships_mv.customer_relationship_id ] }

Fragment 52208 (Actor 737452,737453)
StreamTableScan { table: party_active_customer_relationships_mv, columns: [party_id, customer_relationship_id] } { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id ], stream key: [ party_active_customer_relationships_mv.customer_relationship_id ] }
├── Upstream { output: [ party_id, customer_relationship_id ], stream key: [] }
└── BatchPlanNode { output: [ party_id, customer_relationship_id ], stream key: [] }

Fragment 52209 (Actor 737424,737425)
StreamLocalityProvider { locality_columns: [party_involvements_dm.customer_relationship_id, party_involvements_dm.party_id] } { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.id ], stream key: [ party_involvements_dm.customer_relationship_id, party_involvements_dm.party_id, party_involvements_dm.id ] }
└── MergeExecutor { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.id ], stream key: [ party_involvements_dm.id ] }

Fragment 52210 (Actor 737455,737454)
StreamProject { exprs: [party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.id] } { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.id ], stream key: [ party_involvements_dm.id ] }
└── StreamDynamicFilter { predicate: ($expr4 > now), output_watermarks: [[$expr4]], output: [party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, $expr4, party_involvements_dm.id], cleaned_by_watermark: true } { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, $expr4, party_involvements_dm.id ], stream key: [ party_involvements_dm.id ] }
    ├── StreamProject { exprs: [party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, AtTimeZone(Coalesce(party_involvements_dm.effective_to, '9999-12-31':Date)::Timestamp, 'UTC':Varchar) as $expr4, party_involvements_dm.id] } { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, $expr4, party_involvements_dm.id ], stream key: [ party_involvements_dm.id ] }
    │   └── StreamDynamicFilter { predicate: ($expr3 <= now), output: [party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.effective_to, $expr3, party_involvements_dm.id], cleaned_by_watermark: true } { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.effective_to, $expr3, party_involvements_dm.id ], stream key: [ party_involvements_dm.id ] }
    │       ├── StreamProject { exprs: [party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.effective_to, AtTimeZone(party_involvements_dm.effective_from::Timestamp, 'UTC':Varchar) as $expr3, party_involvements_dm.id] } { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.effective_to, $expr3, party_involvements_dm.id ], stream key: [ party_involvements_dm.id ] }
    │       │   └── StreamFilter { predicate: (party_involvements_dm.entity_type = 'PORTFOLIO':Varchar) AND In(party_involvements_dm.involvement_type, 'PORTFOLIO_HOLDER':Varchar, 'JOINT_PORTFOLIO_HOLDER':Varchar) AND (party_involvements_dm.status = 'ACTIVE':Varchar) AND IsNull(party_involvements_dm.disabled_at) } { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.effective_from, party_involvements_dm.effective_to, party_involvements_dm.id, party_involvements_dm.involvement_type, party_involvements_dm.entity_type, party_involvements_dm.status, party_involvements_dm.disabled_at ], stream key: [ party_involvements_dm.id ] }
    │       │       └── StreamTableScan { table: party_involvements_dm, columns: [party_id, customer_relationship_id, entity_id, effective_from, effective_to, id, involvement_type, entity_type, status, disabled_at] } { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.effective_from, party_involvements_dm.effective_to, party_involvements_dm.id, party_involvements_dm.involvement_type, party_involvements_dm.entity_type, party_involvements_dm.status, party_involvements_dm.disabled_at ], stream key: [ party_involvements_dm.id ] }
    │       │           ├── Upstream { output: [ party_id, customer_relationship_id, entity_id, effective_from, effective_to, id, involvement_type, entity_type, status, disabled_at ], stream key: [] }
    │       │           └── BatchPlanNode { output: [ party_id, customer_relationship_id, entity_id, effective_from, effective_to, id, involvement_type, entity_type, status, disabled_at ], stream key: [] }
    │       └── MergeExecutor { output: [ now ], stream key: [] }
    └── MergeExecutor { output: [ now ], stream key: [] }

Fragment 52211 (Actor 737426)
StreamNow { output: [ now ], stream key: [] }

Fragment 52212 (Actor 737427)
StreamNow { output: [ now ], stream key: [] }

Fragment 52213 (Actor 737429,737428)
StreamLocalityProvider { locality_columns: [account_to_portfolios_dm.portfolio_id] } { output: [ open_accounts_mv.is_restricted, account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date ], stream key: [ account_to_portfolios_dm.portfolio_id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date ] }
└── MergeExecutor { output: [ open_accounts_mv.is_restricted, account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date ], stream key: [ open_accounts_mv.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ] }

Fragment 52214 (Actor 737435,737434)
StreamSyncLogStore { output: [ open_accounts_mv.is_restricted, account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date ], stream key: [ open_accounts_mv.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ] }
└── MergeExecutor { output: [ open_accounts_mv.is_restricted, account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date ], stream key: [ open_accounts_mv.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ] }

Fragment 52215 (Actor 737433,737432)
StreamSyncLogStore { output: [ open_accounts_mv.is_restricted, account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date ], stream key: [ open_accounts_mv.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ] }
└── MergeExecutor { output: [ open_accounts_mv.is_restricted, account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date ], stream key: [ open_accounts_mv.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ] }

Fragment 52216 (Actor 737431,737430)
StreamSyncLogStore { output: [ open_accounts_mv.is_restricted, account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date ], stream key: [ open_accounts_mv.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ] }
└── StreamHashJoin { type: Inner, predicate: open_accounts_mv.account_id = account_to_portfolios_dm.account_id } { output: [ open_accounts_mv.is_restricted, account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date ], stream key: [ open_accounts_mv.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ] }
    ├── MergeExecutor { output: [ open_accounts_mv.account_id, open_accounts_mv.is_restricted ], stream key: [ open_accounts_mv.account_id ] }
    └── 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 ] }

Fragment 52217 (Actor 737457,737456)
StreamTableScan { table: open_accounts_mv, columns: [account_id, is_restricted] } { output: [ open_accounts_mv.account_id, open_accounts_mv.is_restricted ], stream key: [ open_accounts_mv.account_id ] }
├── Upstream { output: [ account_id, is_restricted ], stream key: [] }
└── BatchPlanNode { output: [ account_id, is_restricted ], stream key: [] }

Fragment 52218 (Actor 737436,737437)
StreamLocalityProvider { locality_columns: [account_to_portfolios_dm.account_id] } { 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: [ 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 ] }

Fragment 52219 (Actor 737458,737459)
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 ] }
└── StreamDynamicFilter { predicate: ($expr6 > now), output_watermarks: [[$expr6]], output: [account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, $expr6, account_to_portfolios_dm.effective_start_date], cleaned_by_watermark: true } { output: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, $expr6, 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 ] }
    ├── StreamProject { exprs: [account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, AtTimeZone(Coalesce(account_to_portfolios_dm.effective_end_date, '9999-12-31':Date)::Timestamp, 'UTC':Varchar) as $expr6, account_to_portfolios_dm.effective_start_date] } { output: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, $expr6, 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 ] }
    │   └── StreamDynamicFilter { predicate: ($expr5 <= now), output: [account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_end_date, $expr5, account_to_portfolios_dm.effective_start_date], cleaned_by_watermark: true } { output: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_end_date, $expr5, 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 ] }
    │       ├── StreamProject { exprs: [account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_end_date, AtTimeZone(account_to_portfolios_dm.effective_start_date::Timestamp, 'UTC':Varchar) as $expr5, 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_end_date, $expr5, 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.effective_end_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, effective_end_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.effective_end_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, effective_end_date, disabled_at ], stream key: [] }
    │       │           └── BatchPlanNode { output: [ account_id, portfolio_id, effective_start_date, effective_end_date, disabled_at ], stream key: [] }
    │       └── MergeExecutor { output: [ now ], stream key: [] }
    └── MergeExecutor { output: [ now ], stream key: [] }

Fragment 52220 (Actor 737440)
StreamNow { output: [ now ], stream key: [] }

Fragment 52221 (Actor 737441)
StreamNow { output: [ now ], stream key: [] }

Fragment 52222 (Actor 737372,737373)
StreamProject { exprs: [party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, 'restricted':Varchar, open_accounts_mv.is_restricted, party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, 1:Int32] } { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, 'restricted':Varchar, open_accounts_mv.is_restricted, party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, 1:Int32 ], stream key: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted ] }
└── StreamFilter { predicate: (open_accounts_mv.is_restricted = true:Boolean) } { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted ], stream key: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted ] }
    └── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted ], stream key: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted ] }

Fragment 52223 (Actor 737369,737368)
StreamProject { exprs: [party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, 'un_restricted':Varchar, open_accounts_mv.is_restricted, party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, 2:Int32] } { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, 'un_restricted':Varchar, open_accounts_mv.is_restricted, party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, 2:Int32 ], stream key: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted ] }
└── StreamFilter { predicate: Not(IsTrue(open_accounts_mv.is_restricted)) } { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted ], stream key: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted ] }
    └── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted ], stream key: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, party_involvements_dm.entity_id, open_accounts_mv.is_restricted ] }

Fragment 52224 (Actor 737444,737445)
StreamLocalityProvider { locality_columns: [party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id] } { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, lifecycle_profiles.base_currency_code, lifecycle_profiles.customer_relationship_id, lifecycle_profiles.id ], stream key: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, lifecycle_profiles.id ] }
└── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, lifecycle_profiles.base_currency_code, lifecycle_profiles.customer_relationship_id, lifecycle_profiles.id ], stream key: [ party_active_customer_relationships_mv.customer_relationship_id, lifecycle_profiles.id ] }

Fragment 52225 (Actor 737447,737446)
StreamSyncLogStore { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, lifecycle_profiles.base_currency_code, lifecycle_profiles.customer_relationship_id, lifecycle_profiles.id ], stream key: [ party_active_customer_relationships_mv.customer_relationship_id, lifecycle_profiles.id ] }
└── StreamHashJoin { type: LeftOuter, predicate: party_active_customer_relationships_mv.customer_relationship_id = lifecycle_profiles.customer_relationship_id } { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id, lifecycle_profiles.base_currency_code, lifecycle_profiles.customer_relationship_id, lifecycle_profiles.id ], stream key: [ party_active_customer_relationships_mv.customer_relationship_id, lifecycle_profiles.id ] }
    ├── MergeExecutor { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id ], stream key: [ party_active_customer_relationships_mv.customer_relationship_id ] }
    └── MergeExecutor { output: [ lifecycle_profiles.customer_relationship_id, lifecycle_profiles.base_currency_code, lifecycle_profiles.id ], stream key: [ lifecycle_profiles.customer_relationship_id, lifecycle_profiles.id ] }

Fragment 52226 (Actor 737460,737461)
StreamTableScan { table: party_active_customer_relationships_mv, columns: [party_id, customer_relationship_id] } { output: [ party_active_customer_relationships_mv.party_id, party_active_customer_relationships_mv.customer_relationship_id ], stream key: [ party_active_customer_relationships_mv.customer_relationship_id ] }
├── Upstream { output: [ party_id, customer_relationship_id ], stream key: [] }
└── BatchPlanNode { output: [ party_id, customer_relationship_id ], stream key: [] }

Fragment 52227 (Actor 737450,737451)
StreamLocalityProvider { locality_columns: [lifecycle_profiles.customer_relationship_id] } { output: [ lifecycle_profiles.customer_relationship_id, lifecycle_profiles.base_currency_code, lifecycle_profiles.id ], stream key: [ lifecycle_profiles.customer_relationship_id, lifecycle_profiles.id ] }
└── MergeExecutor { output: [ lifecycle_profiles.customer_relationship_id, lifecycle_profiles.base_currency_code, lifecycle_profiles.id ], stream key: [ lifecycle_profiles.id ] }

Fragment 52228 (Actor 737462,737463)
StreamProject { exprs: [lifecycle_profiles.customer_relationship_id, lifecycle_profiles.base_currency_code, lifecycle_profiles.id] } { output: [ lifecycle_profiles.customer_relationship_id, lifecycle_profiles.base_currency_code, lifecycle_profiles.id ], stream key: [ lifecycle_profiles.id ] }
└── StreamFilter { predicate: IsNull(lifecycle_profiles.disabled_at) } { output: [ lifecycle_profiles.customer_relationship_id, lifecycle_profiles.base_currency_code, lifecycle_profiles.id, lifecycle_profiles.disabled_at ], stream key: [ lifecycle_profiles.id ] }
    └── StreamTableScan { table: lifecycle_profiles, columns: [customer_relationship_id, base_currency_code, id, disabled_at] } { output: [ lifecycle_profiles.customer_relationship_id, lifecycle_profiles.base_currency_code, lifecycle_profiles.id, lifecycle_profiles.disabled_at ], stream key: [ lifecycle_profiles.id ] }
        ├── Upstream { output: [ customer_relationship_id, base_currency_code, id, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ customer_relationship_id, base_currency_code, id, disabled_at ], stream key: [] }