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

← cluster insights objects account_to_account_groups_mv explain
Overview Objects Graph History
materialized view · insights.account_to_account_groups_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 lookupsWindow state — add a WHERE rank <= N to bound it
92 operators
Materialize · insights.account_to_account_groups_mv
0% idle 2 actors
Project
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · $expr1 = account_groups_mv_next.account_group_id
2 actors
HashJoin · Inner · $expr1 = account_groups_mv_next.account_group_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
LocalityProvider
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · account_groups_mv_next
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
OverWindow Window state — add a WHERE rank <= N to bound it
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 · party_account_via_portfolio_mv_next
2 actors
StreamScan · party_account_via_portfolio_mv_next
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · party_account_direct_mv_next
2 actors
StreamScan · party_account_direct_mv_next
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · IsNull(account_to_portfolios_dm.disabled_at) AND (IsNull(ac…
2 actors
Filter · IsNull(account_to_portfolios_dm.disabled_at) AND (IsNull(ac…
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 · account_to_portfolios_dm.account_id = open_accounts_mv.acco…
2 actors
HashJoin · Inner · account_to_portfolios_dm.account_id = open_accounts_mv.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
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
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 · IsNull(account_to_portfolios_dm.disabled_at) AND (IsNull(ac…
2 actors
Filter · IsNull(account_to_portfolios_dm.disabled_at) AND (IsNull(ac…
0% idle 2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · IsNull(account_to_portfolios_dm.disabled_at) AND (IsNull(ac…
2 actors
Filter · IsNull(account_to_portfolios_dm.disabled_at) AND (IsNull(ac…
0% idle 2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · client_account_via_portfolio_mv
2 actors
StreamScan · client_account_via_portfolio_mv
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · client_account_direct_mv
2 actors
StreamScan · client_account_direct_mv
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · open_accounts_mv
2 actors
StreamScan · open_accounts_mv
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 · insights.account_to_account_groups_mv Materialize insights.account_to_acc… idle · 2 actors Project Project — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · $expr1 = account_groups_mv_next.account_group_id SyncLogStore Inner · $expr1 = accoun… — · 2 actors HashJoin · Inner · $expr1 = account_groups_mv_next.account_group_id HashJoin Inner · $expr1 = accoun… 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 StreamScan · account_groups_mv_next StreamScan account_groups_mv_next 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 OverWindow OverWindow 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 · party_account_via_portfolio_mv_next Project party_account_via_portf… — · 2 actors StreamScan · party_account_via_portfolio_mv_next StreamScan party_account_via_portf… idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · party_account_direct_mv_next Project party_account_direct_mv… — · 2 actors StreamScan · party_account_direct_mv_next StreamScan party_account_direct_mv… idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · IsNull(account_to_portfolios_dm.disabled_at) AND (IsNull(ac… Project IsNull(account_to_portf… — · 2 actors Filter · IsNull(account_to_portfolios_dm.disabled_at) AND (IsNull(ac… Filter IsNull(account_to_portf… 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 · account_to_portfolios_dm.account_id = open_accounts_mv.acco… SyncLogStore Inner · account_to_port… — · 2 actors HashJoin · Inner · account_to_portfolios_dm.account_id = open_accounts_mv.acco… HashJoin Inner · account_to_port… 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 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 · IsNull(account_to_portfolios_dm.disabled_at) AND (IsNull(ac… Project IsNull(account_to_portf… — · 2 actors Filter · IsNull(account_to_portfolios_dm.disabled_at) AND (IsNull(ac… Filter IsNull(account_to_portf… idle · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · IsNull(account_to_portfolios_dm.disabled_at) AND (IsNull(ac… Project IsNull(account_to_portf… — · 2 actors Filter · IsNull(account_to_portfolios_dm.disabled_at) AND (IsNull(ac… Filter IsNull(account_to_portf… idle · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · client_account_via_portfolio_mv Project client_account_via_port… — · 2 actors StreamScan · client_account_via_portfolio_mv StreamScan client_account_via_port… idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · client_account_direct_mv Project client_account_direct_mv — · 2 actors StreamScan · client_account_direct_mv StreamScan client_account_direct_mv idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · open_accounts_mv Project open_accounts_mv — · 2 actors StreamScan · open_accounts_mv StreamScan open_accounts_mv 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 61164 (Actor 738466,738467)
StreamMaterialize { columns: [account_id, account_group_id, effective_start_date, effective_end_date, base_currency, opening_date, source_entity_type, open_accounts_mv.account_id(hidden), null:Varchar(hidden), null:Date(hidden), null:Int32(hidden), null:Varchar#1(hidden), null:Date#1(hidden), null:Varchar#2(hidden), null:Varchar#3(hidden), null:Varchar#4(hidden), $src(hidden), account_groups_mv_next.open_accounts_mv.account_id(hidden), account_groups_mv_next.null:Varchar(hidden), account_groups_mv_next.null:Int32(hidden), account_groups_mv_next.null:Varchar#1(hidden), account_groups_mv_next.$src(hidden)], stream_key: [account_group_id, account_id, open_accounts_mv.account_id, null:Varchar, null:Date, null:Int32, null:Varchar#1, null:Date#1, null:Varchar#2, null:Varchar#3, null:Varchar#4, $src, account_groups_mv_next.open_accounts_mv.account_id, account_groups_mv_next.null:Varchar, account_groups_mv_next.null:Int32, account_groups_mv_next.null:Varchar#1, account_groups_mv_next.$src], pk_columns: [account_group_id, account_id, open_accounts_mv.account_id, null:Varchar, null:Date, null:Int32, null:Varchar#1, null:Date#1, null:Varchar#2, null:Varchar#3, null:Varchar#4, $src, account_groups_mv_next.open_accounts_mv.account_id, account_groups_mv_next.null:Varchar, account_groups_mv_next.null:Int32, account_groups_mv_next.null:Varchar#1, account_groups_mv_next.$src], pk_conflict: NoCheck }
├── output: [ open_accounts_mv.account_id, $expr1, '1970-01-01':Date, $expr9, account_groups_mv_next.base_currency, account_groups_mv_next.opening_date, account_groups_mv_next.source_entity_type, open_accounts_mv.account_id, null:Varchar, null:Date, null:Int32, null:Varchar, null:Date, null:Varchar, null:Varchar, null:Varchar, $src, account_groups_mv_next.open_accounts_mv.account_id, account_groups_mv_next.null:Varchar, account_groups_mv_next.null:Int32, account_groups_mv_next.null:Varchar#1, account_groups_mv_next.$src ]
├── stream key: [ $expr1, open_accounts_mv.account_id, open_accounts_mv.account_id, null:Varchar, null:Date, null:Int32, null:Varchar, null:Date, null:Varchar, null:Varchar, null:Varchar, $src, account_groups_mv_next.open_accounts_mv.account_id, account_groups_mv_next.null:Varchar, account_groups_mv_next.null:Int32, account_groups_mv_next.null:Varchar#1, account_groups_mv_next.$src ]
└── StreamProject { exprs: [open_accounts_mv.account_id, $expr1, '1970-01-01':Date, Case((IsNull(null:Date) AND IsNull(first_value)), null:Date, IsNull(null:Date), first_value, IsNull(first_value), null:Date, Least(null:Date, first_value)) as $expr9, account_groups_mv_next.base_currency, account_groups_mv_next.opening_date, account_groups_mv_next.source_entity_type, open_accounts_mv.account_id, null:Varchar, null:Date, null:Int32, null:Varchar, null:Date, null:Varchar, null:Varchar, null:Varchar, $src, account_groups_mv_next.open_accounts_mv.account_id, account_groups_mv_next.null:Varchar, account_groups_mv_next.null:Int32, account_groups_mv_next.null:Varchar#1, account_groups_mv_next.$src] }
    ├── output: [ open_accounts_mv.account_id, $expr1, '1970-01-01':Date, $expr9, account_groups_mv_next.base_currency, account_groups_mv_next.opening_date, account_groups_mv_next.source_entity_type, open_accounts_mv.account_id, null:Varchar, null:Date, null:Int32, null:Varchar, null:Date, null:Varchar, null:Varchar, null:Varchar, $src, account_groups_mv_next.open_accounts_mv.account_id, account_groups_mv_next.null:Varchar, account_groups_mv_next.null:Int32, account_groups_mv_next.null:Varchar#1, account_groups_mv_next.$src ]
    ├── stream key: [ $expr1, open_accounts_mv.account_id, open_accounts_mv.account_id, null:Varchar, null:Date, null:Int32, null:Varchar, null:Date, null:Varchar, null:Varchar, null:Varchar, $src, account_groups_mv_next.open_accounts_mv.account_id, account_groups_mv_next.null:Varchar, account_groups_mv_next.null:Int32, account_groups_mv_next.null:Varchar#1, account_groups_mv_next.$src ]
    └── MergeExecutor { output: [ open_accounts_mv.account_id, $expr1, '1970-01-01':Date, null:Date, first_value, account_groups_mv_next.base_currency, account_groups_mv_next.opening_date, account_groups_mv_next.source_entity_type, open_accounts_mv.account_id, null:Varchar, null:Date, null:Int32, null:Varchar, null:Date, null:Varchar, null:Varchar, null:Varchar, $src, account_groups_mv_next.account_group_id, account_groups_mv_next.open_accounts_mv.account_id, account_groups_mv_next.null:Varchar, account_groups_mv_next.null:Int32, account_groups_mv_next.null:Varchar#1, account_groups_mv_next.$src ], stream key: [ $expr1, open_accounts_mv.account_id, open_accounts_mv.account_id, null:Varchar, null:Date, null:Int32, null:Varchar, null:Date, null:Varchar, null:Varchar, null:Varchar, $src, account_groups_mv_next.open_accounts_mv.account_id, account_groups_mv_next.null:Varchar, account_groups_mv_next.null:Int32, account_groups_mv_next.null:Varchar#1, account_groups_mv_next.$src ] }

Fragment 61165 (Actor 738465,738464)
StreamSyncLogStore { output: [ open_accounts_mv.account_id, $expr1, '1970-01-01':Date, null:Date, first_value, account_groups_mv_next.base_currency, account_groups_mv_next.opening_date, account_groups_mv_next.source_entity_type, open_accounts_mv.account_id, null:Varchar, null:Date, null:Int32, null:Varchar, null:Date, null:Varchar, null:Varchar, null:Varchar, $src, account_groups_mv_next.account_group_id, account_groups_mv_next.open_accounts_mv.account_id, account_groups_mv_next.null:Varchar, account_groups_mv_next.null:Int32, account_groups_mv_next.null:Varchar#1, account_groups_mv_next.$src ], stream key: [ $expr1, open_accounts_mv.account_id, open_accounts_mv.account_id, null:Varchar, null:Date, null:Int32, null:Varchar, null:Date, null:Varchar, null:Varchar, null:Varchar, $src, account_groups_mv_next.open_accounts_mv.account_id, account_groups_mv_next.null:Varchar, account_groups_mv_next.null:Int32, account_groups_mv_next.null:Varchar#1, account_groups_mv_next.$src ] }
└── StreamHashJoin { type: Inner, predicate: $expr1 = account_groups_mv_next.account_group_id } { output: [ open_accounts_mv.account_id, $expr1, '1970-01-01':Date, null:Date, first_value, account_groups_mv_next.base_currency, account_groups_mv_next.opening_date, account_groups_mv_next.source_entity_type, open_accounts_mv.account_id, null:Varchar, null:Date, null:Int32, null:Varchar, null:Date, null:Varchar, null:Varchar, null:Varchar, $src, account_groups_mv_next.account_group_id, account_groups_mv_next.open_accounts_mv.account_id, account_groups_mv_next.null:Varchar, account_groups_mv_next.null:Int32, account_groups_mv_next.null:Varchar#1, account_groups_mv_next.$src ], stream key: [ $expr1, open_accounts_mv.account_id, open_accounts_mv.account_id, null:Varchar, null:Date, null:Int32, null:Varchar, null:Date, null:Varchar, null:Varchar, null:Varchar, $src, account_groups_mv_next.open_accounts_mv.account_id, account_groups_mv_next.null:Varchar, account_groups_mv_next.null:Int32, account_groups_mv_next.null:Varchar#1, account_groups_mv_next.$src ] }
    ├── MergeExecutor { output: [ open_accounts_mv.account_id, $expr1, '1970-01-01':Date, null:Date, null:Date, null:Date, open_accounts_mv.account_id, null:Varchar, null:Varchar, null:Varchar, null:Varchar, null:Varchar, null:Int32, $src, first_value ], stream key: [ $expr1, open_accounts_mv.account_id, open_accounts_mv.account_id, null:Varchar, null:Date, null:Int32, null:Varchar, null:Date, null:Varchar, null:Varchar, null:Varchar, $src ] }
    └── MergeExecutor { output: [ account_groups_mv_next.account_group_id, account_groups_mv_next.base_currency, account_groups_mv_next.opening_date, account_groups_mv_next.source_entity_type, account_groups_mv_next.open_accounts_mv.account_id, account_groups_mv_next.null:Varchar, account_groups_mv_next.null:Int32, account_groups_mv_next.null:Varchar#1, account_groups_mv_next.$src ], stream key: [ account_groups_mv_next.account_group_id, account_groups_mv_next.open_accounts_mv.account_id, account_groups_mv_next.null:Varchar, account_groups_mv_next.null:Int32, account_groups_mv_next.null:Varchar#1, account_groups_mv_next.$src ] }

Fragment 61166 (Actor 738469,738468)
StreamLocalityProvider { locality_columns: [$expr1] } { output: [ open_accounts_mv.account_id, $expr1, '1970-01-01':Date, null:Date, null:Date, null:Date, open_accounts_mv.account_id, null:Varchar, null:Varchar, null:Varchar, null:Varchar, null:Varchar, null:Int32, $src, first_value ], stream key: [ $expr1, open_accounts_mv.account_id, open_accounts_mv.account_id, null:Varchar, null:Date, null:Int32, null:Varchar, null:Date, null:Varchar, null:Varchar, null:Varchar, $src ] }
└── MergeExecutor { output: [ open_accounts_mv.account_id, $expr1, '1970-01-01':Date, null:Date, null:Date, null:Date, open_accounts_mv.account_id, null:Varchar, null:Varchar, null:Varchar, null:Varchar, null:Varchar, null:Int32, $src, first_value ], stream key: [ open_accounts_mv.account_id, $expr1, open_accounts_mv.account_id, null:Varchar, null:Date, null:Int32, null:Varchar, null:Date, null:Varchar, null:Varchar, null:Varchar, $src ] }

Fragment 61167 (Actor 738471,738470)
StreamOverWindow { window_functions: [first_value('1970-01-01':Date) OVER(PARTITION BY open_accounts_mv.account_id, $expr1 ORDER BY '1970-01-01':Date ASC ROWS BETWEEN 1 FOLLOWING AND 1 FOLLOWING)] } { output: [ open_accounts_mv.account_id, $expr1, '1970-01-01':Date, null:Date, null:Date, null:Date, open_accounts_mv.account_id, null:Varchar, null:Varchar, null:Varchar, null:Varchar, null:Varchar, null:Int32, $src, first_value ], stream key: [ open_accounts_mv.account_id, $expr1, open_accounts_mv.account_id, null:Varchar, null:Date, null:Int32, null:Varchar, null:Date, null:Varchar, null:Varchar, null:Varchar, $src ] }
└── StreamLocalityProvider { locality_columns: [open_accounts_mv.account_id, $expr1] } { output: [ open_accounts_mv.account_id, $expr1, '1970-01-01':Date, null:Date, null:Date, null:Date, open_accounts_mv.account_id, null:Varchar, null:Varchar, null:Varchar, null:Varchar, null:Varchar, null:Int32, $src ], stream key: [ open_accounts_mv.account_id, $expr1, open_accounts_mv.account_id, null:Varchar, null:Date, null:Int32, null:Varchar, null:Date, null:Varchar, null:Varchar, null:Varchar, $src ] }
    └── MergeExecutor { output: [ open_accounts_mv.account_id, $expr1, '1970-01-01':Date, null:Date, null:Date, null:Date, open_accounts_mv.account_id, null:Varchar, null:Varchar, null:Varchar, null:Varchar, null:Varchar, null:Int32, $src ], stream key: [ open_accounts_mv.account_id, null:Varchar, null:Date, null:Int32, null:Varchar, null:Date, null:Varchar, null:Varchar, null:Varchar, $src ] }

Fragment 61168 (Actor 738473,738472)
StreamUnion { all: true } { output: [ open_accounts_mv.account_id, $expr1, '1970-01-01':Date, null:Date, null:Date, null:Date, open_accounts_mv.account_id, null:Varchar, null:Varchar, null:Varchar, null:Varchar, null:Varchar, null:Int32, $src ], stream key: [ open_accounts_mv.account_id, null:Varchar, null:Date, null:Int32, null:Varchar, null:Date, null:Varchar, null:Varchar, null:Varchar, $src ] }
├── MergeExecutor { output: [ open_accounts_mv.account_id, $expr1, '1970-01-01':Date, null:Date, null:Date, null:Date, open_accounts_mv.account_id, null:Varchar, null:Varchar, null:Varchar, null:Varchar, null:Varchar, null:Int32, 0:Int32 ], stream key: [ open_accounts_mv.account_id ] }
├── MergeExecutor { output: [ client_account_direct_mv.account_id, $expr2, client_account_direct_mv.effective_start_date, client_account_direct_mv.effective_end_date, client_account_direct_mv.effective_start_date, null:Date, client_account_direct_mv.account_id, client_account_direct_mv.client_id, null:Varchar, null:Varchar, null:Varchar, null:Varchar, client_account_direct_mv.$src, 1:Int32 ], stream key: [ client_account_direct_mv.account_id, client_account_direct_mv.client_id, client_account_direct_mv.effective_start_date, client_account_direct_mv.$src ] }
├── MergeExecutor { output: [ client_account_via_portfolio_mv.account_id, $expr3, client_account_via_portfolio_mv.effective_start_date, client_account_via_portfolio_mv.effective_end_date, client_account_via_portfolio_mv.account_to_portfolios_dm.effective_start_date, client_account_via_portfolio_mv.clients_portfolios_dm.effective_start_date, client_account_via_portfolio_mv.account_to_portfolios_dm.account_id, client_account_via_portfolio_mv.clients_portfolios_dm.client_id, client_account_via_portfolio_mv.account_to_portfolios_dm.portfolio_id, null:Varchar, null:Varchar, null:Varchar, client_account_via_portfolio_mv.$src, 2:Int32 ], stream key: [ client_account_via_portfolio_mv.account_to_portfolios_dm.account_id, client_account_via_portfolio_mv.clients_portfolios_dm.client_id, client_account_via_portfolio_mv.account_to_portfolios_dm.portfolio_id, client_account_via_portfolio_mv.account_to_portfolios_dm.effective_start_date, client_account_via_portfolio_mv.clients_portfolios_dm.effective_start_date, client_account_via_portfolio_mv.$src ] }
├── MergeExecutor { output: [ account_to_portfolios_dm.account_id, $expr4, account_to_portfolios_dm.effective_start_date, account_to_portfolios_dm.effective_end_date, account_to_portfolios_dm.effective_start_date, null:Date, account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, null:Varchar, null:Varchar, null:Varchar, null:Varchar, null:Int32, 3: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, $expr5, account_to_portfolios_dm.effective_start_date, account_to_portfolios_dm.effective_end_date, account_to_portfolios_dm.effective_start_date, null:Date, account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, null:Varchar, null:Varchar, null:Varchar, null:Varchar, null:Int32, 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, $expr6, account_to_portfolios_dm.effective_start_date, account_to_portfolios_dm.effective_end_date, account_to_portfolios_dm.effective_start_date, null:Date, account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, null:Varchar, null:Varchar, null:Varchar, null:Varchar, null:Int32, 5: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_account_direct_mv_next.account_id, $expr7, party_account_direct_mv_next.effective_start_date, party_account_direct_mv_next.effective_end_date, null:Date, null:Date, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, null:Varchar, null:Varchar, null:Varchar, party_account_direct_mv_next.$src, 6:Int32 ], stream key: [ party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src ] }
└── MergeExecutor { output: [ party_account_via_portfolio_mv_next.account_id, $expr8, party_account_via_portfolio_mv_next.effective_start_date, party_account_via_portfolio_mv_next.effective_end_date, party_account_via_portfolio_mv_next.account_to_portfolios_dm.effective_start_date, null:Date, party_account_via_portfolio_mv_next.account_to_portfolios_dm.account_id, party_account_via_portfolio_mv_next.party_involvements_dm.party_id, party_account_via_portfolio_mv_next.account_to_portfolios_dm.portfolio_id, party_account_via_portfolio_mv_next.open_accounts_mv.account_id, party_account_via_portfolio_mv_next.customer_relationships.id, party_account_via_portfolio_mv_next.party_involvements_dm.id, party_account_via_portfolio_mv_next.$src, 7:Int32 ], stream key: [ party_account_via_portfolio_mv_next.account_to_portfolios_dm.account_id, party_account_via_portfolio_mv_next.party_involvements_dm.party_id, party_account_via_portfolio_mv_next.account_to_portfolios_dm.portfolio_id, party_account_via_portfolio_mv_next.open_accounts_mv.account_id, party_account_via_portfolio_mv_next.account_to_portfolios_dm.effective_start_date, party_account_via_portfolio_mv_next.customer_relationships.id, party_account_via_portfolio_mv_next.party_involvements_dm.id, party_account_via_portfolio_mv_next.$src ] }

Fragment 61169 (Actor 738475,738474)
StreamProject { exprs: [open_accounts_mv.account_id, ConcatOp('account_group_':Varchar, Md5(ConcatOp(open_accounts_mv.account_id, 'all':Varchar)::Bytea)) as $expr1, '1970-01-01':Date, null:Date, null:Date, null:Date, open_accounts_mv.account_id, null:Varchar, null:Varchar, null:Varchar, null:Varchar, null:Varchar, null:Int32, 0:Int32] } { output: [ open_accounts_mv.account_id, $expr1, '1970-01-01':Date, null:Date, null:Date, null:Date, open_accounts_mv.account_id, null:Varchar, null:Varchar, null:Varchar, null:Varchar, null:Varchar, null:Int32, 0:Int32 ], stream key: [ open_accounts_mv.account_id ] }
└── StreamTableScan { table: open_accounts_mv, columns: [account_id] } { output: [ open_accounts_mv.account_id ], stream key: [ open_accounts_mv.account_id ] }
    ├── Upstream { output: [ account_id ], stream key: [] }
    └── BatchPlanNode { output: [ account_id ], stream key: [] }

Fragment 61170 (Actor 738477,738476)
StreamProject { exprs: [client_account_direct_mv.account_id, ConcatOp('account_group_':Varchar, Md5(ConcatOp(client_account_direct_mv.client_id, client_account_direct_mv.type)::Bytea)) as $expr2, client_account_direct_mv.effective_start_date, client_account_direct_mv.effective_end_date, client_account_direct_mv.effective_start_date, null:Date, client_account_direct_mv.account_id, client_account_direct_mv.client_id, null:Varchar, null:Varchar, null:Varchar, null:Varchar, client_account_direct_mv.$src, 1:Int32] } { output: [ client_account_direct_mv.account_id, $expr2, client_account_direct_mv.effective_start_date, client_account_direct_mv.effective_end_date, client_account_direct_mv.effective_start_date, null:Date, client_account_direct_mv.account_id, client_account_direct_mv.client_id, null:Varchar, null:Varchar, null:Varchar, null:Varchar, client_account_direct_mv.$src, 1:Int32 ], stream key: [ client_account_direct_mv.account_id, client_account_direct_mv.client_id, client_account_direct_mv.effective_start_date, client_account_direct_mv.$src ] }
└── StreamTableScan { table: client_account_direct_mv, columns: [client_id, account_id, effective_start_date, effective_end_date, type, $src] } { output: [ client_account_direct_mv.client_id, client_account_direct_mv.account_id, client_account_direct_mv.effective_start_date, client_account_direct_mv.effective_end_date, client_account_direct_mv.type, client_account_direct_mv.$src ], stream key: [ client_account_direct_mv.account_id, client_account_direct_mv.client_id, client_account_direct_mv.effective_start_date, client_account_direct_mv.$src ] }
    ├── Upstream { output: [ client_id, account_id, effective_start_date, effective_end_date, type, $src ], stream key: [] }
    └── BatchPlanNode { output: [ client_id, account_id, effective_start_date, effective_end_date, type, $src ], stream key: [] }

Fragment 61171 (Actor 738510,738511)
StreamProject { exprs: [client_account_via_portfolio_mv.account_id, ConcatOp('account_group_':Varchar, Md5(ConcatOp(client_account_via_portfolio_mv.client_id, client_account_via_portfolio_mv.type)::Bytea)) as $expr3, client_account_via_portfolio_mv.effective_start_date, client_account_via_portfolio_mv.effective_end_date, client_account_via_portfolio_mv.account_to_portfolios_dm.effective_start_date, client_account_via_portfolio_mv.clients_portfolios_dm.effective_start_date, client_account_via_portfolio_mv.account_to_portfolios_dm.account_id, client_account_via_portfolio_mv.clients_portfolios_dm.client_id, client_account_via_portfolio_mv.account_to_portfolios_dm.portfolio_id, null:Varchar, null:Varchar, null:Varchar, client_account_via_portfolio_mv.$src, 2:Int32] }
├── output: [ client_account_via_portfolio_mv.account_id, $expr3, client_account_via_portfolio_mv.effective_start_date, client_account_via_portfolio_mv.effective_end_date, client_account_via_portfolio_mv.account_to_portfolios_dm.effective_start_date, client_account_via_portfolio_mv.clients_portfolios_dm.effective_start_date, client_account_via_portfolio_mv.account_to_portfolios_dm.account_id, client_account_via_portfolio_mv.clients_portfolios_dm.client_id, client_account_via_portfolio_mv.account_to_portfolios_dm.portfolio_id, null:Varchar, null:Varchar, null:Varchar, client_account_via_portfolio_mv.$src, 2:Int32 ]
├── stream key: [ client_account_via_portfolio_mv.account_to_portfolios_dm.account_id, client_account_via_portfolio_mv.clients_portfolios_dm.client_id, client_account_via_portfolio_mv.account_to_portfolios_dm.portfolio_id, client_account_via_portfolio_mv.account_to_portfolios_dm.effective_start_date, client_account_via_portfolio_mv.clients_portfolios_dm.effective_start_date, client_account_via_portfolio_mv.$src ]
└── StreamTableScan { table: client_account_via_portfolio_mv, columns: [client_id, account_id, effective_start_date, effective_end_date, type, account_to_portfolios_dm.account_id, clients_portfolios_dm.client_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date, clients_portfolios_dm.effective_start_date, $src] }
    ├── output: [ client_account_via_portfolio_mv.client_id, client_account_via_portfolio_mv.account_id, client_account_via_portfolio_mv.effective_start_date, client_account_via_portfolio_mv.effective_end_date, client_account_via_portfolio_mv.type, client_account_via_portfolio_mv.account_to_portfolios_dm.account_id, client_account_via_portfolio_mv.clients_portfolios_dm.client_id, client_account_via_portfolio_mv.account_to_portfolios_dm.portfolio_id, client_account_via_portfolio_mv.account_to_portfolios_dm.effective_start_date, client_account_via_portfolio_mv.clients_portfolios_dm.effective_start_date, client_account_via_portfolio_mv.$src ]
    ├── stream key: [ client_account_via_portfolio_mv.account_to_portfolios_dm.account_id, client_account_via_portfolio_mv.clients_portfolios_dm.client_id, client_account_via_portfolio_mv.account_to_portfolios_dm.portfolio_id, client_account_via_portfolio_mv.account_to_portfolios_dm.effective_start_date, client_account_via_portfolio_mv.clients_portfolios_dm.effective_start_date, client_account_via_portfolio_mv.$src ]
    ├── Upstream { output: [ client_id, account_id, effective_start_date, effective_end_date, type, account_to_portfolios_dm.account_id, clients_portfolios_dm.client_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date, clients_portfolios_dm.effective_start_date, $src ], stream key: [] }
    └── BatchPlanNode { output: [ client_id, account_id, effective_start_date, effective_end_date, type, account_to_portfolios_dm.account_id, clients_portfolios_dm.client_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date, clients_portfolios_dm.effective_start_date, $src ], stream key: [] }

Fragment 61172 (Actor 738480,738481)
StreamProject { exprs: [account_to_portfolios_dm.account_id, ConcatOp('account_group_':Varchar, Md5(ConcatOp(account_to_portfolios_dm.portfolio_id, 'all':Varchar)::Bytea)) as $expr4, account_to_portfolios_dm.effective_start_date, account_to_portfolios_dm.effective_end_date, account_to_portfolios_dm.effective_start_date, null:Date, account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, null:Varchar, null:Varchar, null:Varchar, null:Varchar, null:Int32, 3:Int32] } { output: [ account_to_portfolios_dm.account_id, $expr4, account_to_portfolios_dm.effective_start_date, account_to_portfolios_dm.effective_end_date, account_to_portfolios_dm.effective_start_date, null:Date, account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, null:Varchar, null:Varchar, null:Varchar, null:Varchar, null:Int32, 3:Int32 ], 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) OR (account_to_portfolios_dm.effective_end_date > account_to_portfolios_dm.effective_start_date)) } { output: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.disabled_at, account_to_portfolios_dm.effective_start_date, account_to_portfolios_dm.effective_end_date, open_accounts_mv.is_restricted, open_accounts_mv.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.disabled_at, account_to_portfolios_dm.effective_start_date, account_to_portfolios_dm.effective_end_date, open_accounts_mv.is_restricted, open_accounts_mv.account_id ], stream key: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ] }

Fragment 61173 (Actor 738488,738489)
StreamSyncLogStore { output: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.disabled_at, account_to_portfolios_dm.effective_start_date, account_to_portfolios_dm.effective_end_date, open_accounts_mv.is_restricted, open_accounts_mv.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.disabled_at, account_to_portfolios_dm.effective_start_date, account_to_portfolios_dm.effective_end_date, open_accounts_mv.is_restricted, open_accounts_mv.account_id ], stream key: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ] }

Fragment 61174 (Actor 738484,738485)
StreamSyncLogStore { output: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.disabled_at, account_to_portfolios_dm.effective_start_date, account_to_portfolios_dm.effective_end_date, open_accounts_mv.is_restricted, open_accounts_mv.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.disabled_at, account_to_portfolios_dm.effective_start_date, account_to_portfolios_dm.effective_end_date, open_accounts_mv.is_restricted, open_accounts_mv.account_id ], stream key: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ] }

Fragment 61175 (Actor 738486,738487)
StreamSyncLogStore { output: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.disabled_at, account_to_portfolios_dm.effective_start_date, account_to_portfolios_dm.effective_end_date, open_accounts_mv.is_restricted, open_accounts_mv.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 = open_accounts_mv.account_id } { output: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.disabled_at, account_to_portfolios_dm.effective_start_date, account_to_portfolios_dm.effective_end_date, open_accounts_mv.is_restricted, open_accounts_mv.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.disabled_at, account_to_portfolios_dm.effective_start_date, 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 ] }
    └── MergeExecutor { output: [ open_accounts_mv.account_id, open_accounts_mv.is_restricted ], stream key: [ open_accounts_mv.account_id ] }

Fragment 61176 (Actor 738490,738491)
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.disabled_at, account_to_portfolios_dm.effective_start_date, 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 ] }
└── MergeExecutor { output: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.disabled_at, account_to_portfolios_dm.effective_start_date, 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 ] }

Fragment 61177 (Actor 738513,738512)
StreamFilter { predicate: IsNull(account_to_portfolios_dm.disabled_at) AND (IsNull(account_to_portfolios_dm.effective_end_date) OR (account_to_portfolios_dm.effective_end_date > account_to_portfolios_dm.effective_start_date)) } { output: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.disabled_at, account_to_portfolios_dm.effective_start_date, 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, disabled_at, effective_start_date, effective_end_date] } { output: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.disabled_at, account_to_portfolios_dm.effective_start_date, 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, disabled_at, effective_start_date, effective_end_date ], stream key: [] }
    └── BatchPlanNode { output: [ account_id, portfolio_id, disabled_at, effective_start_date, effective_end_date ], stream key: [] }

Fragment 61178 (Actor 738515,738514)
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 61179 (Actor 738483,738482)
StreamProject { exprs: [account_to_portfolios_dm.account_id, ConcatOp('account_group_':Varchar, Md5(ConcatOp(account_to_portfolios_dm.portfolio_id, 'restricted':Varchar)::Bytea)) as $expr5, account_to_portfolios_dm.effective_start_date, account_to_portfolios_dm.effective_end_date, account_to_portfolios_dm.effective_start_date, null:Date, account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, null:Varchar, null:Varchar, null:Varchar, null:Varchar, null:Int32, 4:Int32] } { output: [ account_to_portfolios_dm.account_id, $expr5, account_to_portfolios_dm.effective_start_date, account_to_portfolios_dm.effective_end_date, account_to_portfolios_dm.effective_start_date, null:Date, account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, null:Varchar, null:Varchar, null:Varchar, null:Varchar, null:Int32, 4:Int32 ], 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) OR (account_to_portfolios_dm.effective_end_date > account_to_portfolios_dm.effective_start_date)) AND (open_accounts_mv.is_restricted = true:Boolean) } { output: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.disabled_at, account_to_portfolios_dm.effective_start_date, account_to_portfolios_dm.effective_end_date, open_accounts_mv.is_restricted, open_accounts_mv.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.disabled_at, account_to_portfolios_dm.effective_start_date, account_to_portfolios_dm.effective_end_date, open_accounts_mv.is_restricted, open_accounts_mv.account_id ], stream key: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ] }

Fragment 61180 (Actor 738479,738478)
StreamProject { exprs: [account_to_portfolios_dm.account_id, ConcatOp('account_group_':Varchar, Md5(ConcatOp(account_to_portfolios_dm.portfolio_id, 'un_restricted':Varchar)::Bytea)) as $expr6, account_to_portfolios_dm.effective_start_date, account_to_portfolios_dm.effective_end_date, account_to_portfolios_dm.effective_start_date, null:Date, account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, null:Varchar, null:Varchar, null:Varchar, null:Varchar, null:Int32, 5:Int32] } { output: [ account_to_portfolios_dm.account_id, $expr6, account_to_portfolios_dm.effective_start_date, account_to_portfolios_dm.effective_end_date, account_to_portfolios_dm.effective_start_date, null:Date, account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, null:Varchar, null:Varchar, null:Varchar, null:Varchar, null:Int32, 5:Int32 ], 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) OR (account_to_portfolios_dm.effective_end_date > account_to_portfolios_dm.effective_start_date)) AND Not(IsTrue(open_accounts_mv.is_restricted)) } { output: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.disabled_at, account_to_portfolios_dm.effective_start_date, account_to_portfolios_dm.effective_end_date, open_accounts_mv.is_restricted, open_accounts_mv.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.disabled_at, account_to_portfolios_dm.effective_start_date, account_to_portfolios_dm.effective_end_date, open_accounts_mv.is_restricted, open_accounts_mv.account_id ], stream key: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ] }

Fragment 61181 (Actor 738492,738493)
StreamProject { exprs: [party_account_direct_mv_next.account_id, ConcatOp('account_group_':Varchar, Md5(ConcatOp(party_account_direct_mv_next.party_id, party_account_direct_mv_next.type)::Bytea)) as $expr7, party_account_direct_mv_next.effective_start_date, party_account_direct_mv_next.effective_end_date, null:Date, null:Date, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, null:Varchar, null:Varchar, null:Varchar, party_account_direct_mv_next.$src, 6:Int32] } { output: [ party_account_direct_mv_next.account_id, $expr7, party_account_direct_mv_next.effective_start_date, party_account_direct_mv_next.effective_end_date, null:Date, null:Date, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, null:Varchar, null:Varchar, null:Varchar, party_account_direct_mv_next.$src, 6:Int32 ], stream key: [ party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src ] }
└── StreamTableScan { table: party_account_direct_mv_next, columns: [party_id, account_id, effective_start_date, effective_end_date, type, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.id, $src] } { output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.account_id, party_account_direct_mv_next.effective_start_date, party_account_direct_mv_next.effective_end_date, party_account_direct_mv_next.type, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src ], stream key: [ party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.$src ] }
    ├── Upstream { output: [ party_id, account_id, effective_start_date, effective_end_date, type, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.id, $src ], stream key: [] }
    └── BatchPlanNode { output: [ party_id, account_id, effective_start_date, effective_end_date, type, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.id, $src ], stream key: [] }

Fragment 61182 (Actor 738516,738517)
StreamProject { exprs: [party_account_via_portfolio_mv_next.account_id, ConcatOp('account_group_':Varchar, Md5(ConcatOp(party_account_via_portfolio_mv_next.party_id, party_account_via_portfolio_mv_next.type)::Bytea)) as $expr8, party_account_via_portfolio_mv_next.effective_start_date, party_account_via_portfolio_mv_next.effective_end_date, party_account_via_portfolio_mv_next.account_to_portfolios_dm.effective_start_date, null:Date, party_account_via_portfolio_mv_next.account_to_portfolios_dm.account_id, party_account_via_portfolio_mv_next.party_involvements_dm.party_id, party_account_via_portfolio_mv_next.account_to_portfolios_dm.portfolio_id, party_account_via_portfolio_mv_next.open_accounts_mv.account_id, party_account_via_portfolio_mv_next.customer_relationships.id, party_account_via_portfolio_mv_next.party_involvements_dm.id, party_account_via_portfolio_mv_next.$src, 7:Int32] }
├── output: [ party_account_via_portfolio_mv_next.account_id, $expr8, party_account_via_portfolio_mv_next.effective_start_date, party_account_via_portfolio_mv_next.effective_end_date, party_account_via_portfolio_mv_next.account_to_portfolios_dm.effective_start_date, null:Date, party_account_via_portfolio_mv_next.account_to_portfolios_dm.account_id, party_account_via_portfolio_mv_next.party_involvements_dm.party_id, party_account_via_portfolio_mv_next.account_to_portfolios_dm.portfolio_id, party_account_via_portfolio_mv_next.open_accounts_mv.account_id, party_account_via_portfolio_mv_next.customer_relationships.id, party_account_via_portfolio_mv_next.party_involvements_dm.id, party_account_via_portfolio_mv_next.$src, 7:Int32 ]
├── stream key: [ party_account_via_portfolio_mv_next.account_to_portfolios_dm.account_id, party_account_via_portfolio_mv_next.party_involvements_dm.party_id, party_account_via_portfolio_mv_next.account_to_portfolios_dm.portfolio_id, party_account_via_portfolio_mv_next.open_accounts_mv.account_id, party_account_via_portfolio_mv_next.account_to_portfolios_dm.effective_start_date, party_account_via_portfolio_mv_next.customer_relationships.id, party_account_via_portfolio_mv_next.party_involvements_dm.id, party_account_via_portfolio_mv_next.$src ]
└── StreamTableScan { table: party_account_via_portfolio_mv_next, columns: [party_id, account_id, effective_start_date, effective_end_date, type, account_to_portfolios_dm.account_id, party_involvements_dm.party_id, account_to_portfolios_dm.portfolio_id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date, customer_relationships.id, party_involvements_dm.id, $src] }
    ├── output: [ party_account_via_portfolio_mv_next.party_id, party_account_via_portfolio_mv_next.account_id, party_account_via_portfolio_mv_next.effective_start_date, party_account_via_portfolio_mv_next.effective_end_date, party_account_via_portfolio_mv_next.type, party_account_via_portfolio_mv_next.account_to_portfolios_dm.account_id, party_account_via_portfolio_mv_next.party_involvements_dm.party_id, party_account_via_portfolio_mv_next.account_to_portfolios_dm.portfolio_id, party_account_via_portfolio_mv_next.open_accounts_mv.account_id, party_account_via_portfolio_mv_next.account_to_portfolios_dm.effective_start_date, party_account_via_portfolio_mv_next.customer_relationships.id, party_account_via_portfolio_mv_next.party_involvements_dm.id, party_account_via_portfolio_mv_next.$src ]
    ├── stream key: [ party_account_via_portfolio_mv_next.account_to_portfolios_dm.account_id, party_account_via_portfolio_mv_next.party_involvements_dm.party_id, party_account_via_portfolio_mv_next.account_to_portfolios_dm.portfolio_id, party_account_via_portfolio_mv_next.open_accounts_mv.account_id, party_account_via_portfolio_mv_next.account_to_portfolios_dm.effective_start_date, party_account_via_portfolio_mv_next.customer_relationships.id, party_account_via_portfolio_mv_next.party_involvements_dm.id, party_account_via_portfolio_mv_next.$src ]
    ├── Upstream { output: [ party_id, account_id, effective_start_date, effective_end_date, type, account_to_portfolios_dm.account_id, party_involvements_dm.party_id, account_to_portfolios_dm.portfolio_id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date, customer_relationships.id, party_involvements_dm.id, $src ], stream key: [] }
    └── BatchPlanNode { output: [ party_id, account_id, effective_start_date, effective_end_date, type, account_to_portfolios_dm.account_id, party_involvements_dm.party_id, account_to_portfolios_dm.portfolio_id, open_accounts_mv.account_id, account_to_portfolios_dm.effective_start_date, customer_relationships.id, party_involvements_dm.id, $src ], stream key: [] }

Fragment 61183 (Actor 738495,738494)
StreamLocalityProvider { locality_columns: [account_groups_mv_next.account_group_id] } { output: [ account_groups_mv_next.account_group_id, account_groups_mv_next.base_currency, account_groups_mv_next.opening_date, account_groups_mv_next.source_entity_type, account_groups_mv_next.open_accounts_mv.account_id, account_groups_mv_next.null:Varchar, account_groups_mv_next.null:Int32, account_groups_mv_next.null:Varchar#1, account_groups_mv_next.$src ], stream key: [ account_groups_mv_next.account_group_id, account_groups_mv_next.open_accounts_mv.account_id, account_groups_mv_next.null:Varchar, account_groups_mv_next.null:Int32, account_groups_mv_next.null:Varchar#1, account_groups_mv_next.$src ] }
└── MergeExecutor { output: [ account_groups_mv_next.account_group_id, account_groups_mv_next.base_currency, account_groups_mv_next.opening_date, account_groups_mv_next.source_entity_type, account_groups_mv_next.open_accounts_mv.account_id, account_groups_mv_next.null:Varchar, account_groups_mv_next.null:Int32, account_groups_mv_next.null:Varchar#1, account_groups_mv_next.$src ], stream key: [ account_groups_mv_next.open_accounts_mv.account_id, account_groups_mv_next.null:Varchar, account_groups_mv_next.null:Int32, account_groups_mv_next.null:Varchar#1, account_groups_mv_next.$src ] }

Fragment 61184 (Actor 738497,738496)
StreamTableScan { table: account_groups_mv_next, columns: [account_group_id, base_currency, opening_date, source_entity_type, open_accounts_mv.account_id, null:Varchar, null:Int32, null:Varchar#1, $src] } { output: [ account_groups_mv_next.account_group_id, account_groups_mv_next.base_currency, account_groups_mv_next.opening_date, account_groups_mv_next.source_entity_type, account_groups_mv_next.open_accounts_mv.account_id, account_groups_mv_next.null:Varchar, account_groups_mv_next.null:Int32, account_groups_mv_next.null:Varchar#1, account_groups_mv_next.$src ], stream key: [ account_groups_mv_next.open_accounts_mv.account_id, account_groups_mv_next.null:Varchar, account_groups_mv_next.null:Int32, account_groups_mv_next.null:Varchar#1, account_groups_mv_next.$src ] }
├── Upstream { output: [ account_group_id, base_currency, opening_date, source_entity_type, open_accounts_mv.account_id, null:Varchar, null:Int32, null:Varchar#1, $src ], stream key: [] }
└── BatchPlanNode { output: [ account_group_id, base_currency, opening_date, source_entity_type, open_accounts_mv.account_id, null:Varchar, null:Int32, null:Varchar#1, $src ], stream key: [] }