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

← cluster opportunity objects portfolio_extreme_holding_breaches_mv explain
Overview Objects Graph History
materialized view · opportunity.portfolio_extreme_holding_breaches_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
80 operators
Materialize · opportunity.portfolio_extreme_holding_breaches_mv
0% idle 2 actors
Project
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · portfolio_to_account_groups_mv.account_group_id = position_…
2 actors
HashJoin · Inner · portfolio_to_account_groups_mv.account_group_id = position_… 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 · opportunity_conditions_mv_next.activity_name = 'GET_PORTFOL…
2 actors
HashJoin · Inner · opportunity_conditions_mv_next.activity_name = 'GET_PORTFOL… 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
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · assets_dm_next.id = position_by_asset_mv_next.asset_id
2 actors
HashJoin · Inner · assets_dm_next.id = position_by_asset_mv_next.asset_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
Filter · (sum(position_by_asset_mv_next.market_value) > 0:Decimal) A…
0% idle 2 actors
Project · (sum(position_by_asset_mv_next.market_value) > 0:Decimal) A…
2 actors
HashAgg · (sum(position_by_asset_mv_next.market_value) > 0:Decimal) A… Aggregation state — unbounded unless keyed or temporally filtered
0% idle 2 actors
LocalityProvider · (sum(position_by_asset_mv_next.market_value) > 0:Decimal) A…
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · position_by_asset_mv_next
2 actors
DynamicFilter · position_by_asset_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_by_asset_mv_next
2 actors
Filter · position_by_asset_mv_next
2% idle 2 actors
StreamScan · position_by_asset_mv_next
2% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · assets_dm_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
StreamScan · opportunity_conditions_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 · portfolios_dm.portfolio_id = portfolio_to_account_groups_mv…
2 actors
HashJoin · Inner · portfolios_dm.portfolio_id = portfolio_to_account_groups_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
LocalityProvider
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · portfolio_to_account_groups_mv
2 actors
Filter · portfolio_to_account_groups_mv
0% idle 2 actors
StreamScan · portfolio_to_account_groups_mv
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · portfolios_dm
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.portfolio_extreme_holding_breaches_mv Materialize opportunity.portfolio_e… idle · 2 actors Project Project — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · portfolio_to_account_groups_mv.account_group_id = position_… SyncLogStore Inner · portfolio_to_ac… — · 2 actors HashJoin · Inner · portfolio_to_account_groups_mv.account_group_id = position_… HashJoin Inner · portfolio_to_ac… 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 · opportunity_conditions_mv_next.activity_name = 'GET_PORTFOL… SyncLogStore Inner · opportunity_con… — · 2 actors HashJoin · Inner · opportunity_conditions_mv_next.activity_name = 'GET_PORTFOL… HashJoin Inner · opportunity_con… 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 Project — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · assets_dm_next.id = position_by_asset_mv_next.asset_id SyncLogStore Inner · assets_dm_next.… — · 2 actors HashJoin · Inner · assets_dm_next.id = position_by_asset_mv_next.asset_id HashJoin Inner · assets_dm_next.… 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 · (sum(position_by_asset_mv_next.market_value) > 0:Decimal) A… Filter (sum(position_by_asset_… idle · 2 actors Project · (sum(position_by_asset_mv_next.market_value) > 0:Decimal) A… Project (sum(position_by_asset_… — · 2 actors HashAgg · (sum(position_by_asset_mv_next.market_value) > 0:Decimal) A… HashAgg (sum(position_by_asset_… idle · 2 actors LocalityProvider · (sum(position_by_asset_mv_next.market_value) > 0:Decimal) A… LocalityProvider (sum(position_by_asset_… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · position_by_asset_mv_next Project position_by_asset_mv_ne… — · 2 actors DynamicFilter · position_by_asset_mv_next DynamicFilter position_by_asset_mv_ne… 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_by_asset_mv_next Project position_by_asset_mv_ne… — · 2 actors Filter · position_by_asset_mv_next Filter position_by_asset_mv_ne… idle · 2 actors StreamScan · position_by_asset_mv_next StreamScan position_by_asset_mv_ne… idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · assets_dm_next StreamScan assets_dm_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 StreamScan · opportunity_conditions_mv_next StreamScan opportunity_conditions_… 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 · portfolios_dm.portfolio_id = portfolio_to_account_groups_mv… SyncLogStore Inner · portfolios_dm.p… — · 2 actors HashJoin · Inner · portfolios_dm.portfolio_id = portfolio_to_account_groups_mv… HashJoin Inner · portfolios_dm.p… 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 · portfolio_to_account_groups_mv Project portfolio_to_account_gr… — · 2 actors Filter · portfolio_to_account_groups_mv Filter portfolio_to_account_gr… idle · 2 actors StreamScan · portfolio_to_account_groups_mv StreamScan portfolio_to_account_gr… idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · portfolios_dm StreamScan portfolios_dm 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 64430 (Actor 747623,747624)
StreamMaterialize { columns: [opportunity_id, resource_id, fact_date, asset_id, asset_type, asset_name, asset_market_value, asset_currency_code, asset_percentage_value, asset_weight, portfolio_currency_code, portfolio_market_value, portfolio_to_account_groups_mv.account_group_id(hidden), portfolios_dm.portfolio_id(hidden), portfolio_to_account_groups_mv.$src(hidden), opportunity_conditions_mv_next.activity_name(hidden), opportunity_conditions_mv_next._rw_projected_row_id(hidden), opportunity_conditions_mv_next._rw_projected_row_id#1(hidden), assets_dm_next.id(hidden)], stream_key: [portfolio_to_account_groups_mv.account_group_id, portfolios_dm.portfolio_id, portfolio_to_account_groups_mv.$src, opportunity_conditions_mv_next.activity_name, opportunity_id, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1, assets_dm_next.id, fact_date, asset_currency_code], pk_columns: [portfolio_to_account_groups_mv.account_group_id, portfolios_dm.portfolio_id, portfolio_to_account_groups_mv.$src, opportunity_conditions_mv_next.activity_name, opportunity_id, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1, assets_dm_next.id, fact_date, asset_currency_code], pk_conflict: NoCheck }
├── output: [ opportunity_conditions_mv_next.opportunity_id, portfolio_to_account_groups_mv.portfolio_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, assets_dm_next.type, assets_dm_next.name_en, sum(position_by_asset_mv_next.market_value), position_by_asset_mv_next.currency_code, $expr3, sum(position_by_asset_mv_next.weight), portfolios_dm.base_currency_code, $expr4, portfolio_to_account_groups_mv.account_group_id, portfolios_dm.portfolio_id, portfolio_to_account_groups_mv.$src, opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1, assets_dm_next.id ]
├── stream key: [ portfolio_to_account_groups_mv.account_group_id, portfolios_dm.portfolio_id, portfolio_to_account_groups_mv.$src, opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1, assets_dm_next.id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.currency_code ]
└── StreamProject { exprs: [opportunity_conditions_mv_next.opportunity_id, portfolio_to_account_groups_mv.portfolio_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, assets_dm_next.type, assets_dm_next.name_en, sum(position_by_asset_mv_next.market_value), position_by_asset_mv_next.currency_code, (sum(position_by_asset_mv_next.weight) * 100:Decimal) as $expr3, sum(position_by_asset_mv_next.weight), portfolios_dm.base_currency_code, (sum(position_by_asset_mv_next.market_value) / sum(position_by_asset_mv_next.weight)) as $expr4, portfolio_to_account_groups_mv.account_group_id, portfolios_dm.portfolio_id, portfolio_to_account_groups_mv.$src, opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1, assets_dm_next.id] }
    ├── output: [ opportunity_conditions_mv_next.opportunity_id, portfolio_to_account_groups_mv.portfolio_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, assets_dm_next.type, assets_dm_next.name_en, sum(position_by_asset_mv_next.market_value), position_by_asset_mv_next.currency_code, $expr3, sum(position_by_asset_mv_next.weight), portfolios_dm.base_currency_code, $expr4, portfolio_to_account_groups_mv.account_group_id, portfolios_dm.portfolio_id, portfolio_to_account_groups_mv.$src, opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1, assets_dm_next.id ]
    ├── stream key: [ portfolio_to_account_groups_mv.account_group_id, portfolios_dm.portfolio_id, portfolio_to_account_groups_mv.$src, opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1, assets_dm_next.id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.currency_code ]
    └── MergeExecutor { output: [ position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, sum(position_by_asset_mv_next.market_value), sum(position_by_asset_mv_next.weight), portfolio_to_account_groups_mv.portfolio_id, portfolios_dm.base_currency_code, assets_dm_next.name_en, assets_dm_next.type, opportunity_conditions_mv_next.opportunity_id, portfolio_to_account_groups_mv.account_group_id, portfolios_dm.portfolio_id, portfolio_to_account_groups_mv.$src, position_by_asset_mv_next.account_group_id, opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1, assets_dm_next.id ], stream key: [ portfolio_to_account_groups_mv.account_group_id, portfolios_dm.portfolio_id, portfolio_to_account_groups_mv.$src, opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1, assets_dm_next.id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.currency_code ] }

Fragment 64431 (Actor 747622,747621)
StreamSyncLogStore { output: [ position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, sum(position_by_asset_mv_next.market_value), sum(position_by_asset_mv_next.weight), portfolio_to_account_groups_mv.portfolio_id, portfolios_dm.base_currency_code, assets_dm_next.name_en, assets_dm_next.type, opportunity_conditions_mv_next.opportunity_id, portfolio_to_account_groups_mv.account_group_id, portfolios_dm.portfolio_id, portfolio_to_account_groups_mv.$src, position_by_asset_mv_next.account_group_id, opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1, assets_dm_next.id ], stream key: [ portfolio_to_account_groups_mv.account_group_id, portfolios_dm.portfolio_id, portfolio_to_account_groups_mv.$src, opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1, assets_dm_next.id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.currency_code ] }
└── StreamHashJoin { type: Inner, predicate: portfolio_to_account_groups_mv.account_group_id = position_by_asset_mv_next.account_group_id }
    ├── output: [ position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, sum(position_by_asset_mv_next.market_value), sum(position_by_asset_mv_next.weight), portfolio_to_account_groups_mv.portfolio_id, portfolios_dm.base_currency_code, assets_dm_next.name_en, assets_dm_next.type, opportunity_conditions_mv_next.opportunity_id, portfolio_to_account_groups_mv.account_group_id, portfolios_dm.portfolio_id, portfolio_to_account_groups_mv.$src, position_by_asset_mv_next.account_group_id, opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1, assets_dm_next.id ]
    ├── stream key: [ portfolio_to_account_groups_mv.account_group_id, portfolios_dm.portfolio_id, portfolio_to_account_groups_mv.$src, opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1, assets_dm_next.id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.currency_code ]
    ├── MergeExecutor { output: [ portfolios_dm.base_currency_code, portfolio_to_account_groups_mv.portfolio_id, portfolio_to_account_groups_mv.account_group_id, portfolios_dm.portfolio_id, portfolio_to_account_groups_mv.$src ], stream key: [ portfolio_to_account_groups_mv.account_group_id, portfolios_dm.portfolio_id, portfolio_to_account_groups_mv.$src ] }
    └── MergeExecutor { output: [ opportunity_conditions_mv_next.opportunity_id, assets_dm_next.name_en, assets_dm_next.type, position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, sum(position_by_asset_mv_next.market_value), sum(position_by_asset_mv_next.weight), opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1, 'GET_PORTFOLIO_EXTREME_SINGLE_HOLDING_PERCENTAGE':Varchar, assets_dm_next.id ], stream key: [ position_by_asset_mv_next.account_group_id, opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1, assets_dm_next.id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.currency_code ] }

Fragment 64432 (Actor 747675,747676)
StreamLocalityProvider { locality_columns: [portfolio_to_account_groups_mv.account_group_id] } { output: [ portfolios_dm.base_currency_code, portfolio_to_account_groups_mv.portfolio_id, portfolio_to_account_groups_mv.account_group_id, portfolios_dm.portfolio_id, portfolio_to_account_groups_mv.$src ], stream key: [ portfolio_to_account_groups_mv.account_group_id, portfolios_dm.portfolio_id, portfolio_to_account_groups_mv.$src ] }
└── MergeExecutor { output: [ portfolios_dm.base_currency_code, portfolio_to_account_groups_mv.portfolio_id, portfolio_to_account_groups_mv.account_group_id, portfolios_dm.portfolio_id, portfolio_to_account_groups_mv.$src ], stream key: [ portfolios_dm.portfolio_id, portfolio_to_account_groups_mv.$src ] }

Fragment 64433 (Actor 747677,747678)
StreamSyncLogStore { output: [ portfolios_dm.base_currency_code, portfolio_to_account_groups_mv.portfolio_id, portfolio_to_account_groups_mv.account_group_id, portfolios_dm.portfolio_id, portfolio_to_account_groups_mv.$src ], stream key: [ portfolios_dm.portfolio_id, portfolio_to_account_groups_mv.$src ] }
└── StreamHashJoin { type: Inner, predicate: portfolios_dm.portfolio_id = portfolio_to_account_groups_mv.portfolio_id } { output: [ portfolios_dm.base_currency_code, portfolio_to_account_groups_mv.portfolio_id, portfolio_to_account_groups_mv.account_group_id, portfolios_dm.portfolio_id, portfolio_to_account_groups_mv.$src ], stream key: [ portfolios_dm.portfolio_id, portfolio_to_account_groups_mv.$src ] }
    ├── MergeExecutor { output: [ portfolios_dm.portfolio_id, portfolios_dm.base_currency_code ], stream key: [ portfolios_dm.portfolio_id ] }
    └── MergeExecutor { output: [ portfolio_to_account_groups_mv.portfolio_id, portfolio_to_account_groups_mv.account_group_id, portfolio_to_account_groups_mv.$src ], stream key: [ portfolio_to_account_groups_mv.portfolio_id, portfolio_to_account_groups_mv.$src ] }

Fragment 64434 (Actor 747604,747603)
StreamTableScan { table: portfolios_dm, columns: [portfolio_id, base_currency_code] } { output: [ portfolios_dm.portfolio_id, portfolios_dm.base_currency_code ], stream key: [ portfolios_dm.portfolio_id ] }
├── Upstream { output: [ portfolio_id, base_currency_code ], stream key: [] }
└── BatchPlanNode { output: [ portfolio_id, base_currency_code ], stream key: [] }

Fragment 64435 (Actor 747679,747680)
StreamLocalityProvider { locality_columns: [portfolio_to_account_groups_mv.portfolio_id] } { output: [ portfolio_to_account_groups_mv.portfolio_id, portfolio_to_account_groups_mv.account_group_id, portfolio_to_account_groups_mv.$src ], stream key: [ portfolio_to_account_groups_mv.portfolio_id, portfolio_to_account_groups_mv.$src ] }
└── MergeExecutor { output: [ portfolio_to_account_groups_mv.portfolio_id, portfolio_to_account_groups_mv.account_group_id, portfolio_to_account_groups_mv.$src ], stream key: [ portfolio_to_account_groups_mv.portfolio_id, portfolio_to_account_groups_mv.$src ] }

Fragment 64436 (Actor 747606,747605)
StreamProject { exprs: [portfolio_to_account_groups_mv.portfolio_id, portfolio_to_account_groups_mv.account_group_id, portfolio_to_account_groups_mv.$src] } { output: [ portfolio_to_account_groups_mv.portfolio_id, portfolio_to_account_groups_mv.account_group_id, portfolio_to_account_groups_mv.$src ], stream key: [ portfolio_to_account_groups_mv.portfolio_id, portfolio_to_account_groups_mv.$src ] }
└── StreamFilter { predicate: (portfolio_to_account_groups_mv.type = 'all':Varchar) } { output: [ portfolio_to_account_groups_mv.portfolio_id, portfolio_to_account_groups_mv.account_group_id, portfolio_to_account_groups_mv.$src, portfolio_to_account_groups_mv.type ], stream key: [ portfolio_to_account_groups_mv.portfolio_id, portfolio_to_account_groups_mv.$src ] }
    └── StreamTableScan { table: portfolio_to_account_groups_mv, columns: [portfolio_id, account_group_id, $src, type] } { output: [ portfolio_to_account_groups_mv.portfolio_id, portfolio_to_account_groups_mv.account_group_id, portfolio_to_account_groups_mv.$src, portfolio_to_account_groups_mv.type ], stream key: [ portfolio_to_account_groups_mv.portfolio_id, portfolio_to_account_groups_mv.$src ] }
        ├── Upstream { output: [ portfolio_id, account_group_id, $src, type ], stream key: [] }
        └── BatchPlanNode { output: [ portfolio_id, account_group_id, $src, type ], stream key: [] }

Fragment 64437 (Actor 747682,747681)
StreamLocalityProvider { locality_columns: [position_by_asset_mv_next.account_group_id] } { output: [ opportunity_conditions_mv_next.opportunity_id, assets_dm_next.name_en, assets_dm_next.type, position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, sum(position_by_asset_mv_next.market_value), sum(position_by_asset_mv_next.weight), opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1, 'GET_PORTFOLIO_EXTREME_SINGLE_HOLDING_PERCENTAGE':Varchar, assets_dm_next.id ], stream key: [ position_by_asset_mv_next.account_group_id, opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1, assets_dm_next.id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.currency_code ] }
└── MergeExecutor { output: [ opportunity_conditions_mv_next.opportunity_id, assets_dm_next.name_en, assets_dm_next.type, position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, sum(position_by_asset_mv_next.market_value), sum(position_by_asset_mv_next.weight), opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1, 'GET_PORTFOLIO_EXTREME_SINGLE_HOLDING_PERCENTAGE':Varchar, assets_dm_next.id ], stream key: [ opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1, assets_dm_next.id, position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.currency_code ] }

Fragment 64438 (Actor 747683,747684)
StreamSyncLogStore { output: [ opportunity_conditions_mv_next.opportunity_id, assets_dm_next.name_en, assets_dm_next.type, position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, sum(position_by_asset_mv_next.market_value), sum(position_by_asset_mv_next.weight), opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1, 'GET_PORTFOLIO_EXTREME_SINGLE_HOLDING_PERCENTAGE':Varchar, assets_dm_next.id ], stream key: [ opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1, assets_dm_next.id, position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.currency_code ] }
└── StreamHashJoin { type: Inner, predicate: opportunity_conditions_mv_next.activity_name = 'GET_PORTFOLIO_EXTREME_SINGLE_HOLDING_PERCENTAGE':Varchar AND Case((opportunity_conditions_mv_next.op = 'GTE':Varchar), (sum(position_by_asset_mv_next.weight) >= opportunity_conditions_mv_next.threshold), (opportunity_conditions_mv_next.op = 'GT':Varchar), (sum(position_by_asset_mv_next.weight) > opportunity_conditions_mv_next.threshold), (opportunity_conditions_mv_next.op = 'LTE':Varchar), (sum(position_by_asset_mv_next.weight) <= opportunity_conditions_mv_next.threshold), (opportunity_conditions_mv_next.op = 'LT':Varchar), (sum(position_by_asset_mv_next.weight) < opportunity_conditions_mv_next.threshold), (opportunity_conditions_mv_next.op = 'EQ':Varchar), (sum(position_by_asset_mv_next.weight) = opportunity_conditions_mv_next.threshold), false:Boolean) }
    ├── output: [ opportunity_conditions_mv_next.opportunity_id, assets_dm_next.name_en, assets_dm_next.type, position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, sum(position_by_asset_mv_next.market_value), sum(position_by_asset_mv_next.weight), opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1, 'GET_PORTFOLIO_EXTREME_SINGLE_HOLDING_PERCENTAGE':Varchar, assets_dm_next.id ]
    ├── stream key: [ opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1, assets_dm_next.id, position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.currency_code ]
    ├── MergeExecutor { output: [ opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next.op, opportunity_conditions_mv_next.threshold, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1 ], stream key: [ opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1 ] }
    └── MergeExecutor { output: [ assets_dm_next.name_en, assets_dm_next.type, position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, sum(position_by_asset_mv_next.market_value), sum(position_by_asset_mv_next.weight), 'GET_PORTFOLIO_EXTREME_SINGLE_HOLDING_PERCENTAGE':Varchar, assets_dm_next.id ], stream key: [ 'GET_PORTFOLIO_EXTREME_SINGLE_HOLDING_PERCENTAGE':Varchar, assets_dm_next.id, position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.currency_code ] }

Fragment 64439 (Actor 747686,747685)
StreamLocalityProvider { locality_columns: [opportunity_conditions_mv_next.activity_name] } { output: [ opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next.op, opportunity_conditions_mv_next.threshold, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1 ], stream key: [ opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1 ] }
└── MergeExecutor { output: [ opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next.op, opportunity_conditions_mv_next.threshold, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1 ], stream key: [ opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1 ] }

Fragment 64440 (Actor 741901,741900)
StreamTableScan { table: opportunity_conditions_mv_next, columns: [opportunity_id, activity_name, op, threshold, _rw_projected_row_id, _rw_projected_row_id#1] } { output: [ opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next.op, opportunity_conditions_mv_next.threshold, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1 ], stream key: [ opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1 ] }
├── Upstream { output: [ opportunity_id, activity_name, op, threshold, _rw_projected_row_id, _rw_projected_row_id#1 ], stream key: [] }
└── BatchPlanNode { output: [ opportunity_id, activity_name, op, threshold, _rw_projected_row_id, _rw_projected_row_id#1 ], stream key: [] }

Fragment 64441 (Actor 747687,747688)
StreamLocalityProvider { locality_columns: ['GET_PORTFOLIO_EXTREME_SINGLE_HOLDING_PERCENTAGE':Varchar] } { output: [ assets_dm_next.name_en, assets_dm_next.type, position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, sum(position_by_asset_mv_next.market_value), sum(position_by_asset_mv_next.weight), 'GET_PORTFOLIO_EXTREME_SINGLE_HOLDING_PERCENTAGE':Varchar, assets_dm_next.id ], stream key: [ 'GET_PORTFOLIO_EXTREME_SINGLE_HOLDING_PERCENTAGE':Varchar, assets_dm_next.id, position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.currency_code ] }
└── MergeExecutor { output: [ assets_dm_next.name_en, assets_dm_next.type, position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, sum(position_by_asset_mv_next.market_value), sum(position_by_asset_mv_next.weight), 'GET_PORTFOLIO_EXTREME_SINGLE_HOLDING_PERCENTAGE':Varchar, assets_dm_next.id ], stream key: [ assets_dm_next.id, position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.currency_code ] }

Fragment 64442 (Actor 747625,747626)
StreamProject { exprs: [assets_dm_next.name_en, assets_dm_next.type, position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, sum(position_by_asset_mv_next.market_value), sum(position_by_asset_mv_next.weight), 'GET_PORTFOLIO_EXTREME_SINGLE_HOLDING_PERCENTAGE':Varchar, assets_dm_next.id] } { output: [ assets_dm_next.name_en, assets_dm_next.type, position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, sum(position_by_asset_mv_next.market_value), sum(position_by_asset_mv_next.weight), 'GET_PORTFOLIO_EXTREME_SINGLE_HOLDING_PERCENTAGE':Varchar, assets_dm_next.id ], stream key: [ assets_dm_next.id, position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.currency_code ] }
└── MergeExecutor { output: [ assets_dm_next.id, assets_dm_next.name_en, assets_dm_next.type, position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, sum(position_by_asset_mv_next.market_value), sum(position_by_asset_mv_next.weight) ], stream key: [ assets_dm_next.id, position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.currency_code ] }

Fragment 64443 (Actor 747628,747627)
StreamSyncLogStore { output: [ assets_dm_next.id, assets_dm_next.name_en, assets_dm_next.type, position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, sum(position_by_asset_mv_next.market_value), sum(position_by_asset_mv_next.weight) ], stream key: [ assets_dm_next.id, position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.currency_code ] }
└── StreamHashJoin { type: Inner, predicate: assets_dm_next.id = position_by_asset_mv_next.asset_id } { output: [ assets_dm_next.id, assets_dm_next.name_en, assets_dm_next.type, position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, sum(position_by_asset_mv_next.market_value), sum(position_by_asset_mv_next.weight) ], stream key: [ assets_dm_next.id, position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.currency_code ] }
    ├── MergeExecutor { output: [ assets_dm_next.id, assets_dm_next.name_en, assets_dm_next.type ], stream key: [ assets_dm_next.id ] }
    └── MergeExecutor { output: [ position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, sum(position_by_asset_mv_next.market_value), sum(position_by_asset_mv_next.weight) ], stream key: [ position_by_asset_mv_next.asset_id, position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.currency_code ] }

Fragment 64444 (Actor 747694,747695)
StreamTableScan { table: assets_dm_next, columns: [id, name_en, type] } { output: [ assets_dm_next.id, assets_dm_next.name_en, assets_dm_next.type ], stream key: [ assets_dm_next.id ] }
├── Upstream { output: [ id, name_en, type ], stream key: [] }
└── BatchPlanNode { output: [ id, name_en, type ], stream key: [] }

Fragment 64445 (Actor 747689,747690)
StreamLocalityProvider { locality_columns: [position_by_asset_mv_next.asset_id] } { output: [ position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, sum(position_by_asset_mv_next.market_value), sum(position_by_asset_mv_next.weight) ], stream key: [ position_by_asset_mv_next.asset_id, position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.currency_code ] }
└── MergeExecutor { output: [ position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, sum(position_by_asset_mv_next.market_value), sum(position_by_asset_mv_next.weight) ], stream key: [ position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code ] }

Fragment 64446 (Actor 747692,747691)
StreamFilter { predicate: (sum(position_by_asset_mv_next.market_value) > 0:Decimal) AND (sum(position_by_asset_mv_next.weight) > 0:Decimal) AND (sum(position_by_asset_mv_next.weight) < 1:Decimal) } { output: [ position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, sum(position_by_asset_mv_next.market_value), sum(position_by_asset_mv_next.weight) ], stream key: [ position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code ] }
└── StreamProject { exprs: [position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, sum(position_by_asset_mv_next.market_value), sum(position_by_asset_mv_next.weight)] } { output: [ position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, sum(position_by_asset_mv_next.market_value), sum(position_by_asset_mv_next.weight) ], stream key: [ position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code ] }
    └── StreamHashAgg { group_key: [position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code], aggs: [sum(position_by_asset_mv_next.market_value), sum(position_by_asset_mv_next.weight), count] } { output: [ position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, sum(position_by_asset_mv_next.market_value), sum(position_by_asset_mv_next.weight), count ], stream key: [ position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code ] }
        └── StreamLocalityProvider { locality_columns: [position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code] }
            ├── output: [ position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, position_by_asset_mv_next.market_value, position_by_asset_mv_next.weight, position_by_asset_mv_next.position_type, position_by_asset_mv_next.position_values_mv_next.holding_currency_expanded, position_by_asset_mv_next.asset_currency, position_by_asset_mv_next.source_entity_type, position_by_asset_mv_next.position_values_mv_next.position_type_expanded, position_by_asset_mv_next.flag, position_by_asset_mv_next.position_summary_mv_next.source_entity_type ]
            ├── stream key: [ position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, position_by_asset_mv_next.position_type, position_by_asset_mv_next.position_values_mv_next.holding_currency_expanded, position_by_asset_mv_next.asset_currency, position_by_asset_mv_next.source_entity_type, position_by_asset_mv_next.position_values_mv_next.position_type_expanded, position_by_asset_mv_next.flag, position_by_asset_mv_next.position_summary_mv_next.source_entity_type ]
            └── MergeExecutor { output: [ position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, position_by_asset_mv_next.market_value, position_by_asset_mv_next.weight, position_by_asset_mv_next.position_type, position_by_asset_mv_next.position_values_mv_next.holding_currency_expanded, position_by_asset_mv_next.asset_currency, position_by_asset_mv_next.source_entity_type, position_by_asset_mv_next.position_values_mv_next.position_type_expanded, position_by_asset_mv_next.flag, position_by_asset_mv_next.position_summary_mv_next.source_entity_type ], stream key: [ position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.position_type, position_by_asset_mv_next.currency_code, position_by_asset_mv_next.position_values_mv_next.holding_currency_expanded, position_by_asset_mv_next.asset_currency, position_by_asset_mv_next.source_entity_type, position_by_asset_mv_next.position_values_mv_next.position_type_expanded, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.flag, position_by_asset_mv_next.position_summary_mv_next.source_entity_type ] }

Fragment 64447 (Actor 747696,747697)
StreamProject { exprs: [position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, position_by_asset_mv_next.market_value, position_by_asset_mv_next.weight, position_by_asset_mv_next.position_type, position_by_asset_mv_next.position_values_mv_next.holding_currency_expanded, position_by_asset_mv_next.asset_currency, position_by_asset_mv_next.source_entity_type, position_by_asset_mv_next.position_values_mv_next.position_type_expanded, position_by_asset_mv_next.flag, position_by_asset_mv_next.position_summary_mv_next.source_entity_type] }
├── output: [ position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, position_by_asset_mv_next.market_value, position_by_asset_mv_next.weight, position_by_asset_mv_next.position_type, position_by_asset_mv_next.position_values_mv_next.holding_currency_expanded, position_by_asset_mv_next.asset_currency, position_by_asset_mv_next.source_entity_type, position_by_asset_mv_next.position_values_mv_next.position_type_expanded, position_by_asset_mv_next.flag, position_by_asset_mv_next.position_summary_mv_next.source_entity_type ]
├── stream key: [ position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.position_type, position_by_asset_mv_next.currency_code, position_by_asset_mv_next.position_values_mv_next.holding_currency_expanded, position_by_asset_mv_next.asset_currency, position_by_asset_mv_next.source_entity_type, position_by_asset_mv_next.position_values_mv_next.position_type_expanded, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.flag, position_by_asset_mv_next.position_summary_mv_next.source_entity_type ]
└── StreamDynamicFilter { predicate: ($expr1 >= $expr2), output_watermarks: [[$expr1]], output: [position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, position_by_asset_mv_next.market_value, position_by_asset_mv_next.weight, $expr1, position_by_asset_mv_next.position_type, position_by_asset_mv_next.position_values_mv_next.holding_currency_expanded, position_by_asset_mv_next.asset_currency, position_by_asset_mv_next.source_entity_type, position_by_asset_mv_next.position_values_mv_next.position_type_expanded, position_by_asset_mv_next.flag, position_by_asset_mv_next.position_summary_mv_next.source_entity_type], cleaned_by_watermark: true }
    ├── output: [ position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, position_by_asset_mv_next.market_value, position_by_asset_mv_next.weight, $expr1, position_by_asset_mv_next.position_type, position_by_asset_mv_next.position_values_mv_next.holding_currency_expanded, position_by_asset_mv_next.asset_currency, position_by_asset_mv_next.source_entity_type, position_by_asset_mv_next.position_values_mv_next.position_type_expanded, position_by_asset_mv_next.flag, position_by_asset_mv_next.position_summary_mv_next.source_entity_type ]
    ├── stream key: [ position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.position_type, position_by_asset_mv_next.currency_code, position_by_asset_mv_next.position_values_mv_next.holding_currency_expanded, position_by_asset_mv_next.asset_currency, position_by_asset_mv_next.source_entity_type, position_by_asset_mv_next.position_values_mv_next.position_type_expanded, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.flag, position_by_asset_mv_next.position_summary_mv_next.source_entity_type ]
    ├── StreamProject { exprs: [position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, position_by_asset_mv_next.market_value, position_by_asset_mv_next.weight, AtTimeZone(position_by_asset_mv_next.dim_balance_date::Timestamp, 'UTC':Varchar) as $expr1, position_by_asset_mv_next.position_type, position_by_asset_mv_next.position_values_mv_next.holding_currency_expanded, position_by_asset_mv_next.asset_currency, position_by_asset_mv_next.source_entity_type, position_by_asset_mv_next.position_values_mv_next.position_type_expanded, position_by_asset_mv_next.flag, position_by_asset_mv_next.position_summary_mv_next.source_entity_type] }
    │   ├── output: [ position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, position_by_asset_mv_next.market_value, position_by_asset_mv_next.weight, $expr1, position_by_asset_mv_next.position_type, position_by_asset_mv_next.position_values_mv_next.holding_currency_expanded, position_by_asset_mv_next.asset_currency, position_by_asset_mv_next.source_entity_type, position_by_asset_mv_next.position_values_mv_next.position_type_expanded, position_by_asset_mv_next.flag, position_by_asset_mv_next.position_summary_mv_next.source_entity_type ]
    │   ├── stream key: [ position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.position_type, position_by_asset_mv_next.currency_code, position_by_asset_mv_next.position_values_mv_next.holding_currency_expanded, position_by_asset_mv_next.asset_currency, position_by_asset_mv_next.source_entity_type, position_by_asset_mv_next.position_values_mv_next.position_type_expanded, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.flag, position_by_asset_mv_next.position_summary_mv_next.source_entity_type ]
    │   └── StreamFilter { predicate: (position_by_asset_mv_next.source_entity_type = 'portfolio':Varchar) AND (position_by_asset_mv_next.position_type = 'POSITION':Varchar) AND Not(IsNull(position_by_asset_mv_next.asset_id)) }
    │       ├── output: [ position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, position_by_asset_mv_next.market_value, position_by_asset_mv_next.weight, position_by_asset_mv_next.position_type, position_by_asset_mv_next.position_values_mv_next.holding_currency_expanded, position_by_asset_mv_next.asset_currency, position_by_asset_mv_next.source_entity_type, position_by_asset_mv_next.position_values_mv_next.position_type_expanded, position_by_asset_mv_next.flag, position_by_asset_mv_next.position_summary_mv_next.source_entity_type ]
    │       ├── stream key: [ position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.position_type, position_by_asset_mv_next.currency_code, position_by_asset_mv_next.position_values_mv_next.holding_currency_expanded, position_by_asset_mv_next.asset_currency, position_by_asset_mv_next.source_entity_type, position_by_asset_mv_next.position_values_mv_next.position_type_expanded, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.flag, position_by_asset_mv_next.position_summary_mv_next.source_entity_type ]
    │       └── StreamTableScan { table: position_by_asset_mv_next, columns: [account_group_id, dim_balance_date, asset_id, currency_code, market_value, weight, position_type, position_values_mv_next.holding_currency_expanded, asset_currency, source_entity_type, position_values_mv_next.position_type_expanded, flag, position_summary_mv_next.source_entity_type] }
    │           ├── output: [ position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.currency_code, position_by_asset_mv_next.market_value, position_by_asset_mv_next.weight, position_by_asset_mv_next.position_type, position_by_asset_mv_next.position_values_mv_next.holding_currency_expanded, position_by_asset_mv_next.asset_currency, position_by_asset_mv_next.source_entity_type, position_by_asset_mv_next.position_values_mv_next.position_type_expanded, position_by_asset_mv_next.flag, position_by_asset_mv_next.position_summary_mv_next.source_entity_type ]
    │           ├── stream key: [ position_by_asset_mv_next.account_group_id, position_by_asset_mv_next.dim_balance_date, position_by_asset_mv_next.position_type, position_by_asset_mv_next.currency_code, position_by_asset_mv_next.position_values_mv_next.holding_currency_expanded, position_by_asset_mv_next.asset_currency, position_by_asset_mv_next.source_entity_type, position_by_asset_mv_next.position_values_mv_next.position_type_expanded, position_by_asset_mv_next.asset_id, position_by_asset_mv_next.flag, position_by_asset_mv_next.position_summary_mv_next.source_entity_type ]
    │           ├── Upstream { output: [ account_group_id, dim_balance_date, asset_id, currency_code, market_value, weight, position_type, position_values_mv_next.holding_currency_expanded, asset_currency, source_entity_type, position_values_mv_next.position_type_expanded, flag, position_summary_mv_next.source_entity_type ], stream key: [] }
    │           └── BatchPlanNode { output: [ account_group_id, dim_balance_date, asset_id, currency_code, market_value, weight, position_type, position_values_mv_next.holding_currency_expanded, asset_currency, source_entity_type, position_values_mv_next.position_type_expanded, flag, position_summary_mv_next.source_entity_type ], stream key: [] }
    └── MergeExecutor { output: [ $expr2 ], stream key: [] }

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