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

← cluster adib_rm objects party_portfolios_mv explain
Overview Objects Graph History
materialized view · adib_rm.party_portfolios_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
122 operators
Materialize · adib_rm.party_portfolios_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
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · LeftOuter · portfolio_to_account_groups_mv.account_group_id = pnl_snaps…
2 actors
HashJoin · LeftOuter · portfolio_to_account_groups_mv.account_group_id = pnl_snaps… 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
Filter · pnl_snapshot_mv_next
0% idle 2 actors
StreamScan · pnl_snapshot_mv_next
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · LeftOuter · portfolio_to_account_groups_mv.account_group_id = intraday_…
2 actors
HashJoin · LeftOuter · portfolio_to_account_groups_mv.account_group_id = intraday_… 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
Filter · intraday_position_summary_mv_next
0% idle 2 actors
StreamScan · intraday_position_summary_mv_next
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · LeftOuter · portfolio_to_account_groups_mv.account_group_id = position_…
2 actors
HashJoin · LeftOuter · 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
Project · position_snapshot_mv_next
2 actors
Filter · position_snapshot_mv_next
0% idle 2 actors
StreamScan · position_snapshot_mv_next
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · LeftOuter · portfolios_plain_mv.portfolio_id = portfolio_to_account_gro…
2 actors
HashJoin · LeftOuter · portfolios_plain_mv.portfolio_id = portfolio_to_account_gro… 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 · 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
SyncLogStore · LeftOuter · min(party_involvements_dm.customer_relationship_id) = lifec…
2 actors
HashJoin · LeftOuter · min(party_involvements_dm.customer_relationship_id) = lifec… 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 · lifecycle_profiles
2 actors
Filter · lifecycle_profiles
0% idle 2 actors
StreamScan · lifecycle_profiles
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
TemporalJoin · Inner · portfolios_plain_mv.service_type_id = service_types_dm.serv…
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · service_types_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · party_involvements_dm.entity_id = portfolios_plain_mv.portf…
2 actors
HashJoin · Inner · party_involvements_dm.entity_id = portfolios_plain_mv.portf… 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 · portfolios_plain_mv
2 actors
Filter · portfolios_plain_mv
0% idle 2 actors
StreamScan · portfolios_plain_mv
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 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
SyncLogStore · Inner · party_involvements_dm.customer_relationship_id = customer_r…
2 actors
HashJoin · Inner · party_involvements_dm.customer_relationship_id = customer_r… 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 · customer_relationships_next
2 actors
Filter · customer_relationships_next
0% idle 2 actors
StreamScan · customer_relationships_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_involvements_dm.effective_to)
2 actors
Filter · IsNull(party_involvements_dm.effective_to)
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · party_involvements_dm
2 actors
DynamicFilter · party_involvements_dm Dynamic filter — verify it pairs with a temporal condition to clean state
0% idle 2 actors
Merge
2 actors
Exchange
0% 2/s 0 actors
Now
0% 2/s 1 actor
Project · party_involvements_dm
2 actors
Filter · party_involvements_dm
0% idle 2 actors
StreamScan · party_involvements_dm
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_portfolios_mv Materialize adib_rm.party_portfolio… 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 Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · LeftOuter · portfolio_to_account_groups_mv.account_group_id = pnl_snaps… SyncLogStore LeftOuter · portfolio_t… — · 2 actors HashJoin · LeftOuter · portfolio_to_account_groups_mv.account_group_id = pnl_snaps… HashJoin LeftOuter · portfolio_t… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Filter · pnl_snapshot_mv_next Filter pnl_snapshot_mv_next idle · 2 actors StreamScan · pnl_snapshot_mv_next StreamScan pnl_snapshot_mv_next idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · LeftOuter · portfolio_to_account_groups_mv.account_group_id = intraday_… SyncLogStore LeftOuter · portfolio_t… — · 2 actors HashJoin · LeftOuter · portfolio_to_account_groups_mv.account_group_id = intraday_… HashJoin LeftOuter · portfolio_t… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Filter · intraday_position_summary_mv_next Filter intraday_position_summa… idle · 2 actors StreamScan · intraday_position_summary_mv_next StreamScan intraday_position_summa… idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · LeftOuter · portfolio_to_account_groups_mv.account_group_id = position_… SyncLogStore LeftOuter · portfolio_t… — · 2 actors HashJoin · LeftOuter · portfolio_to_account_groups_mv.account_group_id = position_… HashJoin LeftOuter · portfolio_t… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · position_snapshot_mv_next Project position_snapshot_mv_ne… — · 2 actors Filter · position_snapshot_mv_next Filter position_snapshot_mv_ne… idle · 2 actors StreamScan · position_snapshot_mv_next StreamScan position_snapshot_mv_ne… idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · LeftOuter · portfolios_plain_mv.portfolio_id = portfolio_to_account_gro… SyncLogStore LeftOuter · portfolios_… — · 2 actors HashJoin · LeftOuter · portfolios_plain_mv.portfolio_id = portfolio_to_account_gro… HashJoin LeftOuter · portfolios_… 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 SyncLogStore · LeftOuter · min(party_involvements_dm.customer_relationship_id) = lifec… SyncLogStore LeftOuter · min(party_i… — · 2 actors HashJoin · LeftOuter · min(party_involvements_dm.customer_relationship_id) = lifec… HashJoin LeftOuter · min(party_i… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · lifecycle_profiles Project lifecycle_profiles — · 2 actors Filter · lifecycle_profiles Filter lifecycle_profiles idle · 2 actors StreamScan · lifecycle_profiles StreamScan lifecycle_profiles idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors TemporalJoin · Inner · portfolios_plain_mv.service_type_id = service_types_dm.serv… TemporalJoin Inner · portfolios_plai… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · service_types_dm StreamScan service_types_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · party_involvements_dm.entity_id = portfolios_plain_mv.portf… SyncLogStore Inner · party_involveme… — · 2 actors HashJoin · Inner · party_involvements_dm.entity_id = portfolios_plain_mv.portf… HashJoin Inner · party_involveme… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · portfolios_plain_mv Project portfolios_plain_mv — · 2 actors Filter · portfolios_plain_mv Filter portfolios_plain_mv idle · 2 actors StreamScan · portfolios_plain_mv StreamScan portfolios_plain_mv idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 2 actors HashAgg HashAgg idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · party_involvements_dm.customer_relationship_id = customer_r… SyncLogStore Inner · party_involveme… — · 2 actors HashJoin · Inner · party_involvements_dm.customer_relationship_id = customer_r… HashJoin Inner · party_involveme… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · customer_relationships_next Project customer_relationships_… — · 2 actors Filter · customer_relationships_next Filter customer_relationships_… idle · 2 actors StreamScan · customer_relationships_next StreamScan customer_relationships_… 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_involvements_dm.effective_to) Project IsNull(party_involvemen… — · 2 actors Filter · IsNull(party_involvements_dm.effective_to) Filter IsNull(party_involvemen… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · party_involvements_dm Project party_involvements_dm — · 2 actors DynamicFilter · party_involvements_dm DynamicFilter party_involvements_dm idle · 2 actors Merge Merge — · 2 actors Exchange Exchange 2/s · 0 actors Now Now 2/s · 1 actor Project · party_involvements_dm Project party_involvements_dm — · 2 actors Filter · party_involvements_dm Filter party_involvements_dm idle · 2 actors StreamScan · party_involvements_dm StreamScan party_involvements_dm 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 26086 (Actor 132834,132835)
StreamMaterialize { columns: [party_id, portfolios], stream_key: [party_id], pk_columns: [party_id], pk_conflict: NoCheck }
├── output: [ party_involvements_dm.party_id, jsonb_agg($expr3 order_by(portfolios_plain_mv.portfolio_id ASC)) ]
├── stream key: [ party_involvements_dm.party_id ]
└── StreamProject { exprs: [party_involvements_dm.party_id, jsonb_agg($expr3 order_by(portfolios_plain_mv.portfolio_id ASC))] }
    ├── output: [ party_involvements_dm.party_id, jsonb_agg($expr3 order_by(portfolios_plain_mv.portfolio_id ASC)) ]
    ├── stream key: [ party_involvements_dm.party_id ]
    └── StreamHashAgg { group_key: [party_involvements_dm.party_id], aggs: [jsonb_agg($expr3 order_by(portfolios_plain_mv.portfolio_id ASC)), count] }
        ├── output: [ party_involvements_dm.party_id, jsonb_agg($expr3 order_by(portfolios_plain_mv.portfolio_id ASC)), count ]
        ├── stream key: [ party_involvements_dm.party_id ]
        └── MergeExecutor
            ├── output:
            │   ┌── party_involvements_dm.party_id
            │   ├── $expr3
            │   ├── portfolios_plain_mv.portfolio_id
            │   ├── party_involvements_dm.entity_id
            │   ├── portfolios_plain_mv.service_type_id
            │   ├── lifecycle_profiles.id
            │   ├── min(party_involvements_dm.customer_relationship_id)
            │   ├── position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded
            │   ├── position_snapshot_mv_next.flag
            │   ├── portfolio_to_account_groups_mv.account_group_id
            │   ├── portfolios_plain_mv.base_currency_code
            │   ├── intraday_position_summary_mv_next.position_type
            │   └── pnl_snapshot_mv_next.position_type
            └── stream key:
                ┌── party_involvements_dm.party_id
                ├── party_involvements_dm.entity_id
                ├── portfolios_plain_mv.service_type_id
                ├── lifecycle_profiles.id
                ├── min(party_involvements_dm.customer_relationship_id)
                ├── portfolios_plain_mv.portfolio_id
                ├── position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded
                ├── position_snapshot_mv_next.flag
                ├── portfolio_to_account_groups_mv.account_group_id
                ├── portfolios_plain_mv.base_currency_code
                ├── intraday_position_summary_mv_next.position_type
                └── pnl_snapshot_mv_next.position_type

Fragment 26087 (Actor 132843,132842)
StreamProject
└─exprs:
  ┌─party_involvements_dm.party_id
  ├─JsonbBuildObject('id':Varchar, portfolios_plain_mv.portfolio_id, 'name':Varchar, portfolios_plain_mv.name, 'number':Varchar, portfolios_plain_mv.number, 'serviceType':Varchar, service_types_dm.type, 'marketValue':Varchar, JsonbBuildObject('amount':Varchar, Coalesce(intraday_position_summary_mv_next.market_value, position_snapshot_mv_next.market_value)::Varchar, 'currencyCode':Varchar, portfolios_plain_mv.base_currency_code), 'marketValueSystemCurrency':Varchar, JsonbBuildObject('amount':Varchar, Case(Not(IsNull(intraday_position_summary_mv_next.market_value)), intraday_position_summary_mv_next.market_value_system_currency::Varchar, position_snapshot_mv_next.market_value_system_currency::Varchar), 'currencyCode':Varchar, lifecycle_profiles.base_currency_code), 'unrealizedGainLoss':Varchar, JsonbBuildObject('value':Varchar, JsonbBuildObject('amount':Varchar, Case(Not(IsNull(intraday_position_summary_mv_next.market_value)), (intraday_position_summary_mv_next.market_value - intraday_position_summary_mv_next.total_average_cost)::Varchar, pnl_snapshot_mv_next.unrealized_gain_loss::Varchar), 'currencyCode':Varchar, portfolios_plain_mv.base_currency_code), 'percentage':Varchar, Case(Not(IsNull(intraday_position_summary_mv_next.market_value)), Case((IsNull(intraday_position_summary_mv_next.total_average_cost) OR (intraday_position_summary_mv_next.total_average_cost = 0:Decimal)), null:Varchar, ((intraday_position_summary_mv_next.market_value - intraday_position_summary_mv_next.total_average_cost) / intraday_position_summary_mv_next.total_average_cost)::Varchar), Case((IsNull(pnl_snapshot_mv_next.total_average_cost) OR (pnl_snapshot_mv_next.total_average_cost = 0:Decimal)), null:Varchar, (pnl_snapshot_mv_next.unrealized_gain_loss / pnl_snapshot_mv_next.total_average_cost)::Varchar)))) as $expr3
  ├─portfolios_plain_mv.portfolio_id
  ├─party_involvements_dm.entity_id
  ├─portfolios_plain_mv.service_type_id
  ├─lifecycle_profiles.id
  ├─min(party_involvements_dm.customer_relationship_id)
  ├─position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded
  ├─position_snapshot_mv_next.flag
  ├─portfolio_to_account_groups_mv.account_group_id
  ├─portfolios_plain_mv.base_currency_code
  ├─intraday_position_summary_mv_next.position_type
  └─pnl_snapshot_mv_next.position_type
├── output: [ party_involvements_dm.party_id, $expr3, portfolios_plain_mv.portfolio_id, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, lifecycle_profiles.id, min(party_involvements_dm.customer_relationship_id), position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded, position_snapshot_mv_next.flag, portfolio_to_account_groups_mv.account_group_id, portfolios_plain_mv.base_currency_code, intraday_position_summary_mv_next.position_type, pnl_snapshot_mv_next.position_type ]
├── stream key: [ party_involvements_dm.party_id, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, lifecycle_profiles.id, min(party_involvements_dm.customer_relationship_id), portfolios_plain_mv.portfolio_id, position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded, position_snapshot_mv_next.flag, portfolio_to_account_groups_mv.account_group_id, portfolios_plain_mv.base_currency_code, intraday_position_summary_mv_next.position_type, pnl_snapshot_mv_next.position_type ]
└── MergeExecutor { output: [ party_involvements_dm.party_id, portfolios_plain_mv.portfolio_id, portfolios_plain_mv.name, portfolios_plain_mv.number, portfolios_plain_mv.base_currency_code, service_types_dm.type, lifecycle_profiles.base_currency_code, position_snapshot_mv_next.market_value, position_snapshot_mv_next.market_value_system_currency, intraday_position_summary_mv_next.market_value, intraday_position_summary_mv_next.total_average_cost, intraday_position_summary_mv_next.market_value_system_currency, pnl_snapshot_mv_next.unrealized_gain_loss, pnl_snapshot_mv_next.total_average_cost, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, lifecycle_profiles.id, min(party_involvements_dm.customer_relationship_id), position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded, position_snapshot_mv_next.flag, portfolio_to_account_groups_mv.account_group_id, intraday_position_summary_mv_next.position_type, pnl_snapshot_mv_next.account_group_id, pnl_snapshot_mv_next.position_type, pnl_snapshot_mv_next.currency_code ], stream key: [ party_involvements_dm.party_id, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, lifecycle_profiles.id, min(party_involvements_dm.customer_relationship_id), portfolios_plain_mv.portfolio_id, position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded, position_snapshot_mv_next.flag, portfolio_to_account_groups_mv.account_group_id, portfolios_plain_mv.base_currency_code, intraday_position_summary_mv_next.position_type, pnl_snapshot_mv_next.position_type ] }

Fragment 26088 (Actor 132837,132836)
StreamSyncLogStore { output: [ party_involvements_dm.party_id, portfolios_plain_mv.portfolio_id, portfolios_plain_mv.name, portfolios_plain_mv.number, portfolios_plain_mv.base_currency_code, service_types_dm.type, lifecycle_profiles.base_currency_code, position_snapshot_mv_next.market_value, position_snapshot_mv_next.market_value_system_currency, intraday_position_summary_mv_next.market_value, intraday_position_summary_mv_next.total_average_cost, intraday_position_summary_mv_next.market_value_system_currency, pnl_snapshot_mv_next.unrealized_gain_loss, pnl_snapshot_mv_next.total_average_cost, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, lifecycle_profiles.id, min(party_involvements_dm.customer_relationship_id), position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded, position_snapshot_mv_next.flag, portfolio_to_account_groups_mv.account_group_id, intraday_position_summary_mv_next.position_type, pnl_snapshot_mv_next.account_group_id, pnl_snapshot_mv_next.position_type, pnl_snapshot_mv_next.currency_code ], stream key: [ party_involvements_dm.party_id, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, lifecycle_profiles.id, min(party_involvements_dm.customer_relationship_id), portfolios_plain_mv.portfolio_id, position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded, position_snapshot_mv_next.flag, portfolio_to_account_groups_mv.account_group_id, portfolios_plain_mv.base_currency_code, intraday_position_summary_mv_next.position_type, pnl_snapshot_mv_next.position_type ] }
└── StreamHashJoin { type: LeftOuter, predicate: portfolio_to_account_groups_mv.account_group_id = pnl_snapshot_mv_next.account_group_id AND portfolios_plain_mv.base_currency_code = pnl_snapshot_mv_next.currency_code } { output: [ party_involvements_dm.party_id, portfolios_plain_mv.portfolio_id, portfolios_plain_mv.name, portfolios_plain_mv.number, portfolios_plain_mv.base_currency_code, service_types_dm.type, lifecycle_profiles.base_currency_code, position_snapshot_mv_next.market_value, position_snapshot_mv_next.market_value_system_currency, intraday_position_summary_mv_next.market_value, intraday_position_summary_mv_next.total_average_cost, intraday_position_summary_mv_next.market_value_system_currency, pnl_snapshot_mv_next.unrealized_gain_loss, pnl_snapshot_mv_next.total_average_cost, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, lifecycle_profiles.id, min(party_involvements_dm.customer_relationship_id), position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded, position_snapshot_mv_next.flag, portfolio_to_account_groups_mv.account_group_id, intraday_position_summary_mv_next.position_type, pnl_snapshot_mv_next.account_group_id, pnl_snapshot_mv_next.position_type, pnl_snapshot_mv_next.currency_code ], stream key: [ party_involvements_dm.party_id, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, lifecycle_profiles.id, min(party_involvements_dm.customer_relationship_id), portfolios_plain_mv.portfolio_id, position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded, position_snapshot_mv_next.flag, portfolio_to_account_groups_mv.account_group_id, portfolios_plain_mv.base_currency_code, intraday_position_summary_mv_next.position_type, pnl_snapshot_mv_next.position_type ] }
    ├── MergeExecutor { output: [ party_involvements_dm.party_id, portfolios_plain_mv.portfolio_id, portfolios_plain_mv.name, portfolios_plain_mv.number, portfolios_plain_mv.base_currency_code, service_types_dm.type, lifecycle_profiles.base_currency_code, portfolio_to_account_groups_mv.account_group_id, position_snapshot_mv_next.market_value, position_snapshot_mv_next.market_value_system_currency, intraday_position_summary_mv_next.market_value, intraday_position_summary_mv_next.total_average_cost, intraday_position_summary_mv_next.market_value_system_currency, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, lifecycle_profiles.id, min(party_involvements_dm.customer_relationship_id), position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded, position_snapshot_mv_next.flag, intraday_position_summary_mv_next.account_group_id, intraday_position_summary_mv_next.currency_code, intraday_position_summary_mv_next.position_type ], stream key: [ party_involvements_dm.party_id, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, lifecycle_profiles.id, min(party_involvements_dm.customer_relationship_id), portfolios_plain_mv.portfolio_id, position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded, position_snapshot_mv_next.flag, portfolio_to_account_groups_mv.account_group_id, portfolios_plain_mv.base_currency_code, intraday_position_summary_mv_next.position_type ] }
    └── MergeExecutor { output: [ pnl_snapshot_mv_next.account_group_id, pnl_snapshot_mv_next.currency_code, pnl_snapshot_mv_next.unrealized_gain_loss, pnl_snapshot_mv_next.total_average_cost, pnl_snapshot_mv_next.position_type ], stream key: [ pnl_snapshot_mv_next.account_group_id, pnl_snapshot_mv_next.position_type, pnl_snapshot_mv_next.currency_code ] }

Fragment 26089 (Actor 132839,132838)
StreamSyncLogStore { output: [ party_involvements_dm.party_id, portfolios_plain_mv.portfolio_id, portfolios_plain_mv.name, portfolios_plain_mv.number, portfolios_plain_mv.base_currency_code, service_types_dm.type, lifecycle_profiles.base_currency_code, portfolio_to_account_groups_mv.account_group_id, position_snapshot_mv_next.market_value, position_snapshot_mv_next.market_value_system_currency, intraday_position_summary_mv_next.market_value, intraday_position_summary_mv_next.total_average_cost, intraday_position_summary_mv_next.market_value_system_currency, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, lifecycle_profiles.id, min(party_involvements_dm.customer_relationship_id), position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded, position_snapshot_mv_next.flag, intraday_position_summary_mv_next.account_group_id, intraday_position_summary_mv_next.currency_code, intraday_position_summary_mv_next.position_type ], stream key: [ party_involvements_dm.party_id, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, lifecycle_profiles.id, min(party_involvements_dm.customer_relationship_id), portfolios_plain_mv.portfolio_id, position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded, position_snapshot_mv_next.flag, portfolio_to_account_groups_mv.account_group_id, portfolios_plain_mv.base_currency_code, intraday_position_summary_mv_next.position_type ] }
└── StreamHashJoin { type: LeftOuter, predicate: portfolio_to_account_groups_mv.account_group_id = intraday_position_summary_mv_next.account_group_id AND portfolios_plain_mv.base_currency_code = intraday_position_summary_mv_next.currency_code } { output: [ party_involvements_dm.party_id, portfolios_plain_mv.portfolio_id, portfolios_plain_mv.name, portfolios_plain_mv.number, portfolios_plain_mv.base_currency_code, service_types_dm.type, lifecycle_profiles.base_currency_code, portfolio_to_account_groups_mv.account_group_id, position_snapshot_mv_next.market_value, position_snapshot_mv_next.market_value_system_currency, intraday_position_summary_mv_next.market_value, intraday_position_summary_mv_next.total_average_cost, intraday_position_summary_mv_next.market_value_system_currency, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, lifecycle_profiles.id, min(party_involvements_dm.customer_relationship_id), position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded, position_snapshot_mv_next.flag, intraday_position_summary_mv_next.account_group_id, intraday_position_summary_mv_next.currency_code, intraday_position_summary_mv_next.position_type ], stream key: [ party_involvements_dm.party_id, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, lifecycle_profiles.id, min(party_involvements_dm.customer_relationship_id), portfolios_plain_mv.portfolio_id, position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded, position_snapshot_mv_next.flag, portfolio_to_account_groups_mv.account_group_id, portfolios_plain_mv.base_currency_code, intraday_position_summary_mv_next.position_type ] }
    ├── MergeExecutor { output: [ party_involvements_dm.party_id, portfolios_plain_mv.portfolio_id, portfolios_plain_mv.name, portfolios_plain_mv.number, portfolios_plain_mv.base_currency_code, service_types_dm.type, lifecycle_profiles.base_currency_code, portfolio_to_account_groups_mv.account_group_id, position_snapshot_mv_next.market_value, position_snapshot_mv_next.market_value_system_currency, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, lifecycle_profiles.id, min(party_involvements_dm.customer_relationship_id), position_snapshot_mv_next.account_group_id, position_snapshot_mv_next.currency_code, position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded, position_snapshot_mv_next.flag ], stream key: [ party_involvements_dm.party_id, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, lifecycle_profiles.id, min(party_involvements_dm.customer_relationship_id), portfolios_plain_mv.portfolio_id, position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded, position_snapshot_mv_next.flag, portfolio_to_account_groups_mv.account_group_id, portfolios_plain_mv.base_currency_code ] }
    └── MergeExecutor { output: [ intraday_position_summary_mv_next.account_group_id, intraday_position_summary_mv_next.currency_code, intraday_position_summary_mv_next.market_value, intraday_position_summary_mv_next.total_average_cost, intraday_position_summary_mv_next.market_value_system_currency, intraday_position_summary_mv_next.position_type ], stream key: [ intraday_position_summary_mv_next.account_group_id, intraday_position_summary_mv_next.currency_code, intraday_position_summary_mv_next.position_type ] }

Fragment 26090 (Actor 132841,132840)
StreamSyncLogStore { output: [ party_involvements_dm.party_id, portfolios_plain_mv.portfolio_id, portfolios_plain_mv.name, portfolios_plain_mv.number, portfolios_plain_mv.base_currency_code, service_types_dm.type, lifecycle_profiles.base_currency_code, portfolio_to_account_groups_mv.account_group_id, position_snapshot_mv_next.market_value, position_snapshot_mv_next.market_value_system_currency, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, lifecycle_profiles.id, min(party_involvements_dm.customer_relationship_id), position_snapshot_mv_next.account_group_id, position_snapshot_mv_next.currency_code, position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded, position_snapshot_mv_next.flag ], stream key: [ party_involvements_dm.party_id, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, lifecycle_profiles.id, min(party_involvements_dm.customer_relationship_id), portfolios_plain_mv.portfolio_id, position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded, position_snapshot_mv_next.flag, portfolio_to_account_groups_mv.account_group_id, portfolios_plain_mv.base_currency_code ] }
└── StreamHashJoin { type: LeftOuter, predicate: portfolio_to_account_groups_mv.account_group_id = position_snapshot_mv_next.account_group_id AND portfolios_plain_mv.base_currency_code = position_snapshot_mv_next.currency_code } { output: [ party_involvements_dm.party_id, portfolios_plain_mv.portfolio_id, portfolios_plain_mv.name, portfolios_plain_mv.number, portfolios_plain_mv.base_currency_code, service_types_dm.type, lifecycle_profiles.base_currency_code, portfolio_to_account_groups_mv.account_group_id, position_snapshot_mv_next.market_value, position_snapshot_mv_next.market_value_system_currency, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, lifecycle_profiles.id, min(party_involvements_dm.customer_relationship_id), position_snapshot_mv_next.account_group_id, position_snapshot_mv_next.currency_code, position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded, position_snapshot_mv_next.flag ], stream key: [ party_involvements_dm.party_id, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, lifecycle_profiles.id, min(party_involvements_dm.customer_relationship_id), portfolios_plain_mv.portfolio_id, position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded, position_snapshot_mv_next.flag, portfolio_to_account_groups_mv.account_group_id, portfolios_plain_mv.base_currency_code ] }
    ├── MergeExecutor { output: [ party_involvements_dm.party_id, portfolios_plain_mv.portfolio_id, portfolios_plain_mv.name, portfolios_plain_mv.number, portfolios_plain_mv.base_currency_code, service_types_dm.type, lifecycle_profiles.base_currency_code, portfolio_to_account_groups_mv.account_group_id, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, lifecycle_profiles.id, min(party_involvements_dm.customer_relationship_id), portfolio_to_account_groups_mv.portfolio_id ], stream key: [ party_involvements_dm.party_id, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, lifecycle_profiles.id, min(party_involvements_dm.customer_relationship_id), portfolios_plain_mv.portfolio_id ] }
    └── MergeExecutor { output: [ position_snapshot_mv_next.account_group_id, position_snapshot_mv_next.currency_code, position_snapshot_mv_next.market_value, position_snapshot_mv_next.market_value_system_currency, position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded, position_snapshot_mv_next.flag ], stream key: [ position_snapshot_mv_next.account_group_id, position_snapshot_mv_next.currency_code, position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded, position_snapshot_mv_next.flag ] }

Fragment 26091 (Actor 132844,132845)
StreamSyncLogStore { output: [ party_involvements_dm.party_id, portfolios_plain_mv.portfolio_id, portfolios_plain_mv.name, portfolios_plain_mv.number, portfolios_plain_mv.base_currency_code, service_types_dm.type, lifecycle_profiles.base_currency_code, portfolio_to_account_groups_mv.account_group_id, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, lifecycle_profiles.id, min(party_involvements_dm.customer_relationship_id), portfolio_to_account_groups_mv.portfolio_id ], stream key: [ party_involvements_dm.party_id, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, lifecycle_profiles.id, min(party_involvements_dm.customer_relationship_id), portfolios_plain_mv.portfolio_id ] }
└── StreamHashJoin { type: LeftOuter, predicate: portfolios_plain_mv.portfolio_id = portfolio_to_account_groups_mv.portfolio_id } { output: [ party_involvements_dm.party_id, portfolios_plain_mv.portfolio_id, portfolios_plain_mv.name, portfolios_plain_mv.number, portfolios_plain_mv.base_currency_code, service_types_dm.type, lifecycle_profiles.base_currency_code, portfolio_to_account_groups_mv.account_group_id, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, lifecycle_profiles.id, min(party_involvements_dm.customer_relationship_id), portfolio_to_account_groups_mv.portfolio_id ], stream key: [ party_involvements_dm.party_id, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, lifecycle_profiles.id, min(party_involvements_dm.customer_relationship_id), portfolios_plain_mv.portfolio_id ] }
    ├── MergeExecutor { output: [ party_involvements_dm.party_id, portfolios_plain_mv.portfolio_id, portfolios_plain_mv.name, portfolios_plain_mv.number, portfolios_plain_mv.base_currency_code, service_types_dm.type, lifecycle_profiles.base_currency_code, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, min(party_involvements_dm.customer_relationship_id), lifecycle_profiles.id ], stream key: [ party_involvements_dm.party_id, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, lifecycle_profiles.id, min(party_involvements_dm.customer_relationship_id) ] }
    └── MergeExecutor { output: [ portfolio_to_account_groups_mv.portfolio_id, portfolio_to_account_groups_mv.account_group_id ], stream key: [ portfolio_to_account_groups_mv.portfolio_id ] }

Fragment 26092 (Actor 132846,132847)
StreamSyncLogStore { output: [ party_involvements_dm.party_id, portfolios_plain_mv.portfolio_id, portfolios_plain_mv.name, portfolios_plain_mv.number, portfolios_plain_mv.base_currency_code, service_types_dm.type, lifecycle_profiles.base_currency_code, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, min(party_involvements_dm.customer_relationship_id), lifecycle_profiles.id ], stream key: [ party_involvements_dm.party_id, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, lifecycle_profiles.id, min(party_involvements_dm.customer_relationship_id) ] }
└── StreamHashJoin { type: LeftOuter, predicate: min(party_involvements_dm.customer_relationship_id) = lifecycle_profiles.customer_relationship_id } { output: [ party_involvements_dm.party_id, portfolios_plain_mv.portfolio_id, portfolios_plain_mv.name, portfolios_plain_mv.number, portfolios_plain_mv.base_currency_code, service_types_dm.type, lifecycle_profiles.base_currency_code, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, min(party_involvements_dm.customer_relationship_id), lifecycle_profiles.id ], stream key: [ party_involvements_dm.party_id, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, lifecycle_profiles.id, min(party_involvements_dm.customer_relationship_id) ] }
    ├── MergeExecutor { output: [ party_involvements_dm.party_id, min(party_involvements_dm.customer_relationship_id), portfolios_plain_mv.portfolio_id, portfolios_plain_mv.name, portfolios_plain_mv.number, portfolios_plain_mv.base_currency_code, service_types_dm.type, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, service_types_dm.service_type_id ], stream key: [ party_involvements_dm.party_id, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id ] }
    └── MergeExecutor { output: [ lifecycle_profiles.customer_relationship_id, lifecycle_profiles.base_currency_code, lifecycle_profiles.id ], stream key: [ lifecycle_profiles.id ] }

Fragment 26093 (Actor 132702,132703)
StreamTemporalJoin { type: Inner, append_only: false, predicate: portfolios_plain_mv.service_type_id = service_types_dm.service_type_id, nested_loop: false } { output: [ party_involvements_dm.party_id, min(party_involvements_dm.customer_relationship_id), portfolios_plain_mv.portfolio_id, portfolios_plain_mv.name, portfolios_plain_mv.number, portfolios_plain_mv.base_currency_code, service_types_dm.type, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id, service_types_dm.service_type_id ], stream key: [ party_involvements_dm.party_id, party_involvements_dm.entity_id, portfolios_plain_mv.service_type_id ] }
├── MergeExecutor { output: [ party_involvements_dm.party_id, min(party_involvements_dm.customer_relationship_id), portfolios_plain_mv.portfolio_id, portfolios_plain_mv.name, portfolios_plain_mv.number, portfolios_plain_mv.base_currency_code, portfolios_plain_mv.service_type_id, party_involvements_dm.entity_id ], stream key: [ party_involvements_dm.party_id, party_involvements_dm.entity_id ] }
└── MergeExecutor { output: [ service_types_dm.service_type_id, service_types_dm.type ], stream key: [ service_types_dm.service_type_id ] }

Fragment 26094 (Actor 132848,132849)
StreamSyncLogStore { output: [ party_involvements_dm.party_id, min(party_involvements_dm.customer_relationship_id), portfolios_plain_mv.portfolio_id, portfolios_plain_mv.name, portfolios_plain_mv.number, portfolios_plain_mv.base_currency_code, portfolios_plain_mv.service_type_id, party_involvements_dm.entity_id ], stream key: [ party_involvements_dm.party_id, party_involvements_dm.entity_id ] }
└── StreamHashJoin { type: Inner, predicate: party_involvements_dm.entity_id = portfolios_plain_mv.portfolio_id } { output: [ party_involvements_dm.party_id, min(party_involvements_dm.customer_relationship_id), portfolios_plain_mv.portfolio_id, portfolios_plain_mv.name, portfolios_plain_mv.number, portfolios_plain_mv.base_currency_code, portfolios_plain_mv.service_type_id, party_involvements_dm.entity_id ], stream key: [ party_involvements_dm.party_id, party_involvements_dm.entity_id ] }
    ├── MergeExecutor { output: [ party_involvements_dm.party_id, party_involvements_dm.entity_id, min(party_involvements_dm.customer_relationship_id) ], stream key: [ party_involvements_dm.party_id, party_involvements_dm.entity_id ] }
    └── MergeExecutor { output: [ portfolios_plain_mv.portfolio_id, portfolios_plain_mv.name, portfolios_plain_mv.number, portfolios_plain_mv.base_currency_code, portfolios_plain_mv.service_type_id ], stream key: [ portfolios_plain_mv.portfolio_id ] }

Fragment 26095 (Actor 132850,132851)
StreamProject { exprs: [party_involvements_dm.party_id, party_involvements_dm.entity_id, min(party_involvements_dm.customer_relationship_id)] } { output: [ party_involvements_dm.party_id, party_involvements_dm.entity_id, min(party_involvements_dm.customer_relationship_id) ], stream key: [ party_involvements_dm.party_id, party_involvements_dm.entity_id ] }
└── StreamHashAgg { group_key: [party_involvements_dm.party_id, party_involvements_dm.entity_id], aggs: [min(party_involvements_dm.customer_relationship_id), count] } { output: [ party_involvements_dm.party_id, party_involvements_dm.entity_id, min(party_involvements_dm.customer_relationship_id), count ], stream key: [ party_involvements_dm.party_id, party_involvements_dm.entity_id ] }
    └── MergeExecutor { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.id, $src, customer_relationships_next.id ], stream key: [ party_involvements_dm.id, $src, party_involvements_dm.customer_relationship_id ] }

Fragment 26096 (Actor 132852,132853)
StreamSyncLogStore { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.id, $src, customer_relationships_next.id ], stream key: [ party_involvements_dm.id, $src, party_involvements_dm.customer_relationship_id ] }
└── StreamHashJoin { type: Inner, predicate: party_involvements_dm.customer_relationship_id = customer_relationships_next.id } { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.id, $src, customer_relationships_next.id ], stream key: [ party_involvements_dm.id, $src, party_involvements_dm.customer_relationship_id ] }
    ├── MergeExecutor { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.id, $src ], stream key: [ party_involvements_dm.id, $src ] }
    └── MergeExecutor { output: [ customer_relationships_next.id ], stream key: [ customer_relationships_next.id ] }

