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

← cluster alinma_bff objects unlinked_account_pairs_mv explain
Overview Objects Graph History
materialized view · alinma_bff.unlinked_account_pairs_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
113 operators
Materialize · alinma_bff.unlinked_account_pairs_mv
0% idle 2 actors
Project
2 actors
MaterializedExprs
0% idle 2 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
Project · IsNull(account_to_portfolios_dm.account_id)
2 actors
Filter · IsNull(account_to_portfolios_dm.account_id)
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · LeftOuter · accounts_dm.account_id = account_to_portfolios_dm.account_id
2 actors
HashJoin · LeftOuter · accounts_dm.account_id = account_to_portfolios_dm.account_id Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
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
2% 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
2% idle 2 actors
StreamScan · account_to_portfolios_dm
2% 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 · Inner · reference_identifiers.entity_id = accounts_dm.account_id AN…
2 actors
HashJoin · Inner · reference_identifiers.entity_id = accounts_dm.account_id AN… 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 · Inner · account_primary_client_mv_next.account_id = accounts_dm.acc…
2 actors
HashJoin · Inner · account_primary_client_mv_next.account_id = accounts_dm.acc… 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 · Inner · product_types_dm.product_type_id = accounts_dm.product_type…
2 actors
HashJoin · Inner · product_types_dm.product_type_id = accounts_dm.product_type… 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 · accounts_dm
2 actors
Filter · accounts_dm
0% idle 2 actors
StreamScan · accounts_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · product_types_dm
2 actors
Filter · product_types_dm
0% idle 2 actors
StreamScan · product_types_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · account_primary_client_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
SyncLogStore · Inner · accounts_dm.number = reference_identifiers.value
2 actors
HashJoin · Inner · accounts_dm.number = reference_identifiers.value 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 · reference_identifiers
2 actors
Filter · reference_identifiers
0% idle 2 actors
StreamScan · reference_identifiers
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 · Inner · account_primary_client_mv_next.account_id = accounts_dm.acc…
2 actors
HashJoin · Inner · account_primary_client_mv_next.account_id = accounts_dm.acc… Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · accounts_dm
2 actors
Filter · accounts_dm
0% idle 2 actors
StreamScan · accounts_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · account_primary_client_mv_next
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 · alinma_bff.unlinked_account_pairs_mv Materialize alinma_bff.unlinked_acc… idle · 2 actors Project Project — · 2 actors MaterializedExprs MaterializedExprs idle · 2 actors Project Project — · 2 actors HashAgg HashAgg idle · 2 actors LocalityProvider LocalityProvider idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · IsNull(account_to_portfolios_dm.account_id) Project IsNull(account_to_portf… — · 2 actors Filter · IsNull(account_to_portfolios_dm.account_id) Filter IsNull(account_to_portf… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · LeftOuter · accounts_dm.account_id = account_to_portfolios_dm.account_id SyncLogStore LeftOuter · accounts_dm… — · 2 actors HashJoin · LeftOuter · accounts_dm.account_id = account_to_portfolios_dm.account_id HashJoin LeftOuter · accounts_dm… 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 LocalityProvider LocalityProvider idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · reference_identifiers.entity_id = accounts_dm.account_id AN… SyncLogStore Inner · reference_ident… — · 2 actors HashJoin · Inner · reference_identifiers.entity_id = accounts_dm.account_id AN… HashJoin Inner · reference_ident… 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 · Inner · account_primary_client_mv_next.account_id = accounts_dm.acc… SyncLogStore Inner · account_primary… — · 2 actors HashJoin · Inner · account_primary_client_mv_next.account_id = accounts_dm.acc… HashJoin Inner · account_primary… 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 · Inner · product_types_dm.product_type_id = accounts_dm.product_type… SyncLogStore Inner · product_types_d… — · 2 actors HashJoin · Inner · product_types_dm.product_type_id = accounts_dm.product_type… HashJoin Inner · product_types_d… 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 · accounts_dm Project accounts_dm — · 2 actors Filter · accounts_dm Filter accounts_dm idle · 2 actors StreamScan · accounts_dm StreamScan accounts_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · product_types_dm Project product_types_dm — · 2 actors Filter · product_types_dm Filter product_types_dm idle · 2 actors StreamScan · product_types_dm StreamScan product_types_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · account_primary_client_mv_next StreamScan account_primary_client_… 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 · Inner · accounts_dm.number = reference_identifiers.value SyncLogStore Inner · accounts_dm.num… — · 2 actors HashJoin · Inner · accounts_dm.number = reference_identifiers.value HashJoin Inner · accounts_dm.num… 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 · reference_identifiers Project reference_identifiers — · 2 actors Filter · reference_identifiers Filter reference_identifiers idle · 2 actors StreamScan · reference_identifiers StreamScan reference_identifiers 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 · Inner · account_primary_client_mv_next.account_id = accounts_dm.acc… SyncLogStore Inner · account_primary… — · 2 actors HashJoin · Inner · account_primary_client_mv_next.account_id = accounts_dm.acc… HashJoin Inner · account_primary… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · accounts_dm Project accounts_dm — · 2 actors Filter · accounts_dm Filter accounts_dm idle · 2 actors StreamScan · accounts_dm StreamScan accounts_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · account_primary_client_mv_next StreamScan account_primary_client_… 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 54324 (Actor 739163,739162)
StreamMaterialize { columns: [id, security_account_id, cash_account_id, registered_by, security_account, cash_account], stream_key: [security_account_id, cash_account_id], pk_columns: [security_account_id, cash_account_id], pk_conflict: NoCheck }
├── output: [ $expr5, accounts_dm.account_id, accounts_dm.account_id, 'risingwave-pipeline':Varchar, $expr6, $expr7 ]
├── stream key: [ accounts_dm.account_id, accounts_dm.account_id ]
└── StreamProject { exprs: [$expr5, accounts_dm.account_id, accounts_dm.account_id, 'risingwave-pipeline':Varchar, JsonbAccess(jsonb_agg($expr3 order_by(reference_identifiers.id ASC)), 0:Int32) as $expr6, JsonbAccess(jsonb_agg($expr4 order_by(reference_identifiers.id ASC)), 0:Int32) as $expr7] }
    ├── output: [ $expr5, accounts_dm.account_id, accounts_dm.account_id, 'risingwave-pipeline':Varchar, $expr6, $expr7 ]
    ├── stream key: [ accounts_dm.account_id, accounts_dm.account_id ]
    └── StreamMaterializedExprs { exprs: [id_from_string_with_prefix('unlinked_account_pair':Varchar, ConcatOp(ConcatOp(accounts_dm.account_id, ':':Varchar), accounts_dm.account_id)) as $expr5] }
        ├── output: [ accounts_dm.account_id, accounts_dm.account_id, jsonb_agg($expr3 order_by(reference_identifiers.id ASC)), jsonb_agg($expr4 order_by(reference_identifiers.id ASC)), $expr5 ]
        ├── stream key: [ accounts_dm.account_id, accounts_dm.account_id ]
        └── StreamProject { exprs: [accounts_dm.account_id, accounts_dm.account_id, jsonb_agg($expr3 order_by(reference_identifiers.id ASC)), jsonb_agg($expr4 order_by(reference_identifiers.id ASC))] }
            ├── output: [ accounts_dm.account_id, accounts_dm.account_id, jsonb_agg($expr3 order_by(reference_identifiers.id ASC)), jsonb_agg($expr4 order_by(reference_identifiers.id ASC)) ]
            ├── stream key: [ accounts_dm.account_id, accounts_dm.account_id ]
            └── StreamHashAgg { group_key: [accounts_dm.account_id, accounts_dm.account_id], aggs: [jsonb_agg($expr3 order_by(reference_identifiers.id ASC)), jsonb_agg($expr4 order_by(reference_identifiers.id ASC)), count] }
                ├── output: [ accounts_dm.account_id, accounts_dm.account_id, jsonb_agg($expr3 order_by(reference_identifiers.id ASC)), jsonb_agg($expr4 order_by(reference_identifiers.id ASC)), count ]
                ├── stream key: [ accounts_dm.account_id, accounts_dm.account_id ]
                └── StreamLocalityProvider { locality_columns: [accounts_dm.account_id, accounts_dm.account_id] }
                    ├── output:
                    │   ┌── accounts_dm.account_id
                    │   ├── accounts_dm.account_id
                    │   ├── reference_identifiers.id
                    │   ├── $expr3
                    │   ├── $expr4
                    │   ├── reference_identifiers.entity_id
                    │   ├── accounts_dm.number
                    │   ├── account_primary_client_mv_next.account_id
                    │   ├── product_types_dm.product_type_id
                    │   ├── account_to_portfolios_dm.portfolio_id
                    │   └── account_to_portfolios_dm.effective_start_date
                    ├── stream key:
                    │   ┌── accounts_dm.account_id
                    │   ├── accounts_dm.account_id
                    │   ├── reference_identifiers.entity_id
                    │   ├── accounts_dm.number
                    │   ├── account_primary_client_mv_next.account_id
                    │   ├── reference_identifiers.id
                    │   ├── product_types_dm.product_type_id
                    │   ├── account_to_portfolios_dm.portfolio_id
                    │   └── account_to_portfolios_dm.effective_start_date
                    └── MergeExecutor
                        ├── output:
                        │   ┌── accounts_dm.account_id
                        │   ├── accounts_dm.account_id
                        │   ├── reference_identifiers.id
                        │   ├── $expr3
                        │   ├── $expr4
                        │   ├── reference_identifiers.entity_id
                        │   ├── accounts_dm.number
                        │   ├── account_primary_client_mv_next.account_id
                        │   ├── product_types_dm.product_type_id
                        │   ├── account_to_portfolios_dm.portfolio_id
                        │   └── account_to_portfolios_dm.effective_start_date
                        └── stream key:
                            ┌── accounts_dm.account_id
                            ├── reference_identifiers.entity_id
                            ├── accounts_dm.number
                            ├── account_primary_client_mv_next.account_id
                            ├── reference_identifiers.id
                            ├── product_types_dm.product_type_id
                            ├── account_to_portfolios_dm.portfolio_id
                            └── account_to_portfolios_dm.effective_start_date

Fragment 54325 (Actor 739428,739427)
StreamProject { exprs: [accounts_dm.account_id, accounts_dm.account_id, reference_identifiers.id, JsonbBuildObject('id':Varchar, accounts_dm.account_id, 'external_id':Varchar, '':Varchar, 'account_number':Varchar, accounts_dm.number, 'currency_code':Varchar, Coalesce(accounts_dm.base_currency_code, '':Varchar), 'opened_at':Varchar, Case(IsNull(accounts_dm.opening_date), null:Jsonb, JsonbBuildObject('seconds':Varchar, Extract('EPOCH':Varchar, accounts_dm.opening_date::Timestamp)::Int64)), 'client':Varchar, account_primary_client_mv_next.client_json) as $expr3, JsonbBuildObject('id':Varchar, accounts_dm.account_id, 'external_id':Varchar, '':Varchar, 'account_number':Varchar, accounts_dm.number, 'currency_code':Varchar, Coalesce(accounts_dm.base_currency_code, '':Varchar), 'opened_at':Varchar, Case(IsNull(accounts_dm.opening_date), null:Jsonb, JsonbBuildObject('seconds':Varchar, Extract('EPOCH':Varchar, accounts_dm.opening_date::Timestamp)::Int64)), 'client':Varchar, account_primary_client_mv_next.client_json) as $expr4, reference_identifiers.entity_id, accounts_dm.number, account_primary_client_mv_next.account_id, product_types_dm.product_type_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date] }
├── output: [ accounts_dm.account_id, accounts_dm.account_id, reference_identifiers.id, $expr3, $expr4, reference_identifiers.entity_id, accounts_dm.number, account_primary_client_mv_next.account_id, product_types_dm.product_type_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ]
├── stream key: [ accounts_dm.account_id, reference_identifiers.entity_id, accounts_dm.number, account_primary_client_mv_next.account_id, reference_identifiers.id, product_types_dm.product_type_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ]
└── StreamFilter { predicate: IsNull(account_to_portfolios_dm.account_id) } { output: [ reference_identifiers.id, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, account_primary_client_mv_next.client_json, account_primary_client_mv_next.client_json, account_to_portfolios_dm.account_id, reference_identifiers.entity_id, account_primary_client_mv_next.account_id, product_types_dm.product_type_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ], stream key: [ accounts_dm.account_id, reference_identifiers.entity_id, accounts_dm.number, account_primary_client_mv_next.account_id, reference_identifiers.id, product_types_dm.product_type_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ] }
    └── MergeExecutor { output: [ reference_identifiers.id, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, account_primary_client_mv_next.client_json, account_primary_client_mv_next.client_json, account_to_portfolios_dm.account_id, reference_identifiers.entity_id, account_primary_client_mv_next.account_id, product_types_dm.product_type_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ], stream key: [ accounts_dm.account_id, reference_identifiers.entity_id, accounts_dm.number, account_primary_client_mv_next.account_id, reference_identifiers.id, product_types_dm.product_type_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ] }

Fragment 54326 (Actor 739425,739426)
StreamSyncLogStore { output: [ reference_identifiers.id, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, account_primary_client_mv_next.client_json, account_primary_client_mv_next.client_json, account_to_portfolios_dm.account_id, reference_identifiers.entity_id, account_primary_client_mv_next.account_id, product_types_dm.product_type_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ], stream key: [ accounts_dm.account_id, reference_identifiers.entity_id, accounts_dm.number, account_primary_client_mv_next.account_id, reference_identifiers.id, product_types_dm.product_type_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ] }
└── StreamHashJoin { type: LeftOuter, predicate: accounts_dm.account_id = account_to_portfolios_dm.account_id } { output: [ reference_identifiers.id, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, account_primary_client_mv_next.client_json, account_primary_client_mv_next.client_json, account_to_portfolios_dm.account_id, reference_identifiers.entity_id, account_primary_client_mv_next.account_id, product_types_dm.product_type_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ], stream key: [ accounts_dm.account_id, reference_identifiers.entity_id, accounts_dm.number, account_primary_client_mv_next.account_id, reference_identifiers.id, product_types_dm.product_type_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ] }
    ├── MergeExecutor { output: [ reference_identifiers.id, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, account_primary_client_mv_next.client_json, account_primary_client_mv_next.client_json, reference_identifiers.entity_id, account_primary_client_mv_next.account_id, account_primary_client_mv_next.account_id, product_types_dm.product_type_id ], stream key: [ accounts_dm.account_id, reference_identifiers.entity_id, accounts_dm.number, account_primary_client_mv_next.account_id, reference_identifiers.id, product_types_dm.product_type_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 54327 (Actor 739429,739430)
StreamLocalityProvider { locality_columns: [accounts_dm.account_id] } { output: [ reference_identifiers.id, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, account_primary_client_mv_next.client_json, account_primary_client_mv_next.client_json, reference_identifiers.entity_id, account_primary_client_mv_next.account_id, account_primary_client_mv_next.account_id, product_types_dm.product_type_id ], stream key: [ accounts_dm.account_id, reference_identifiers.entity_id, accounts_dm.number, account_primary_client_mv_next.account_id, reference_identifiers.id, product_types_dm.product_type_id ] }
└── MergeExecutor { output: [ reference_identifiers.id, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, account_primary_client_mv_next.client_json, account_primary_client_mv_next.client_json, reference_identifiers.entity_id, account_primary_client_mv_next.account_id, account_primary_client_mv_next.account_id, product_types_dm.product_type_id ], stream key: [ reference_identifiers.entity_id, reference_identifiers.entity_id, accounts_dm.number, account_primary_client_mv_next.account_id, reference_identifiers.id, product_types_dm.product_type_id ] }

Fragment 54328 (Actor 739446,739445)
StreamSyncLogStore { output: [ reference_identifiers.id, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, account_primary_client_mv_next.client_json, account_primary_client_mv_next.client_json, reference_identifiers.entity_id, account_primary_client_mv_next.account_id, account_primary_client_mv_next.account_id, product_types_dm.product_type_id ], stream key: [ reference_identifiers.entity_id, reference_identifiers.entity_id, accounts_dm.number, account_primary_client_mv_next.account_id, reference_identifiers.id, product_types_dm.product_type_id ] }
└── StreamHashJoin { type: Inner, predicate: reference_identifiers.entity_id = accounts_dm.account_id AND reference_identifiers.entity_id = account_primary_client_mv_next.account_id } { output: [ reference_identifiers.id, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, account_primary_client_mv_next.client_json, account_primary_client_mv_next.client_json, reference_identifiers.entity_id, account_primary_client_mv_next.account_id, account_primary_client_mv_next.account_id, product_types_dm.product_type_id ], stream key: [ reference_identifiers.entity_id, reference_identifiers.entity_id, accounts_dm.number, account_primary_client_mv_next.account_id, reference_identifiers.id, product_types_dm.product_type_id ] }
    ├── MergeExecutor { output: [ account_primary_client_mv_next.client_json, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, reference_identifiers.id, reference_identifiers.entity_id, account_primary_client_mv_next.account_id, reference_identifiers.value ], stream key: [ reference_identifiers.entity_id, reference_identifiers.entity_id, accounts_dm.number, account_primary_client_mv_next.account_id, reference_identifiers.id ] }
    └── MergeExecutor { output: [ account_primary_client_mv_next.account_id, account_primary_client_mv_next.client_json, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, product_types_dm.product_type_id ], stream key: [ accounts_dm.account_id, account_primary_client_mv_next.account_id, product_types_dm.product_type_id ] }

Fragment 54329 (Actor 739447,739448)
StreamLocalityProvider { locality_columns: [reference_identifiers.entity_id, reference_identifiers.entity_id] } { output: [ account_primary_client_mv_next.client_json, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, reference_identifiers.id, reference_identifiers.entity_id, account_primary_client_mv_next.account_id, reference_identifiers.value ], stream key: [ reference_identifiers.entity_id, reference_identifiers.entity_id, accounts_dm.number, account_primary_client_mv_next.account_id, reference_identifiers.id ] }
└── MergeExecutor { output: [ account_primary_client_mv_next.client_json, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, reference_identifiers.id, reference_identifiers.entity_id, account_primary_client_mv_next.account_id, reference_identifiers.value ], stream key: [ accounts_dm.number, account_primary_client_mv_next.account_id, reference_identifiers.id ] }

Fragment 54330 (Actor 739455,739456)
StreamSyncLogStore { output: [ account_primary_client_mv_next.client_json, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, reference_identifiers.id, reference_identifiers.entity_id, account_primary_client_mv_next.account_id, reference_identifiers.value ], stream key: [ accounts_dm.number, account_primary_client_mv_next.account_id, reference_identifiers.id ] }
└── StreamHashJoin { type: Inner, predicate: accounts_dm.number = reference_identifiers.value } { output: [ account_primary_client_mv_next.client_json, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, reference_identifiers.id, reference_identifiers.entity_id, account_primary_client_mv_next.account_id, reference_identifiers.value ], stream key: [ accounts_dm.number, account_primary_client_mv_next.account_id, reference_identifiers.id ] }
    ├── MergeExecutor { output: [ account_primary_client_mv_next.client_json, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, account_primary_client_mv_next.account_id ], stream key: [ accounts_dm.number, account_primary_client_mv_next.account_id ] }
    └── MergeExecutor { output: [ reference_identifiers.id, reference_identifiers.entity_id, reference_identifiers.value ], stream key: [ reference_identifiers.value, reference_identifiers.id ] }

Fragment 54331 (Actor 739464,739463)
StreamLocalityProvider { locality_columns: [accounts_dm.number] } { output: [ account_primary_client_mv_next.client_json, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, account_primary_client_mv_next.account_id ], stream key: [ accounts_dm.number, account_primary_client_mv_next.account_id ] }
└── MergeExecutor { output: [ account_primary_client_mv_next.client_json, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, account_primary_client_mv_next.account_id ], stream key: [ account_primary_client_mv_next.account_id ] }

Fragment 54332 (Actor 739467,739468)
StreamSyncLogStore { output: [ account_primary_client_mv_next.client_json, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, account_primary_client_mv_next.account_id ], stream key: [ account_primary_client_mv_next.account_id ] }
└── StreamHashJoin { type: Inner, predicate: account_primary_client_mv_next.account_id = accounts_dm.account_id } { output: [ account_primary_client_mv_next.client_json, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, account_primary_client_mv_next.account_id ], stream key: [ account_primary_client_mv_next.account_id ] }
    ├── MergeExecutor { output: [ account_primary_client_mv_next.account_id, account_primary_client_mv_next.client_json ], stream key: [ account_primary_client_mv_next.account_id ] }
    └── MergeExecutor { output: [ accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date ], stream key: [ accounts_dm.account_id ] }

Fragment 54333 (Actor 738648,738649)
StreamTableScan { table: account_primary_client_mv_next, columns: [account_id, client_json] } { output: [ account_primary_client_mv_next.account_id, account_primary_client_mv_next.client_json ], stream key: [ account_primary_client_mv_next.account_id ] }
├── Upstream { output: [ account_id, client_json ], stream key: [] }
└── BatchPlanNode { output: [ account_id, client_json ], stream key: [] }

Fragment 54334 (Actor 739591,739592)
StreamProject { exprs: [accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date] } { output: [ accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date ], stream key: [ accounts_dm.account_id ] }
└── StreamFilter { predicate: IsNull(accounts_dm.disabled_at) } { output: [ accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, accounts_dm.disabled_at ], stream key: [ accounts_dm.account_id ] }
    └── StreamTableScan { table: accounts_dm, columns: [account_id, number, base_currency_code, opening_date, disabled_at] } { output: [ accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, accounts_dm.disabled_at ], stream key: [ accounts_dm.account_id ] }
        ├── Upstream { output: [ account_id, number, base_currency_code, opening_date, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ account_id, number, base_currency_code, opening_date, disabled_at ], stream key: [] }

Fragment 54335 (Actor 739495,739496)
StreamLocalityProvider { locality_columns: [reference_identifiers.value] } { output: [ reference_identifiers.id, reference_identifiers.entity_id, reference_identifiers.value ], stream key: [ reference_identifiers.value, reference_identifiers.id ] }
└── MergeExecutor { output: [ reference_identifiers.id, reference_identifiers.entity_id, reference_identifiers.value ], stream key: [ reference_identifiers.id ] }

Fragment 54336 (Actor 739599,739600)
StreamProject { exprs: [reference_identifiers.id, reference_identifiers.entity_id, reference_identifiers.value] } { output: [ reference_identifiers.id, reference_identifiers.entity_id, reference_identifiers.value ], stream key: [ reference_identifiers.id ] }
└── StreamFilter { predicate: IsNotNull(reference_identifiers.entity_id) AND (reference_identifiers.entity_type = 'account':Varchar) AND (reference_identifiers.key = 'cash_account_num':Varchar) } { output: [ reference_identifiers.id, reference_identifiers.entity_id, reference_identifiers.value, reference_identifiers.entity_type, reference_identifiers.key ], stream key: [ reference_identifiers.id ] }
    └── StreamTableScan { table: reference_identifiers, columns: [id, entity_id, value, entity_type, key] } { output: [ reference_identifiers.id, reference_identifiers.entity_id, reference_identifiers.value, reference_identifiers.entity_type, reference_identifiers.key ], stream key: [ reference_identifiers.id ] }
        ├── Upstream { output: [ id, entity_id, value, entity_type, key ], stream key: [] }
        └── BatchPlanNode { output: [ id, entity_id, value, entity_type, key ], stream key: [] }

Fragment 54337 (Actor 739525,739526)
StreamLocalityProvider { locality_columns: [accounts_dm.account_id, account_primary_client_mv_next.account_id] } { output: [ account_primary_client_mv_next.account_id, account_primary_client_mv_next.client_json, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, product_types_dm.product_type_id ], stream key: [ accounts_dm.account_id, account_primary_client_mv_next.account_id, product_types_dm.product_type_id ] }
└── MergeExecutor { output: [ account_primary_client_mv_next.account_id, account_primary_client_mv_next.client_json, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, product_types_dm.product_type_id ], stream key: [ account_primary_client_mv_next.account_id, product_types_dm.product_type_id ] }

Fragment 54338 (Actor 739527,739528)
StreamSyncLogStore { output: [ account_primary_client_mv_next.account_id, account_primary_client_mv_next.client_json, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, product_types_dm.product_type_id ], stream key: [ account_primary_client_mv_next.account_id, product_types_dm.product_type_id ] }
└── StreamHashJoin { type: Inner, predicate: account_primary_client_mv_next.account_id = accounts_dm.account_id } { output: [ account_primary_client_mv_next.account_id, account_primary_client_mv_next.client_json, accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, product_types_dm.product_type_id ], stream key: [ account_primary_client_mv_next.account_id, product_types_dm.product_type_id ] }
    ├── MergeExecutor { output: [ account_primary_client_mv_next.account_id, account_primary_client_mv_next.client_json ], stream key: [ account_primary_client_mv_next.account_id ] }
    └── MergeExecutor { output: [ accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, product_types_dm.product_type_id, accounts_dm.product_type_id ], stream key: [ accounts_dm.account_id, product_types_dm.product_type_id ] }

Fragment 54339 (Actor 739530,739529)
StreamTableScan { table: account_primary_client_mv_next, columns: [account_id, client_json] } { output: [ account_primary_client_mv_next.account_id, account_primary_client_mv_next.client_json ], stream key: [ account_primary_client_mv_next.account_id ] }
├── Upstream { output: [ account_id, client_json ], stream key: [] }
└── BatchPlanNode { output: [ account_id, client_json ], stream key: [] }

Fragment 54340 (Actor 739539,739540)
StreamLocalityProvider { locality_columns: [accounts_dm.account_id] } { output: [ accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, product_types_dm.product_type_id, accounts_dm.product_type_id ], stream key: [ accounts_dm.account_id, product_types_dm.product_type_id ] }
└── MergeExecutor { output: [ accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, product_types_dm.product_type_id, accounts_dm.product_type_id ], stream key: [ product_types_dm.product_type_id, accounts_dm.account_id ] }

Fragment 54341 (Actor 739542,739541)
StreamSyncLogStore { output: [ accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, product_types_dm.product_type_id, accounts_dm.product_type_id ], stream key: [ product_types_dm.product_type_id, accounts_dm.account_id ] }
└── StreamHashJoin { type: Inner, predicate: product_types_dm.product_type_id = accounts_dm.product_type_id } { output: [ accounts_dm.account_id, accounts_dm.number, accounts_dm.base_currency_code, accounts_dm.opening_date, product_types_dm.product_type_id, accounts_dm.product_type_id ], stream key: [ product_types_dm.product_type_id, accounts_dm.account_id ] }
    ├── MergeExecutor { output: [ product_types_dm.product_type_id ], stream key: [ product_types_dm.product_type_id ] }
    └── MergeExecutor { output: [ accounts_dm.account_id, accounts_dm.number, accounts_dm.product_type_id, accounts_dm.base_currency_code, accounts_dm.opening_date ], stream key: [ accounts_dm.product_type_id, accounts_dm.account_id ] }

Fragment 54342 (Actor 739553,739554)
StreamProject { exprs: [product_types_dm.product_type_id] } { output: [ product_types_dm.product_type_id ], stream key: [ product_types_dm.product_type_id ] }
└── StreamFilter { predicate: (product_types_dm.type = 'INVESTMENT':Varchar) } { output: [ product_types_dm.product_type_id, product_types_dm.type ], stream key: [ product_types_dm.product_type_id ] }
    └── StreamTableScan { table: product_types_dm, columns: [product_type_id, type] } { output: [ product_types_dm.product_type_id, product_types_dm.type ], stream key: [ product_types_dm.product_type_id ] }
        ├── Upstream { output: [ product_type_id, type ], stream key: [] }
        └── BatchPlanNode { output: [ product_type_id, type ], stream key: [] }

Fragment 54343 (Actor 739550,739549)
StreamLocalityProvider { locality_columns: [accounts_dm.product_type_id] } { output: [ accounts_dm.account_id, accounts_dm.number, accounts_dm.product_type_id, accounts_dm.base_currency_code, accounts_dm.opening_date ], stream key: [ accounts_dm.product_type_id, accounts_dm.account_id ] }
└── MergeExecutor { output: [ accounts_dm.account_id, accounts_dm.number, accounts_dm.product_type_id, accounts_dm.base_currency_code, accounts_dm.opening_date ], stream key: [ accounts_dm.account_id ] }

Fragment 54344 (Actor 739601,739602)
StreamProject { exprs: [accounts_dm.account_id, accounts_dm.number, accounts_dm.product_type_id, accounts_dm.base_currency_code, accounts_dm.opening_date] } { output: [ accounts_dm.account_id, accounts_dm.number, accounts_dm.product_type_id, accounts_dm.base_currency_code, accounts_dm.opening_date ], stream key: [ accounts_dm.account_id ] }
└── StreamFilter { predicate: IsNull(accounts_dm.disabled_at) } { output: [ accounts_dm.account_id, accounts_dm.number, accounts_dm.product_type_id, accounts_dm.base_currency_code, accounts_dm.opening_date, accounts_dm.disabled_at ], stream key: [ accounts_dm.account_id ] }
    └── StreamTableScan { table: accounts_dm, columns: [account_id, number, product_type_id, base_currency_code, opening_date, disabled_at] } { output: [ accounts_dm.account_id, accounts_dm.number, accounts_dm.product_type_id, accounts_dm.base_currency_code, accounts_dm.opening_date, accounts_dm.disabled_at ], stream key: [ accounts_dm.account_id ] }
        ├── Upstream { output: [ account_id, number, product_type_id, base_currency_code, opening_date, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ account_id, number, product_type_id, base_currency_code, opening_date, disabled_at ], stream key: [] }

Fragment 54345 (Actor 739551,739552)
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 54346 (Actor 739603,739604)
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: ($expr2 > now), output_watermarks: [[$expr2]], output: [account_to_portfolios_dm.account_id, $expr2, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date], cleaned_by_watermark: true } { output: [ account_to_portfolios_dm.account_id, $expr2, 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 ] }
    ├── StreamProject { exprs: [account_to_portfolios_dm.account_id, AtTimeZone(Coalesce(account_to_portfolios_dm.effective_end_date, '9999-12-31':Date)::Timestamp, 'UTC':Varchar) as $expr2, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date] } { output: [ account_to_portfolios_dm.account_id, $expr2, 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: ($expr1 <= now), output: [account_to_portfolios_dm.account_id, account_to_portfolios_dm.effective_end_date, $expr1, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date], cleaned_by_watermark: true } { output: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.effective_end_date, $expr1, 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 ] }
    │       ├── StreamProject { exprs: [account_to_portfolios_dm.account_id, account_to_portfolios_dm.effective_end_date, AtTimeZone(account_to_portfolios_dm.effective_start_date::Timestamp, 'UTC':Varchar) as $expr1, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date] } { output: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.effective_end_date, $expr1, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ], stream key: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.portfolio_id, account_to_portfolios_dm.effective_start_date ] }
    │       │   └── StreamFilter { predicate: IsNull(account_to_portfolios_dm.disabled_at) } { output: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.effective_start_date, account_to_portfolios_dm.effective_end_date, account_to_portfolios_dm.portfolio_id, 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, effective_start_date, effective_end_date, portfolio_id, disabled_at] } { output: [ account_to_portfolios_dm.account_id, account_to_portfolios_dm.effective_start_date, account_to_portfolios_dm.effective_end_date, account_to_portfolios_dm.portfolio_id, 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, effective_start_date, effective_end_date, portfolio_id, disabled_at ], stream key: [] }
    │       │           └── BatchPlanNode { output: [ account_id, effective_start_date, effective_end_date, portfolio_id, disabled_at ], stream key: [] }
    │       └── MergeExecutor { output: [ now ], stream key: [] }
    └── MergeExecutor { output: [ now ], stream key: [] }

Fragment 54347 (Actor 739577)
StreamNow { output: [ now ], stream key: [] }

Fragment 54348 (Actor 739578)
StreamNow { output: [ now ], stream key: [] }