Fragment 24598 (Actor 116172,116171)
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.holding_values_latest_mv_next.type_expanded
│ ├── position_snapshot_mv.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.holding_values_latest_mv_next.type_expanded
├── position_snapshot_mv.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 24599 (Actor 116175,116176)
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.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.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.holding_values_latest_mv_next.type_expanded
├─position_snapshot_mv.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.holding_values_latest_mv_next.type_expanded, position_snapshot_mv.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.holding_values_latest_mv_next.type_expanded, position_snapshot_mv.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.market_value, position_snapshot_mv.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.holding_values_latest_mv_next.type_expanded, position_snapshot_mv.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.holding_values_latest_mv_next.type_expanded, position_snapshot_mv.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 24600 (Actor 116179,116180)
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.market_value, position_snapshot_mv.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.holding_values_latest_mv_next.type_expanded, position_snapshot_mv.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.holding_values_latest_mv_next.type_expanded, position_snapshot_mv.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.market_value, position_snapshot_mv.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.holding_values_latest_mv_next.type_expanded, position_snapshot_mv.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.holding_values_latest_mv_next.type_expanded, position_snapshot_mv.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.market_value, position_snapshot_mv.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.holding_values_latest_mv_next.type_expanded, position_snapshot_mv.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.holding_values_latest_mv_next.type_expanded, position_snapshot_mv.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 24601 (Actor 116174,116173)
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.market_value, position_snapshot_mv.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.holding_values_latest_mv_next.type_expanded, position_snapshot_mv.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.holding_values_latest_mv_next.type_expanded, position_snapshot_mv.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.market_value, position_snapshot_mv.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.holding_values_latest_mv_next.type_expanded, position_snapshot_mv.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.holding_values_latest_mv_next.type_expanded, position_snapshot_mv.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.market_value, position_snapshot_mv.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.account_group_id, position_snapshot_mv.currency_code, position_snapshot_mv.holding_values_latest_mv_next.type_expanded, position_snapshot_mv.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.holding_values_latest_mv_next.type_expanded, position_snapshot_mv.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 24602 (Actor 116178,116177)
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.market_value, position_snapshot_mv.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.account_group_id, position_snapshot_mv.currency_code, position_snapshot_mv.holding_values_latest_mv_next.type_expanded, position_snapshot_mv.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.holding_values_latest_mv_next.type_expanded, position_snapshot_mv.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.account_group_id AND portfolios_plain_mv.base_currency_code = position_snapshot_mv.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.market_value, position_snapshot_mv.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.account_group_id, position_snapshot_mv.currency_code, position_snapshot_mv.holding_values_latest_mv_next.type_expanded, position_snapshot_mv.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.holding_values_latest_mv_next.type_expanded, position_snapshot_mv.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.account_group_id, position_snapshot_mv.currency_code, position_snapshot_mv.market_value, position_snapshot_mv.market_value_system_currency, position_snapshot_mv.holding_values_latest_mv_next.type_expanded, position_snapshot_mv.flag ], stream key: [ position_snapshot_mv.account_group_id, position_snapshot_mv.currency_code, position_snapshot_mv.holding_values_latest_mv_next.type_expanded, position_snapshot_mv.flag ] }
Fragment 24603 (Actor 116182,116181)
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 24604 (Actor 116184,116183)
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 24605 (Actor 116041,116042)
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 24606 (Actor 116186,116185)
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 24607 (Actor 116188,116187)
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.id ], stream key: [ party_involvements_dm.id, $src, party_involvements_dm.customer_relationship_id ] }
Fragment 24608 (Actor 116189,116190)
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.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.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.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.id ], stream key: [ customer_relationships.id ] }
Fragment 24609 (Actor 116191,116192)
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 24610 (Actor 116194,116193)
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 24611 (Actor 116196,116195)
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 24612 (Actor 116199)
StreamNow { output: [ now ], stream key: [] }
Fragment 24613 (Actor 116200)
StreamNow { output: [ now ], stream key: [] }
Fragment 24614 (Actor 116198,116197)
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 24615 (Actor 116202,116201)
StreamProject { exprs: [customer_relationships.id] } { output: [ customer_relationships.id ], stream key: [ customer_relationships.id ] }
└── StreamFilter { predicate: (customer_relationships.type = 'CUSTOMER':Varchar) AND (customer_relationships.status = 'ACTIVE':Varchar) AND IsNull(customer_relationships.disabled_at) } { output: [ customer_relationships.id, customer_relationships.type, customer_relationships.status, customer_relationships.disabled_at ], stream key: [ customer_relationships.id ] }
└── StreamTableScan { table: customer_relationships, columns: [id, type, status, disabled_at] } { output: [ customer_relationships.id, customer_relationships.type, customer_relationships.status, customer_relationships.disabled_at ], stream key: [ customer_relationships.id ] }
├── Upstream { output: [ id, type, status, disabled_at ], stream key: [] }
└── BatchPlanNode { output: [ id, type, status, disabled_at ], stream key: [] }
Fragment 24616 (Actor 116203,116204)
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 24617 (Actor 116039,116040)
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 24618 (Actor 116205,116206)
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 24619 (Actor 116207,116208)
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 24620 (Actor 116210,116209)
StreamProject { exprs: [position_snapshot_mv.account_group_id, position_snapshot_mv.currency_code, position_snapshot_mv.market_value, position_snapshot_mv.market_value_system_currency, position_snapshot_mv.holding_values_latest_mv_next.type_expanded, position_snapshot_mv.flag] } { output: [ position_snapshot_mv.account_group_id, position_snapshot_mv.currency_code, position_snapshot_mv.market_value, position_snapshot_mv.market_value_system_currency, position_snapshot_mv.holding_values_latest_mv_next.type_expanded, position_snapshot_mv.flag ], stream key: [ position_snapshot_mv.account_group_id, position_snapshot_mv.currency_code, position_snapshot_mv.holding_values_latest_mv_next.type_expanded, position_snapshot_mv.flag ] }
└── StreamFilter { predicate: (position_snapshot_mv.position_type = 'POSITION':Varchar) } { output: [ position_snapshot_mv.account_group_id, position_snapshot_mv.currency_code, position_snapshot_mv.market_value, position_snapshot_mv.market_value_system_currency, position_snapshot_mv.holding_values_latest_mv_next.type_expanded, position_snapshot_mv.flag, position_snapshot_mv.position_type ], stream key: [ position_snapshot_mv.account_group_id, position_snapshot_mv.currency_code, position_snapshot_mv.holding_values_latest_mv_next.type_expanded, position_snapshot_mv.flag ] }
└── StreamTableScan { table: position_snapshot_mv, 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.account_group_id, position_snapshot_mv.currency_code, position_snapshot_mv.market_value, position_snapshot_mv.market_value_system_currency, position_snapshot_mv.holding_values_latest_mv_next.type_expanded, position_snapshot_mv.flag, position_snapshot_mv.position_type ], stream key: [ position_snapshot_mv.account_group_id, position_snapshot_mv.currency_code, position_snapshot_mv.holding_values_latest_mv_next.type_expanded, position_snapshot_mv.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 24621 (Actor 116211,116212)
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 24622 (Actor 116214,116213)
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: [] }