Fragment 26097 (Actor 132855,132854)
StreamUnion { all: true } { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.id, $src ], stream key: [ party_involvements_dm.id, $src ] }
├── MergeExecutor { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.id, 0:Int32 ], stream key: [ party_involvements_dm.id ] }
└── MergeExecutor { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.id, 1:Int32 ], stream key: [ party_involvements_dm.id ] }

Fragment 26098 (Actor 132856,132857)
StreamProject { exprs: [party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.id, 0:Int32] } { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.id, 0:Int32 ], stream key: [ party_involvements_dm.id ] }
└── StreamDynamicFilter { predicate: ($expr2 > now), output_watermarks: [[$expr2]], output: [party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, $expr2, party_involvements_dm.id], cleaned_by_watermark: true } { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, $expr2, party_involvements_dm.id ], stream key: [ party_involvements_dm.id ] }
    ├── StreamProject { exprs: [party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, AtTimeZone(party_involvements_dm.effective_to::Timestamp, 'UTC':Varchar) as $expr2, party_involvements_dm.id] } { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, $expr2, party_involvements_dm.id ], stream key: [ party_involvements_dm.id ] }
    │   └── StreamFilter { predicate: IsNotTrue(IsNull(party_involvements_dm.effective_to)) } { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.effective_to, party_involvements_dm.id ], stream key: [ party_involvements_dm.id ] }
    │       └── MergeExecutor { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.effective_to, party_involvements_dm.id ], stream key: [ party_involvements_dm.id ] }
    └── MergeExecutor { output: [ now ], stream key: [] }

