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

← cluster opportunity objects account_balance_delta_mv explain
Overview Objects Graph History
materialized view · opportunity.account_balance_delta_mv profiled over 5s
seconds (1–30)
Stateful hash join (4 state tables) — consider a temporal join for dimension lookupsDynamic filter — verify it pairs with a temporal condition to clean state
62 operators
Materialize · opportunity.account_balance_delta_mv
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · $expr1 = position_summary_mv_next.account_group_id
2 actors
HashJoin · Inner · $expr1 = position_summary_mv_next.account_group_id Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
LocalityProvider
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · position_summary_mv_next.account_group_id = position_summar…
2 actors
HashJoin · Inner · position_summary_mv_next.account_group_id = position_summar… 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
Filter · position_summary_mv_next
0% idle 2 actors
StreamScan · position_summary_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
Project · position_summary_mv_next
2 actors
DynamicFilter · position_summary_mv_next 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
Project
1 actor
Now
0% 2/s 1 actor
Project · position_summary_mv_next
2 actors
Filter · position_summary_mv_next
2% idle 2 actors
StreamScan · position_summary_mv_next
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
Project
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · open_accounts_mv.product_type_id = product_types_dm.product…
2 actors
HashJoin · Inner · open_accounts_mv.product_type_id = product_types_dm.product… 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 · 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
LocalityProvider
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · open_accounts_mv
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Heat = the operator's output-buffer backpressure over the sampling window. Click a node to fold its subtree.
Materialize · opportunity.account_balance_delta_mv Materialize opportunity.account_bal… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · $expr1 = position_summary_mv_next.account_group_id SyncLogStore Inner · $expr1 = positi… — · 2 actors HashJoin · Inner · $expr1 = position_summary_mv_next.account_group_id HashJoin Inner · $expr1 = positi… 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 · position_summary_mv_next.account_group_id = position_summar… SyncLogStore Inner · position_summar… — · 2 actors HashJoin · Inner · position_summary_mv_next.account_group_id = position_summar… HashJoin Inner · position_summar… 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 Filter · position_summary_mv_next Filter position_summary_mv_next idle · 2 actors StreamScan · position_summary_mv_next StreamScan position_summary_mv_next idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors LocalityProvider LocalityProvider idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · position_summary_mv_next Project position_summary_mv_next — · 2 actors DynamicFilter · position_summary_mv_next DynamicFilter position_summary_mv_next idle · 2 actors Merge Merge — · 2 actors Exchange Exchange 2/s · 0 actors Project Project — · 1 actor Now Now 2/s · 1 actor Project · position_summary_mv_next Project position_summary_mv_next — · 2 actors Filter · position_summary_mv_next Filter position_summary_mv_next idle · 2 actors StreamScan · position_summary_mv_next StreamScan position_summary_mv_next idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors LocalityProvider LocalityProvider idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · open_accounts_mv.product_type_id = product_types_dm.product… SyncLogStore Inner · open_accounts_m… — · 2 actors HashJoin · Inner · open_accounts_mv.product_type_id = product_types_dm.product… HashJoin Inner · open_accounts_m… idle · 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 LocalityProvider LocalityProvider idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · open_accounts_mv StreamScan open_accounts_mv idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors
Streaming operator plan from EXPLAIN ANALYZE. Node heat = backpressure. Drag to pan, scroll to zoom.
Fragments (DESCRIBE FRAGMENTS) — click to expand
Fragment 64240 (Actor 747305,747304)
StreamMaterialize { columns: [account_id, dim_balance_date, currency_code, market_value, prev_market_value, $expr1(hidden), open_accounts_mv.product_type_id(hidden), position_summary_mv_next.account_group_id(hidden), position_summary_mv_next.position_type(hidden), $expr4(hidden), position_summary_mv_next.source_entity_type(hidden), position_summary_mv_next.source_entity_type#1(hidden)], stream_key: [$expr1, open_accounts_mv.product_type_id, account_id, position_summary_mv_next.position_type, currency_code, $expr4, position_summary_mv_next.source_entity_type, dim_balance_date, position_summary_mv_next.source_entity_type#1], pk_columns: [$expr1, open_accounts_mv.product_type_id, account_id, position_summary_mv_next.position_type, currency_code, $expr4, position_summary_mv_next.source_entity_type, dim_balance_date, position_summary_mv_next.source_entity_type#1], pk_conflict: NoCheck }
├── output: [ open_accounts_mv.account_id, position_summary_mv_next.dim_balance_date, position_summary_mv_next.currency_code, position_summary_mv_next.market_value, position_summary_mv_next.market_value, $expr1, open_accounts_mv.product_type_id, position_summary_mv_next.account_group_id, position_summary_mv_next.position_type, $expr4, position_summary_mv_next.source_entity_type, position_summary_mv_next.source_entity_type ]
├── stream key: [ $expr1, open_accounts_mv.product_type_id, open_accounts_mv.account_id, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, $expr4, position_summary_mv_next.source_entity_type, position_summary_mv_next.dim_balance_date, position_summary_mv_next.source_entity_type ]
└── MergeExecutor { output: [ open_accounts_mv.account_id, position_summary_mv_next.dim_balance_date, position_summary_mv_next.currency_code, position_summary_mv_next.market_value, position_summary_mv_next.market_value, $expr1, open_accounts_mv.product_type_id, position_summary_mv_next.account_group_id, position_summary_mv_next.position_type, $expr4, position_summary_mv_next.source_entity_type, position_summary_mv_next.source_entity_type ], stream key: [ $expr1, open_accounts_mv.product_type_id, open_accounts_mv.account_id, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, $expr4, position_summary_mv_next.source_entity_type, position_summary_mv_next.dim_balance_date, position_summary_mv_next.source_entity_type ] }

Fragment 64241 (Actor 747302,747303)
StreamSyncLogStore { output: [ open_accounts_mv.account_id, position_summary_mv_next.dim_balance_date, position_summary_mv_next.currency_code, position_summary_mv_next.market_value, position_summary_mv_next.market_value, $expr1, open_accounts_mv.product_type_id, position_summary_mv_next.account_group_id, position_summary_mv_next.position_type, $expr4, position_summary_mv_next.source_entity_type, position_summary_mv_next.source_entity_type ], stream key: [ $expr1, open_accounts_mv.product_type_id, open_accounts_mv.account_id, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, $expr4, position_summary_mv_next.source_entity_type, position_summary_mv_next.dim_balance_date, position_summary_mv_next.source_entity_type ] }
└── StreamHashJoin { type: Inner, predicate: $expr1 = position_summary_mv_next.account_group_id } { output: [ open_accounts_mv.account_id, position_summary_mv_next.dim_balance_date, position_summary_mv_next.currency_code, position_summary_mv_next.market_value, position_summary_mv_next.market_value, $expr1, open_accounts_mv.product_type_id, position_summary_mv_next.account_group_id, position_summary_mv_next.position_type, $expr4, position_summary_mv_next.source_entity_type, position_summary_mv_next.source_entity_type ], stream key: [ $expr1, open_accounts_mv.product_type_id, open_accounts_mv.account_id, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, $expr4, position_summary_mv_next.source_entity_type, position_summary_mv_next.dim_balance_date, position_summary_mv_next.source_entity_type ] }
    ├── MergeExecutor { output: [ open_accounts_mv.account_id, $expr1, open_accounts_mv.product_type_id ], stream key: [ $expr1, open_accounts_mv.product_type_id, open_accounts_mv.account_id ] }
    └── MergeExecutor { output: [ position_summary_mv_next.account_group_id, position_summary_mv_next.dim_balance_date, position_summary_mv_next.currency_code, position_summary_mv_next.market_value, position_summary_mv_next.market_value, position_summary_mv_next.position_type, $expr4, position_summary_mv_next.source_entity_type, position_summary_mv_next.account_group_id, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, position_summary_mv_next.dim_balance_date, position_summary_mv_next.source_entity_type ], stream key: [ position_summary_mv_next.account_group_id, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, $expr4, position_summary_mv_next.source_entity_type, position_summary_mv_next.dim_balance_date, position_summary_mv_next.source_entity_type ] }

Fragment 64242 (Actor 747307,747306)
StreamLocalityProvider { locality_columns: [$expr1] } { output: [ open_accounts_mv.account_id, $expr1, open_accounts_mv.product_type_id ], stream key: [ $expr1, open_accounts_mv.product_type_id, open_accounts_mv.account_id ] }
└── MergeExecutor { output: [ open_accounts_mv.account_id, $expr1, open_accounts_mv.product_type_id ], stream key: [ open_accounts_mv.product_type_id, open_accounts_mv.account_id ] }

Fragment 64243 (Actor 747308,747309)
StreamProject { exprs: [open_accounts_mv.account_id, ConcatOp('account_group_':Varchar, Md5(ConcatOp(open_accounts_mv.account_id, 'all':Varchar)::Bytea)) as $expr1, open_accounts_mv.product_type_id] } { output: [ open_accounts_mv.account_id, $expr1, open_accounts_mv.product_type_id ], stream key: [ open_accounts_mv.product_type_id, open_accounts_mv.account_id ] }
└── MergeExecutor { output: [ open_accounts_mv.account_id, open_accounts_mv.product_type_id, product_types_dm.product_type_id ], stream key: [ open_accounts_mv.product_type_id, open_accounts_mv.account_id ] }

Fragment 64244 (Actor 747311,747310)
StreamSyncLogStore { output: [ open_accounts_mv.account_id, open_accounts_mv.product_type_id, product_types_dm.product_type_id ], stream key: [ open_accounts_mv.product_type_id, open_accounts_mv.account_id ] }
└── StreamHashJoin { type: Inner, predicate: open_accounts_mv.product_type_id = product_types_dm.product_type_id } { output: [ open_accounts_mv.account_id, open_accounts_mv.product_type_id, product_types_dm.product_type_id ], stream key: [ open_accounts_mv.product_type_id, open_accounts_mv.account_id ] }
    ├── MergeExecutor { output: [ open_accounts_mv.account_id, open_accounts_mv.product_type_id ], stream key: [ open_accounts_mv.product_type_id, open_accounts_mv.account_id ] }
    └── MergeExecutor { output: [ product_types_dm.product_type_id ], stream key: [ product_types_dm.product_type_id ] }

Fragment 64245 (Actor 747313,747312)
StreamLocalityProvider { locality_columns: [open_accounts_mv.product_type_id] } { output: [ open_accounts_mv.account_id, open_accounts_mv.product_type_id ], stream key: [ open_accounts_mv.product_type_id, open_accounts_mv.account_id ] }
└── MergeExecutor { output: [ open_accounts_mv.account_id, open_accounts_mv.product_type_id ], stream key: [ open_accounts_mv.account_id ] }

Fragment 64246 (Actor 745967,745966)
StreamTableScan { table: open_accounts_mv, columns: [account_id, product_type_id] } { output: [ open_accounts_mv.account_id, open_accounts_mv.product_type_id ], stream key: [ open_accounts_mv.account_id ] }
├── Upstream { output: [ account_id, product_type_id ], stream key: [] }
└── BatchPlanNode { output: [ account_id, product_type_id ], stream key: [] }

Fragment 64247 (Actor 747327,747328)
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 64248 (Actor 745823,745822)
StreamLocalityProvider { locality_columns: [position_summary_mv_next.account_group_id] } { output: [ position_summary_mv_next.account_group_id, position_summary_mv_next.dim_balance_date, position_summary_mv_next.currency_code, position_summary_mv_next.market_value, position_summary_mv_next.market_value, position_summary_mv_next.position_type, $expr4, position_summary_mv_next.source_entity_type, position_summary_mv_next.account_group_id, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, position_summary_mv_next.dim_balance_date, position_summary_mv_next.source_entity_type ], stream key: [ position_summary_mv_next.account_group_id, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, $expr4, position_summary_mv_next.source_entity_type, position_summary_mv_next.dim_balance_date, position_summary_mv_next.source_entity_type ] }
└── MergeExecutor { output: [ position_summary_mv_next.account_group_id, position_summary_mv_next.dim_balance_date, position_summary_mv_next.currency_code, position_summary_mv_next.market_value, position_summary_mv_next.market_value, position_summary_mv_next.position_type, $expr4, position_summary_mv_next.source_entity_type, position_summary_mv_next.account_group_id, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, position_summary_mv_next.dim_balance_date, position_summary_mv_next.source_entity_type ], stream key: [ position_summary_mv_next.account_group_id, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, $expr4, position_summary_mv_next.source_entity_type, position_summary_mv_next.dim_balance_date, position_summary_mv_next.source_entity_type ] }

Fragment 64249 (Actor 747315,747314)
StreamSyncLogStore { output: [ position_summary_mv_next.account_group_id, position_summary_mv_next.dim_balance_date, position_summary_mv_next.currency_code, position_summary_mv_next.market_value, position_summary_mv_next.market_value, position_summary_mv_next.position_type, $expr4, position_summary_mv_next.source_entity_type, position_summary_mv_next.account_group_id, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, position_summary_mv_next.dim_balance_date, position_summary_mv_next.source_entity_type ], stream key: [ position_summary_mv_next.account_group_id, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, $expr4, position_summary_mv_next.source_entity_type, position_summary_mv_next.dim_balance_date, position_summary_mv_next.source_entity_type ] }
└── StreamHashJoin { type: Inner, predicate: position_summary_mv_next.account_group_id = position_summary_mv_next.account_group_id AND position_summary_mv_next.position_type = position_summary_mv_next.position_type AND position_summary_mv_next.currency_code = position_summary_mv_next.currency_code AND $expr4 = position_summary_mv_next.dim_balance_date }
    ├── output: [ position_summary_mv_next.account_group_id, position_summary_mv_next.dim_balance_date, position_summary_mv_next.currency_code, position_summary_mv_next.market_value, position_summary_mv_next.market_value, position_summary_mv_next.position_type, $expr4, position_summary_mv_next.source_entity_type, position_summary_mv_next.account_group_id, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, position_summary_mv_next.dim_balance_date, position_summary_mv_next.source_entity_type ]
    ├── stream key: [ position_summary_mv_next.account_group_id, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, $expr4, position_summary_mv_next.source_entity_type, position_summary_mv_next.dim_balance_date, position_summary_mv_next.source_entity_type ]
    ├── MergeExecutor { output: [ position_summary_mv_next.account_group_id, position_summary_mv_next.dim_balance_date, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, position_summary_mv_next.market_value, $expr4, position_summary_mv_next.source_entity_type ], stream key: [ position_summary_mv_next.account_group_id, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, $expr4, position_summary_mv_next.source_entity_type, position_summary_mv_next.dim_balance_date ] }
    └── MergeExecutor { output: [ position_summary_mv_next.account_group_id, position_summary_mv_next.dim_balance_date, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, position_summary_mv_next.market_value, position_summary_mv_next.source_entity_type ], stream key: [ position_summary_mv_next.account_group_id, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, position_summary_mv_next.dim_balance_date, position_summary_mv_next.source_entity_type ] }

Fragment 64250 (Actor 747317,747316)
StreamLocalityProvider { locality_columns: [position_summary_mv_next.account_group_id, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, $expr4] } { output: [ position_summary_mv_next.account_group_id, position_summary_mv_next.dim_balance_date, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, position_summary_mv_next.market_value, $expr4, position_summary_mv_next.source_entity_type ], stream key: [ position_summary_mv_next.account_group_id, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, $expr4, position_summary_mv_next.source_entity_type, position_summary_mv_next.dim_balance_date ] }
└── MergeExecutor { output: [ position_summary_mv_next.account_group_id, position_summary_mv_next.dim_balance_date, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, position_summary_mv_next.market_value, $expr4, position_summary_mv_next.source_entity_type ], stream key: [ position_summary_mv_next.account_group_id, position_summary_mv_next.source_entity_type, position_summary_mv_next.dim_balance_date, position_summary_mv_next.currency_code, position_summary_mv_next.position_type ] }

Fragment 64251 (Actor 747322,747321)
StreamProject { exprs: [position_summary_mv_next.account_group_id, position_summary_mv_next.dim_balance_date, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, position_summary_mv_next.market_value, (position_summary_mv_next.dim_balance_date - 1:Int32) as $expr4, position_summary_mv_next.source_entity_type] } { output: [ position_summary_mv_next.account_group_id, position_summary_mv_next.dim_balance_date, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, position_summary_mv_next.market_value, $expr4, position_summary_mv_next.source_entity_type ], stream key: [ position_summary_mv_next.account_group_id, position_summary_mv_next.source_entity_type, position_summary_mv_next.dim_balance_date, position_summary_mv_next.currency_code, position_summary_mv_next.position_type ] }
└── StreamDynamicFilter { predicate: ($expr2 >= $expr3), output_watermarks: [[$expr2]], output: [position_summary_mv_next.account_group_id, position_summary_mv_next.dim_balance_date, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, position_summary_mv_next.market_value, $expr2, position_summary_mv_next.source_entity_type], cleaned_by_watermark: true } { output: [ position_summary_mv_next.account_group_id, position_summary_mv_next.dim_balance_date, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, position_summary_mv_next.market_value, $expr2, position_summary_mv_next.source_entity_type ], stream key: [ position_summary_mv_next.account_group_id, position_summary_mv_next.source_entity_type, position_summary_mv_next.dim_balance_date, position_summary_mv_next.currency_code, position_summary_mv_next.position_type ] }
    ├── StreamProject { exprs: [position_summary_mv_next.account_group_id, position_summary_mv_next.dim_balance_date, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, position_summary_mv_next.market_value, AtTimeZone(position_summary_mv_next.dim_balance_date::Timestamp, 'UTC':Varchar) as $expr2, position_summary_mv_next.source_entity_type] } { output: [ position_summary_mv_next.account_group_id, position_summary_mv_next.dim_balance_date, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, position_summary_mv_next.market_value, $expr2, position_summary_mv_next.source_entity_type ], stream key: [ position_summary_mv_next.account_group_id, position_summary_mv_next.source_entity_type, position_summary_mv_next.dim_balance_date, position_summary_mv_next.currency_code, position_summary_mv_next.position_type ] }
    │   └── StreamFilter { predicate: (position_summary_mv_next.source_entity_type = 'account':Varchar) AND (position_summary_mv_next.position_type = 'POSITION':Varchar) } { output: [ position_summary_mv_next.account_group_id, position_summary_mv_next.dim_balance_date, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, position_summary_mv_next.market_value, position_summary_mv_next.source_entity_type ], stream key: [ position_summary_mv_next.account_group_id, position_summary_mv_next.source_entity_type, position_summary_mv_next.dim_balance_date, position_summary_mv_next.currency_code, position_summary_mv_next.position_type ] }
    │       └── StreamTableScan { table: position_summary_mv_next, columns: [account_group_id, dim_balance_date, position_type, currency_code, market_value, source_entity_type] } { output: [ position_summary_mv_next.account_group_id, position_summary_mv_next.dim_balance_date, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, position_summary_mv_next.market_value, position_summary_mv_next.source_entity_type ], stream key: [ position_summary_mv_next.account_group_id, position_summary_mv_next.source_entity_type, position_summary_mv_next.dim_balance_date, position_summary_mv_next.currency_code, position_summary_mv_next.position_type ] }
    │           ├── Upstream { output: [ account_group_id, dim_balance_date, position_type, currency_code, market_value, source_entity_type ], stream key: [] }
    │           └── BatchPlanNode { output: [ account_group_id, dim_balance_date, position_type, currency_code, market_value, source_entity_type ], stream key: [] }
    └── MergeExecutor { output: [ $expr3 ], stream key: [] }

Fragment 64252 (Actor 747318)
StreamProject { exprs: [SubtractWithTimeZone(now, '30 days':Interval, 'UTC':Varchar) as $expr3], output_watermarks: [[$expr3]] } { output: [ $expr3 ], stream key: [] }
└── StreamNow { output: [ now ], stream key: [] }

Fragment 64253 (Actor 747319,747320)
StreamLocalityProvider { locality_columns: [position_summary_mv_next.account_group_id, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, position_summary_mv_next.dim_balance_date] } { output: [ position_summary_mv_next.account_group_id, position_summary_mv_next.dim_balance_date, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, position_summary_mv_next.market_value, position_summary_mv_next.source_entity_type ], stream key: [ position_summary_mv_next.account_group_id, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, position_summary_mv_next.dim_balance_date, position_summary_mv_next.source_entity_type ] }
└── MergeExecutor { output: [ position_summary_mv_next.account_group_id, position_summary_mv_next.dim_balance_date, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, position_summary_mv_next.market_value, position_summary_mv_next.source_entity_type ], stream key: [ position_summary_mv_next.account_group_id, position_summary_mv_next.source_entity_type, position_summary_mv_next.dim_balance_date, position_summary_mv_next.currency_code, position_summary_mv_next.position_type ] }

Fragment 64254 (Actor 747323,747324)
StreamFilter { predicate: (position_summary_mv_next.market_value > 0:Decimal) AND (position_summary_mv_next.position_type = 'POSITION':Varchar) } { output: [ position_summary_mv_next.account_group_id, position_summary_mv_next.dim_balance_date, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, position_summary_mv_next.market_value, position_summary_mv_next.source_entity_type ], stream key: [ position_summary_mv_next.account_group_id, position_summary_mv_next.source_entity_type, position_summary_mv_next.dim_balance_date, position_summary_mv_next.currency_code, position_summary_mv_next.position_type ] }
└── StreamTableScan { table: position_summary_mv_next, columns: [account_group_id, dim_balance_date, position_type, currency_code, market_value, source_entity_type] } { output: [ position_summary_mv_next.account_group_id, position_summary_mv_next.dim_balance_date, position_summary_mv_next.position_type, position_summary_mv_next.currency_code, position_summary_mv_next.market_value, position_summary_mv_next.source_entity_type ], stream key: [ position_summary_mv_next.account_group_id, position_summary_mv_next.source_entity_type, position_summary_mv_next.dim_balance_date, position_summary_mv_next.currency_code, position_summary_mv_next.position_type ] }
    ├── Upstream { output: [ account_group_id, dim_balance_date, position_type, currency_code, market_value, source_entity_type ], stream key: [] }
    └── BatchPlanNode { output: [ account_group_id, dim_balance_date, position_type, currency_code, market_value, source_entity_type ], stream key: [] }