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

← cluster adib_rm objects party_cash_tiles_mv explain
Overview Objects Graph History
materialized view · adib_rm.party_cash_tiles_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
75 operators
Materialize · adib_rm.party_cash_tiles_mv
0% idle 2 actors
Project
2 actors
HashAgg Aggregation state — unbounded unless keyed or temporally filtered
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
2 actors
Project
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · LeftSemi · accounts_plain_mv.product_type_id = cash_product_type_ids_m…
2 actors
HashJoin · LeftSemi · accounts_plain_mv.product_type_id = cash_product_type_ids_m… 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 · cash_product_type_ids_mv
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · LeftOuter · accounts_plain_mv.account_id = investment_account_balance_s…
2 actors
HashJoin · LeftOuter · accounts_plain_mv.account_id = investment_account_balance_s… 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 · investment_account_balance_snapshot_journal_mv_next
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Union
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · IsNull(party_account_direct_mv_next.effective_end_date)
2 actors
Filter · IsNull(party_account_direct_mv_next.effective_end_date)
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · LeftOuter · ($expr1 <= now), output: [party_account_direct_mv_next.part…
2 actors
DynamicFilter · LeftOuter · ($expr1 <= now), output: [party_account_direct_mv_next.part… 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 · LeftOuter · ($expr1 <= now), output: [party_account_direct_mv_next.part…
2 actors
TemporalJoin · LeftOuter · ($expr1 <= now), output: [party_account_direct_mv_next.part…
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · currencies_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · party_account_direct_mv_next.account_id = accounts_plain_mv…
2 actors
HashJoin · Inner · party_account_direct_mv_next.account_id = accounts_plain_mv… 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 · accounts_plain_mv
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
Filter · party_account_direct_mv_next
0% idle 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 · ($expr2 > now), output_watermarks: [[$expr2]], output: [par…
2 actors
DynamicFilter · ($expr2 > now), output_watermarks: [[$expr2]], output: [par… 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 · ($expr2 > now), output_watermarks: [[$expr2]], output: [par…
2 actors
Filter · ($expr2 > now), output_watermarks: [[$expr2]], output: [par…
0% idle 2 actors
Merge
2 actors
Heat = the operator's output-buffer backpressure over the sampling window. Click a node to fold its subtree.
Materialize · adib_rm.party_cash_tiles_mv Materialize adib_rm.party_cash_tile… idle · 2 actors Project Project — · 2 actors HashAgg HashAgg idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 2 actors Project Project — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · LeftSemi · accounts_plain_mv.product_type_id = cash_product_type_ids_m… SyncLogStore LeftSemi · accounts_pla… — · 2 actors HashJoin · LeftSemi · accounts_plain_mv.product_type_id = cash_product_type_ids_m… HashJoin LeftSemi · accounts_pla… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · cash_product_type_ids_mv StreamScan cash_product_type_ids_mv idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · LeftOuter · accounts_plain_mv.account_id = investment_account_balance_s… SyncLogStore LeftOuter · accounts_pl… — · 2 actors HashJoin · LeftOuter · accounts_plain_mv.account_id = investment_account_balance_s… HashJoin LeftOuter · accounts_pl… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · investment_account_balance_snapshot_journal_mv_next StreamScan investment_account_bala… idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 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 · IsNull(party_account_direct_mv_next.effective_end_date) Project IsNull(party_account_di… — · 2 actors Filter · IsNull(party_account_direct_mv_next.effective_end_date) Filter IsNull(party_account_di… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · LeftOuter · ($expr1 <= now), output: [party_account_direct_mv_next.part… Project LeftOuter · ($expr1 <= … — · 2 actors DynamicFilter · LeftOuter · ($expr1 <= now), output: [party_account_direct_mv_next.part… DynamicFilter LeftOuter · ($expr1 <= … idle · 2 actors Merge Merge — · 2 actors Exchange Exchange 2/s · 0 actors Now Now 2/s · 1 actor Project · LeftOuter · ($expr1 <= now), output: [party_account_direct_mv_next.part… Project LeftOuter · ($expr1 <= … — · 2 actors TemporalJoin · LeftOuter · ($expr1 <= now), output: [party_account_direct_mv_next.part… TemporalJoin LeftOuter · ($expr1 <= … idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · currencies_dm StreamScan currencies_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore SyncLogStore — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · party_account_direct_mv_next.account_id = accounts_plain_mv… SyncLogStore Inner · party_account_d… — · 2 actors HashJoin · Inner · party_account_direct_mv_next.account_id = accounts_plain_mv… HashJoin Inner · party_account_d… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · accounts_plain_mv StreamScan accounts_plain_mv 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 Filter · party_account_direct_mv_next Filter party_account_direct_mv… idle · 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 · ($expr2 > now), output_watermarks: [[$expr2]], output: [par… Project ($expr2 > now), output_… — · 2 actors DynamicFilter · ($expr2 > now), output_watermarks: [[$expr2]], output: [par… DynamicFilter ($expr2 > now), output_… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange 2/s · 0 actors Now Now 2/s · 1 actor Project · ($expr2 > now), output_watermarks: [[$expr2]], output: [par… Project ($expr2 > now), output_… — · 2 actors Filter · ($expr2 > now), output_watermarks: [[$expr2]], output: [par… Filter ($expr2 > now), output_… idle · 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 25516 (Actor 125831,125830)
StreamMaterialize { columns: [party_id, cash_tiles], stream_key: [party_id], pk_columns: [party_id], pk_conflict: NoCheck }
├── output: [ party_account_direct_mv_next.party_id, jsonb_agg($expr4) ]
├── stream key: [ party_account_direct_mv_next.party_id ]
└── StreamProject { exprs: [party_account_direct_mv_next.party_id, jsonb_agg($expr4)] }
    ├── output: [ party_account_direct_mv_next.party_id, jsonb_agg($expr4) ]
    ├── stream key: [ party_account_direct_mv_next.party_id ]
    └── StreamHashAgg { group_key: [party_account_direct_mv_next.party_id], aggs: [jsonb_agg($expr4), count] }
        ├── output: [ party_account_direct_mv_next.party_id, jsonb_agg($expr4), count ]
        ├── stream key: [ party_account_direct_mv_next.party_id ]
        └── MergeExecutor
            ├── output:
            │   ┌── party_account_direct_mv_next.party_id
            │   ├── $expr4
            │   ├── party_account_direct_mv_next.party_involvements_dm.id
            │   ├── party_account_direct_mv_next.party_involvements_dm.entity_id
            │   ├── party_account_direct_mv_next.party_involvements_dm.customer_relationship_id
            │   ├── party_account_direct_mv_next.$src
            │   ├── party_account_direct_mv_next.account_id
            │   ├── accounts_plain_mv.base_currency_code
            │   ├── $src
            │   ├── accounts_plain_mv.account_id
            │   └── accounts_plain_mv.product_type_id
            └── stream key:
                ┌── party_account_direct_mv_next.party_involvements_dm.id
                ├── party_account_direct_mv_next.party_involvements_dm.entity_id
                ├── party_account_direct_mv_next.party_involvements_dm.customer_relationship_id
                ├── party_account_direct_mv_next.$src
                ├── party_account_direct_mv_next.account_id
                ├── accounts_plain_mv.base_currency_code
                ├── $src
                ├── accounts_plain_mv.account_id
                └── accounts_plain_mv.product_type_id

Fragment 25517 (Actor 125835,125834)
StreamProject { exprs: [party_account_direct_mv_next.party_id, JsonbBuildObject('accountId':Varchar, accounts_plain_mv.account_id, 'name':Varchar, accounts_plain_mv.name, 'number':Varchar, accounts_plain_mv.number, 'value':Varchar, JsonbBuildObject('amount':Varchar, investment_account_balance_snapshot_journal_mv_next.available_balance::Varchar, 'currency':Varchar, $expr3), 'baseCurrency':Varchar, $expr3) as $expr4, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code, $src, accounts_plain_mv.account_id, accounts_plain_mv.product_type_id] }
├── output: [ party_account_direct_mv_next.party_id, $expr4, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code, $src, accounts_plain_mv.account_id, accounts_plain_mv.product_type_id ]
├── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code, $src, accounts_plain_mv.account_id, accounts_plain_mv.product_type_id ]
└── StreamProject { exprs: [party_account_direct_mv_next.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.base_currency_code, currencies_dm.symbol, investment_account_balance_snapshot_journal_mv_next.available_balance, JsonbBuildObject('code':Varchar, accounts_plain_mv.base_currency_code, 'symbol':Varchar, currencies_dm.symbol) as $expr3, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code, $src, accounts_plain_mv.product_type_id] }
    ├── output: [ party_account_direct_mv_next.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.base_currency_code, currencies_dm.symbol, investment_account_balance_snapshot_journal_mv_next.available_balance, $expr3, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code, $src, accounts_plain_mv.product_type_id ]
    ├── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code, $src, accounts_plain_mv.account_id, accounts_plain_mv.product_type_id ]
    └── MergeExecutor
        ├── output: [ party_account_direct_mv_next.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.base_currency_code, currencies_dm.symbol, investment_account_balance_snapshot_journal_mv_next.available_balance, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code, $src, accounts_plain_mv.product_type_id ]
        └── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code, $src, accounts_plain_mv.account_id, accounts_plain_mv.product_type_id ]

Fragment 25518 (Actor 125833,125832)
StreamSyncLogStore
├── output: [ party_account_direct_mv_next.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.base_currency_code, currencies_dm.symbol, investment_account_balance_snapshot_journal_mv_next.available_balance, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code, $src, accounts_plain_mv.product_type_id ]
├── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code, $src, accounts_plain_mv.account_id, accounts_plain_mv.product_type_id ]
└── StreamHashJoin { type: LeftSemi, predicate: accounts_plain_mv.product_type_id = cash_product_type_ids_mv.product_type_id }
    ├── output: [ party_account_direct_mv_next.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.base_currency_code, currencies_dm.symbol, investment_account_balance_snapshot_journal_mv_next.available_balance, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code, $src, accounts_plain_mv.product_type_id ]
    ├── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code, $src, accounts_plain_mv.account_id, accounts_plain_mv.product_type_id ]
    ├── MergeExecutor
    │   ├── output: [ party_account_direct_mv_next.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, investment_account_balance_snapshot_journal_mv_next.available_balance, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code, $src, investment_account_balance_snapshot_journal_mv_next.account_id ]
    │   └── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code, $src, accounts_plain_mv.account_id ]
    └── MergeExecutor { output: [ cash_product_type_ids_mv.product_type_id, cash_product_type_ids_mv._row_id, cash_product_type_ids_mv.$src ], stream key: [ cash_product_type_ids_mv._row_id, cash_product_type_ids_mv.$src ] }

Fragment 25519 (Actor 125836,125837)
StreamSyncLogStore
├── output: [ party_account_direct_mv_next.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, investment_account_balance_snapshot_journal_mv_next.available_balance, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code, $src, investment_account_balance_snapshot_journal_mv_next.account_id ]
├── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code, $src, accounts_plain_mv.account_id ]
└── StreamHashJoin { type: LeftOuter, predicate: accounts_plain_mv.account_id = investment_account_balance_snapshot_journal_mv_next.account_id }
    ├── output: [ party_account_direct_mv_next.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, investment_account_balance_snapshot_journal_mv_next.available_balance, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code, $src, investment_account_balance_snapshot_journal_mv_next.account_id ]
    ├── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code, $src, accounts_plain_mv.account_id ]
    ├── MergeExecutor
    │   ├── output: [ party_account_direct_mv_next.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code, party_account_direct_mv_next.$src, $src ]
    │   └── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code, $src ]
    └── MergeExecutor { output: [ investment_account_balance_snapshot_journal_mv_next.account_id, investment_account_balance_snapshot_journal_mv_next.available_balance ], stream key: [ investment_account_balance_snapshot_journal_mv_next.account_id ] }

Fragment 25520 (Actor 125839,125838)
StreamUnion { all: true }
├── output: [ party_account_direct_mv_next.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code, party_account_direct_mv_next.$src, $src ]
├── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code, $src ]
├── MergeExecutor
│   ├── output: [ party_account_direct_mv_next.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code, party_account_direct_mv_next.$src, 0:Int32 ]
│   └── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code ]
└── MergeExecutor
    ├── output: [ party_account_direct_mv_next.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code, party_account_direct_mv_next.$src, 1:Int32 ]
    └── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code ]

Fragment 25521 (Actor 125783,125782)
StreamProject { exprs: [party_account_direct_mv_next.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code, party_account_direct_mv_next.$src, 0:Int32] }
├── output: [ party_account_direct_mv_next.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code, party_account_direct_mv_next.$src, 0:Int32 ]
├── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code ]
└── StreamDynamicFilter { predicate: ($expr2 > now), output_watermarks: [[$expr2]], output: [party_account_direct_mv_next.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, $expr2, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id], cleaned_by_watermark: true }
    ├── output: [ party_account_direct_mv_next.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, $expr2, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id ]
    ├── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code ]
    ├── StreamProject { exprs: [party_account_direct_mv_next.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, AtTimeZone(party_account_direct_mv_next.effective_end_date::Timestamp, 'UTC':Varchar) as $expr2, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id] }
    │   ├── output: [ party_account_direct_mv_next.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, $expr2, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id ]
    │   ├── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code ]
    │   └── StreamFilter { predicate: IsNotTrue(IsNull(party_account_direct_mv_next.effective_end_date)) }
    │       ├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id ]
    │       ├── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code ]
    │       └── MergeExecutor
    │           ├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id ]
    │           └── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code ]
    └── MergeExecutor { output: [ now ], stream key: [] }

Fragment 25522 (Actor 125778,125779)
StreamProject { exprs: [party_account_direct_mv_next.party_id, party_account_direct_mv_next.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id] }
├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id ]
├── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code ]
└── StreamDynamicFilter { predicate: ($expr1 <= now), output: [party_account_direct_mv_next.party_id, party_account_direct_mv_next.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, $expr1, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id], cleaned_by_watermark: true }
    ├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, $expr1, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id ]
    ├── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code ]
    ├── StreamProject { exprs: [party_account_direct_mv_next.party_id, party_account_direct_mv_next.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, AtTimeZone(party_account_direct_mv_next.effective_start_date::Timestamp, 'UTC':Varchar) as $expr1, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id] }
    │   ├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, $expr1, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id ]
    │   ├── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code ]
    │   └── StreamTemporalJoin { type: LeftOuter, append_only: false, predicate: accounts_plain_mv.base_currency_code = currencies_dm.code, nested_loop: false }
    │       ├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.effective_start_date, party_account_direct_mv_next.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, currencies_dm.code ]
    │       ├── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code ]
    │       ├── MergeExecutor
    │       │   ├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.effective_start_date, party_account_direct_mv_next.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id ]
    │       │   └── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id ]
    │       └── MergeExecutor { output: [ currencies_dm.code, currencies_dm.symbol ], stream key: [ currencies_dm.code ] }
    └── MergeExecutor { output: [ now ], stream key: [] }