Fragment 26099 (Actor 132859,132858)
StreamProject { exprs: [party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.effective_to, party_involvements_dm.id] } { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.effective_to, party_involvements_dm.id ], stream key: [ party_involvements_dm.id ] }
└── StreamDynamicFilter { predicate: ($expr1 <= now), output: [party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.effective_to, $expr1, party_involvements_dm.id], cleaned_by_watermark: true } { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.effective_to, $expr1, party_involvements_dm.id ], stream key: [ party_involvements_dm.id ] }
    ├── StreamProject { exprs: [party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.effective_to, AtTimeZone(party_involvements_dm.effective_from::Timestamp, 'UTC':Varchar) as $expr1, party_involvements_dm.id] } { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.effective_to, $expr1, party_involvements_dm.id ], stream key: [ party_involvements_dm.id ] }
    │   └── StreamFilter { predicate: (IsNotTrue(IsNull(party_involvements_dm.effective_to)) OR IsNull(party_involvements_dm.effective_to)) AND (party_involvements_dm.entity_type = 'PORTFOLIO':Varchar) AND (party_involvements_dm.status = 'ACTIVE':Varchar) AND IsNull(party_involvements_dm.disabled_at) } { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.effective_from, party_involvements_dm.effective_to, party_involvements_dm.id, party_involvements_dm.entity_type, party_involvements_dm.status, party_involvements_dm.disabled_at ], stream key: [ party_involvements_dm.id ] }
    │       └── StreamTableScan { table: party_involvements_dm, columns: [party_id, customer_relationship_id, entity_id, effective_from, effective_to, id, entity_type, status, disabled_at] } { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.effective_from, party_involvements_dm.effective_to, party_involvements_dm.id, party_involvements_dm.entity_type, party_involvements_dm.status, party_involvements_dm.disabled_at ], stream key: [ party_involvements_dm.id ] }
    │           ├── Upstream { output: [ party_id, customer_relationship_id, entity_id, effective_from, effective_to, id, entity_type, status, disabled_at ], stream key: [] }
    │           └── BatchPlanNode { output: [ party_id, customer_relationship_id, entity_id, effective_from, effective_to, id, entity_type, status, disabled_at ], stream key: [] }
    └── MergeExecutor { output: [ now ], stream key: [] }

Fragment 26100 (Actor 132862)
StreamNow { output: [ now ], stream key: [] }

Fragment 26101 (Actor 132863)
StreamNow { output: [ now ], stream key: [] }

Fragment 26102 (Actor 132860,132861)
StreamProject { exprs: [party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.id, 1:Int32] } { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.id, 1:Int32 ], stream key: [ party_involvements_dm.id ] }
└── StreamFilter { predicate: IsNull(party_involvements_dm.effective_to) } { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.effective_to, party_involvements_dm.id ], stream key: [ party_involvements_dm.id ] }
    └── MergeExecutor { output: [ party_involvements_dm.party_id, party_involvements_dm.customer_relationship_id, party_involvements_dm.entity_id, party_involvements_dm.effective_to, party_involvements_dm.id ], stream key: [ party_involvements_dm.id ] }

Fragment 26103 (Actor 132865,132864)
StreamProject { exprs: [customer_relationships_next.id] } { output: [ customer_relationships_next.id ], stream key: [ customer_relationships_next.id ] }
└── StreamFilter { predicate: (customer_relationships_next.type = 'CUSTOMER':Varchar) AND (customer_relationships_next.status = 'ACTIVE':Varchar) AND IsNull(customer_relationships_next.disabled_at) } { output: [ customer_relationships_next.id, customer_relationships_next.type, customer_relationships_next.status, customer_relationships_next.disabled_at ], stream key: [ customer_relationships_next.id ] }
    └── StreamTableScan { table: customer_relationships_next, columns: [id, type, status, disabled_at] } { output: [ customer_relationships_next.id, customer_relationships_next.type, customer_relationships_next.status, customer_relationships_next.disabled_at ], stream key: [ customer_relationships_next.id ] }
        ├── Upstream { output: [ id, type, status, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ id, type, status, disabled_at ], stream key: [] }