Fragment 25523 (Actor 125843,125842)
StreamSyncLogStore
├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.effective_start_date, party_account_direct_mv_next.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id ]
├── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id ]
└── MergeExecutor
    ├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.effective_start_date, party_account_direct_mv_next.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id ]
    └── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id ]

Fragment 25524 (Actor 125841,125840)
StreamSyncLogStore
├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.effective_start_date, party_account_direct_mv_next.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id ]
├── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id ]
└── StreamHashJoin { type: Inner, predicate: party_account_direct_mv_next.account_id = accounts_plain_mv.account_id }
    ├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.effective_start_date, party_account_direct_mv_next.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id ]
    ├── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id ]
    ├── MergeExecutor { 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.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src ], stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src ] }
    └── MergeExecutor { output: [ accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code ], stream key: [ accounts_plain_mv.account_id ] }

Fragment 25525 (Actor 125845,125844)
StreamProject { exprs: [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.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$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.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src ]
├── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src ]
└── StreamFilter { predicate: (IsNotTrue(IsNull(party_account_direct_mv_next.effective_end_date)) OR IsNull(party_account_direct_mv_next.effective_end_date)) AND (party_account_direct_mv_next.type = 'all':Varchar) }
    ├── 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.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.type ]
    ├── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_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, party_involvements_dm.id, party_involvements_dm.entity_id, party_involvements_dm.customer_relationship_id, $src, type] }
        ├── 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.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.type ]
        ├── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src ]
        ├── Upstream { output: [ party_id, account_id, effective_start_date, effective_end_date, party_involvements_dm.id, party_involvements_dm.entity_id, party_involvements_dm.customer_relationship_id, $src, type ], stream key: [] }
        └── BatchPlanNode { output: [ party_id, account_id, effective_start_date, effective_end_date, party_involvements_dm.id, party_involvements_dm.entity_id, party_involvements_dm.customer_relationship_id, $src, type ], stream key: [] }