Fragment 26104 (Actor 132866,132867)
StreamProject { exprs: [portfolios_plain_mv.portfolio_id, portfolios_plain_mv.name, portfolios_plain_mv.number, portfolios_plain_mv.base_currency_code, portfolios_plain_mv.service_type_id] } { output: [ portfolios_plain_mv.portfolio_id, portfolios_plain_mv.name, portfolios_plain_mv.number, portfolios_plain_mv.base_currency_code, portfolios_plain_mv.service_type_id ], stream key: [ portfolios_plain_mv.portfolio_id ] }
└── StreamFilter { predicate: IsNull(portfolios_plain_mv.closing_date) } { output: [ portfolios_plain_mv.portfolio_id, portfolios_plain_mv.name, portfolios_plain_mv.number, portfolios_plain_mv.base_currency_code, portfolios_plain_mv.service_type_id, portfolios_plain_mv.closing_date ], stream key: [ portfolios_plain_mv.portfolio_id ] }
    └── StreamTableScan { table: portfolios_plain_mv, columns: [portfolio_id, name, number, base_currency_code, service_type_id, closing_date] } { output: [ portfolios_plain_mv.portfolio_id, portfolios_plain_mv.name, portfolios_plain_mv.number, portfolios_plain_mv.base_currency_code, portfolios_plain_mv.service_type_id, portfolios_plain_mv.closing_date ], stream key: [ portfolios_plain_mv.portfolio_id ] }
        ├── Upstream { output: [ portfolio_id, name, number, base_currency_code, service_type_id, closing_date ], stream key: [] }
        └── BatchPlanNode { output: [ portfolio_id, name, number, base_currency_code, service_type_id, closing_date ], stream key: [] }

Fragment 26105 (Actor 132704,132705)
StreamTableScan { table: service_types_dm, columns: [service_type_id, type] } { output: [ service_types_dm.service_type_id, service_types_dm.type ], stream key: [ service_types_dm.service_type_id ] }
├── Upstream { output: [ service_type_id, type ], stream key: [] }
└── BatchPlanNode { output: [ service_type_id, type ], stream key: [] }

Fragment 26106 (Actor 132868,132869)
StreamProject { exprs: [lifecycle_profiles.customer_relationship_id, lifecycle_profiles.base_currency_code, lifecycle_profiles.id] } { output: [ lifecycle_profiles.customer_relationship_id, lifecycle_profiles.base_currency_code, lifecycle_profiles.id ], stream key: [ lifecycle_profiles.id ] }
└── StreamFilter { predicate: IsNull(lifecycle_profiles.disabled_at) } { output: [ lifecycle_profiles.customer_relationship_id, lifecycle_profiles.base_currency_code, lifecycle_profiles.id, lifecycle_profiles.disabled_at ], stream key: [ lifecycle_profiles.id ] }
    └── StreamTableScan { table: lifecycle_profiles, columns: [customer_relationship_id, base_currency_code, id, disabled_at] } { output: [ lifecycle_profiles.customer_relationship_id, lifecycle_profiles.base_currency_code, lifecycle_profiles.id, lifecycle_profiles.disabled_at ], stream key: [ lifecycle_profiles.id ] }
        ├── Upstream { output: [ customer_relationship_id, base_currency_code, id, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ customer_relationship_id, base_currency_code, id, disabled_at ], stream key: [] }