Fragment 25526 (Actor 125846,125847)
StreamTableScan { table: accounts_plain_mv, columns: [account_id, name, number, product_type_id, base_currency_code] } { output: [ accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code ], stream key: [ accounts_plain_mv.account_id ] }
├── Upstream { output: [ account_id, name, number, product_type_id, base_currency_code ], stream key: [] }
└── BatchPlanNode { output: [ account_id, name, number, product_type_id, base_currency_code ], stream key: [] }

Fragment 25527 (Actor 125785,125784)
StreamTableScan { table: currencies_dm, columns: [code, symbol] } { output: [ currencies_dm.code, currencies_dm.symbol ], stream key: [ currencies_dm.code ] }
├── Upstream { output: [ code, symbol ], stream key: [] }
└── BatchPlanNode { output: [ code, symbol ], stream key: [] }

Fragment 25528 (Actor 125848)
StreamNow { output: [ now ], stream key: [] }

Fragment 25529 (Actor 125849)
StreamNow { output: [ now ], stream key: [] }

Fragment 25530 (Actor 125781,125780)
StreamProject { exprs: [party_account_direct_mv_next.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code, party_account_direct_mv_next.$src, 1:Int32] }
├── output: [ party_account_direct_mv_next.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code, party_account_direct_mv_next.$src, 1:Int32 ]
├── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code ]
└── StreamFilter { predicate: IsNull(party_account_direct_mv_next.effective_end_date) }
    ├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id ]
    ├── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code ]
    └── MergeExecutor
        ├── output: [ party_account_direct_mv_next.party_id, party_account_direct_mv_next.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id ]
        └── stream key: [ party_account_direct_mv_next.party_involvements_dm.id, party_account_direct_mv_next.party_involvements_dm.entity_id, party_account_direct_mv_next.party_involvements_dm.customer_relationship_id, party_account_direct_mv_next.$src, party_account_direct_mv_next.account_id, accounts_plain_mv.base_currency_code ]

Fragment 25531 (Actor 125850,125851)
StreamTableScan { table: investment_account_balance_snapshot_journal_mv_next, columns: [account_id, available_balance] } { output: [ investment_account_balance_snapshot_journal_mv_next.account_id, investment_account_balance_snapshot_journal_mv_next.available_balance ], stream key: [ investment_account_balance_snapshot_journal_mv_next.account_id ] }
├── Upstream { output: [ account_id, available_balance ], stream key: [] }
└── BatchPlanNode { output: [ account_id, available_balance ], stream key: [] }

Fragment 25532 (Actor 125853,125852)
StreamTableScan { table: cash_product_type_ids_mv, columns: [product_type_id, _row_id, $src] } { output: [ cash_product_type_ids_mv.product_type_id, cash_product_type_ids_mv._row_id, cash_product_type_ids_mv.$src ], stream key: [ cash_product_type_ids_mv._row_id, cash_product_type_ids_mv.$src ] }
├── Upstream { output: [ product_type_id, _row_id, $src ], stream key: [] }
└── BatchPlanNode { output: [ product_type_id, _row_id, $src ], stream key: [] }