Fragment 26107 (Actor 132870,132871)
StreamProject { exprs: [portfolio_to_account_groups_mv.portfolio_id, portfolio_to_account_groups_mv.account_group_id] } { output: [ portfolio_to_account_groups_mv.portfolio_id, portfolio_to_account_groups_mv.account_group_id ], stream key: [ portfolio_to_account_groups_mv.portfolio_id ] }
└── 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.type ], stream key: [ portfolio_to_account_groups_mv.portfolio_id ] }
    └── StreamTableScan { table: portfolio_to_account_groups_mv, columns: [portfolio_id, account_group_id, type] } { output: [ portfolio_to_account_groups_mv.portfolio_id, portfolio_to_account_groups_mv.account_group_id, portfolio_to_account_groups_mv.type ], stream key: [ portfolio_to_account_groups_mv.portfolio_id ] }
        ├── Upstream { output: [ portfolio_id, account_group_id, type ], stream key: [] }
        └── BatchPlanNode { output: [ portfolio_id, account_group_id, type ], stream key: [] }

Fragment 26108 (Actor 132873,132872)
StreamProject { exprs: [position_snapshot_mv_next.account_group_id, position_snapshot_mv_next.currency_code, position_snapshot_mv_next.market_value, position_snapshot_mv_next.market_value_system_currency, position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded, position_snapshot_mv_next.flag] } { output: [ position_snapshot_mv_next.account_group_id, position_snapshot_mv_next.currency_code, position_snapshot_mv_next.market_value, position_snapshot_mv_next.market_value_system_currency, position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded, position_snapshot_mv_next.flag ], stream key: [ position_snapshot_mv_next.account_group_id, position_snapshot_mv_next.currency_code, position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded, position_snapshot_mv_next.flag ] }
└── StreamFilter { predicate: (position_snapshot_mv_next.position_type = 'POSITION':Varchar) } { output: [ position_snapshot_mv_next.account_group_id, position_snapshot_mv_next.currency_code, position_snapshot_mv_next.market_value, position_snapshot_mv_next.market_value_system_currency, position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded, position_snapshot_mv_next.flag, position_snapshot_mv_next.position_type ], stream key: [ position_snapshot_mv_next.account_group_id, position_snapshot_mv_next.currency_code, position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded, position_snapshot_mv_next.flag ] }
    └── StreamTableScan { table: position_snapshot_mv_next, columns: [account_group_id, currency_code, market_value, market_value_system_currency, holding_values_latest_mv_next.type_expanded, flag, position_type] } { output: [ position_snapshot_mv_next.account_group_id, position_snapshot_mv_next.currency_code, position_snapshot_mv_next.market_value, position_snapshot_mv_next.market_value_system_currency, position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded, position_snapshot_mv_next.flag, position_snapshot_mv_next.position_type ], stream key: [ position_snapshot_mv_next.account_group_id, position_snapshot_mv_next.currency_code, position_snapshot_mv_next.holding_values_latest_mv_next.type_expanded, position_snapshot_mv_next.flag ] }
        ├── Upstream { output: [ account_group_id, currency_code, market_value, market_value_system_currency, holding_values_latest_mv_next.type_expanded, flag, position_type ], stream key: [] }
        └── BatchPlanNode { output: [ account_group_id, currency_code, market_value, market_value_system_currency, holding_values_latest_mv_next.type_expanded, flag, position_type ], stream key: [] }

Fragment 26109 (Actor 132875,132874)
StreamFilter { predicate: (intraday_position_summary_mv_next.position_type = 'POSITION':Varchar) } { output: [ intraday_position_summary_mv_next.account_group_id, intraday_position_summary_mv_next.currency_code, intraday_position_summary_mv_next.market_value, intraday_position_summary_mv_next.total_average_cost, intraday_position_summary_mv_next.market_value_system_currency, intraday_position_summary_mv_next.position_type ], stream key: [ intraday_position_summary_mv_next.account_group_id, intraday_position_summary_mv_next.currency_code, intraday_position_summary_mv_next.position_type ] }
└── StreamTableScan { table: intraday_position_summary_mv_next, columns: [account_group_id, currency_code, market_value, total_average_cost, market_value_system_currency, position_type] } { output: [ intraday_position_summary_mv_next.account_group_id, intraday_position_summary_mv_next.currency_code, intraday_position_summary_mv_next.market_value, intraday_position_summary_mv_next.total_average_cost, intraday_position_summary_mv_next.market_value_system_currency, intraday_position_summary_mv_next.position_type ], stream key: [ intraday_position_summary_mv_next.account_group_id, intraday_position_summary_mv_next.currency_code, intraday_position_summary_mv_next.position_type ] }
    ├── Upstream { output: [ account_group_id, currency_code, market_value, total_average_cost, market_value_system_currency, position_type ], stream key: [] }
    └── BatchPlanNode { output: [ account_group_id, currency_code, market_value, total_average_cost, market_value_system_currency, position_type ], stream key: [] }

Fragment 26110 (Actor 132876,132877)
StreamFilter { predicate: (pnl_snapshot_mv_next.position_type = 'POSITION':Varchar) } { output: [ pnl_snapshot_mv_next.account_group_id, pnl_snapshot_mv_next.currency_code, pnl_snapshot_mv_next.unrealized_gain_loss, pnl_snapshot_mv_next.total_average_cost, pnl_snapshot_mv_next.position_type ], stream key: [ pnl_snapshot_mv_next.account_group_id, pnl_snapshot_mv_next.position_type, pnl_snapshot_mv_next.currency_code ] }
└── StreamTableScan { table: pnl_snapshot_mv_next, columns: [account_group_id, currency_code, unrealized_gain_loss, total_average_cost, position_type] } { output: [ pnl_snapshot_mv_next.account_group_id, pnl_snapshot_mv_next.currency_code, pnl_snapshot_mv_next.unrealized_gain_loss, pnl_snapshot_mv_next.total_average_cost, pnl_snapshot_mv_next.position_type ], stream key: [ pnl_snapshot_mv_next.account_group_id, pnl_snapshot_mv_next.position_type, pnl_snapshot_mv_next.currency_code ] }
    ├── Upstream { output: [ account_group_id, currency_code, unrealized_gain_loss, total_average_cost, position_type ], stream key: [] }
    └── BatchPlanNode { output: [ account_group_id, currency_code, unrealized_gain_loss, total_average_cost, position_type ], stream key: [] }