Fragment 24324 (Actor 112483,112484)
StreamMaterialize { columns: [party_id, cash_tiles], stream_key: [party_id], pk_columns: [party_id], pk_conflict: NoCheck }
├── output: [ party_account_direct_mv.party_id, jsonb_agg($expr4) ]
├── stream key: [ party_account_direct_mv.party_id ]
└── StreamProject { exprs: [party_account_direct_mv.party_id, jsonb_agg($expr4)] }
├── output: [ party_account_direct_mv.party_id, jsonb_agg($expr4) ]
├── stream key: [ party_account_direct_mv.party_id ]
└── StreamHashAgg { group_key: [party_account_direct_mv.party_id], aggs: [jsonb_agg($expr4), count] }
├── output: [ party_account_direct_mv.party_id, jsonb_agg($expr4), count ]
├── stream key: [ party_account_direct_mv.party_id ]
└── MergeExecutor
├── output:
│ ┌── party_account_direct_mv.party_id
│ ├── $expr4
│ ├── party_account_direct_mv.party_involvements_dm.id
│ ├── party_account_direct_mv.party_involvements_dm.entity_id
│ ├── party_account_direct_mv.party_involvements_dm.customer_relationship_id
│ ├── party_account_direct_mv.$src
│ ├── party_account_direct_mv.account_id
│ ├── accounts_plain_mv.base_currency_code
│ ├── $src
│ ├── accounts_plain_mv.account_id
│ └── accounts_plain_mv.product_type_id
└── stream key:
┌── party_account_direct_mv.party_involvements_dm.id
├── party_account_direct_mv.party_involvements_dm.entity_id
├── party_account_direct_mv.party_involvements_dm.customer_relationship_id
├── party_account_direct_mv.$src
├── party_account_direct_mv.account_id
├── accounts_plain_mv.base_currency_code
├── $src
├── accounts_plain_mv.account_id
└── accounts_plain_mv.product_type_id
Fragment 24325 (Actor 112487,112488)
StreamProject { exprs: [party_account_direct_mv.party_id, JsonbBuildObject('accountId':Varchar, accounts_plain_mv.account_id, 'name':Varchar, accounts_plain_mv.name, 'number':Varchar, accounts_plain_mv.number, 'value':Varchar, JsonbBuildObject('amount':Varchar, investment_account_balance_snapshot_journal_mv.available_balance::Varchar, 'currency':Varchar, $expr3), 'baseCurrency':Varchar, $expr3) as $expr4, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code, $src, accounts_plain_mv.account_id, accounts_plain_mv.product_type_id] }
├── output: [ party_account_direct_mv.party_id, $expr4, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code, $src, accounts_plain_mv.account_id, accounts_plain_mv.product_type_id ]
├── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code, $src, accounts_plain_mv.account_id, accounts_plain_mv.product_type_id ]
└── StreamProject { exprs: [party_account_direct_mv.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.base_currency_code, currencies_dm.symbol, investment_account_balance_snapshot_journal_mv.available_balance, JsonbBuildObject('code':Varchar, accounts_plain_mv.base_currency_code, 'symbol':Varchar, currencies_dm.symbol) as $expr3, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code, $src, accounts_plain_mv.product_type_id] }
├── output: [ party_account_direct_mv.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.base_currency_code, currencies_dm.symbol, investment_account_balance_snapshot_journal_mv.available_balance, $expr3, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code, $src, accounts_plain_mv.product_type_id ]
├── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code, $src, accounts_plain_mv.account_id, accounts_plain_mv.product_type_id ]
└── MergeExecutor
├── output: [ party_account_direct_mv.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.base_currency_code, currencies_dm.symbol, investment_account_balance_snapshot_journal_mv.available_balance, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code, $src, accounts_plain_mv.product_type_id ]
└── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code, $src, accounts_plain_mv.account_id, accounts_plain_mv.product_type_id ]
Fragment 24326 (Actor 112485,112486)
StreamSyncLogStore
├── output: [ party_account_direct_mv.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.base_currency_code, currencies_dm.symbol, investment_account_balance_snapshot_journal_mv.available_balance, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code, $src, accounts_plain_mv.product_type_id ]
├── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code, $src, accounts_plain_mv.account_id, accounts_plain_mv.product_type_id ]
└── StreamHashJoin { type: LeftSemi, predicate: accounts_plain_mv.product_type_id = cash_product_type_ids_mv.product_type_id }
├── output: [ party_account_direct_mv.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.base_currency_code, currencies_dm.symbol, investment_account_balance_snapshot_journal_mv.available_balance, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code, $src, accounts_plain_mv.product_type_id ]
├── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code, $src, accounts_plain_mv.account_id, accounts_plain_mv.product_type_id ]
├── MergeExecutor
│ ├── output: [ party_account_direct_mv.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, investment_account_balance_snapshot_journal_mv.available_balance, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code, $src, investment_account_balance_snapshot_journal_mv.account_id ]
│ └── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code, $src, accounts_plain_mv.account_id ]
└── MergeExecutor { output: [ cash_product_type_ids_mv.product_type_id, cash_product_type_ids_mv._row_id, cash_product_type_ids_mv.$src ], stream key: [ cash_product_type_ids_mv._row_id, cash_product_type_ids_mv.$src ] }
Fragment 24327 (Actor 112489,112490)
StreamSyncLogStore
├── output: [ party_account_direct_mv.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, investment_account_balance_snapshot_journal_mv.available_balance, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code, $src, investment_account_balance_snapshot_journal_mv.account_id ]
├── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code, $src, accounts_plain_mv.account_id ]
└── StreamHashJoin { type: LeftOuter, predicate: accounts_plain_mv.account_id = investment_account_balance_snapshot_journal_mv.account_id }
├── output: [ party_account_direct_mv.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, investment_account_balance_snapshot_journal_mv.available_balance, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code, $src, investment_account_balance_snapshot_journal_mv.account_id ]
├── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code, $src, accounts_plain_mv.account_id ]
├── MergeExecutor
│ ├── output: [ party_account_direct_mv.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code, party_account_direct_mv.$src, $src ]
│ └── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code, $src ]
└── MergeExecutor { output: [ investment_account_balance_snapshot_journal_mv.account_id, investment_account_balance_snapshot_journal_mv.available_balance ], stream key: [ investment_account_balance_snapshot_journal_mv.account_id ] }
Fragment 24328 (Actor 112491,112492)
StreamUnion { all: true }
├── output: [ party_account_direct_mv.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code, party_account_direct_mv.$src, $src ]
├── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code, $src ]
├── MergeExecutor
│ ├── output: [ party_account_direct_mv.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code, party_account_direct_mv.$src, 0:Int32 ]
│ └── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code ]
└── MergeExecutor
├── output: [ party_account_direct_mv.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code, party_account_direct_mv.$src, 1:Int32 ]
└── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code ]
Fragment 24329 (Actor 112419,112420)
StreamProject { exprs: [party_account_direct_mv.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code, party_account_direct_mv.$src, 0:Int32] }
├── output: [ party_account_direct_mv.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code, party_account_direct_mv.$src, 0:Int32 ]
├── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code ]
└── StreamDynamicFilter { predicate: ($expr2 > now), output_watermarks: [[$expr2]], output: [party_account_direct_mv.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, $expr2, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id], cleaned_by_watermark: true }
├── output: [ party_account_direct_mv.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, $expr2, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id ]
├── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code ]
├── StreamProject { exprs: [party_account_direct_mv.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, AtTimeZone(party_account_direct_mv.effective_end_date::Timestamp, 'UTC':Varchar) as $expr2, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id] }
│ ├── output: [ party_account_direct_mv.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, $expr2, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id ]
│ ├── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code ]
│ └── StreamFilter { predicate: IsNotTrue(IsNull(party_account_direct_mv.effective_end_date)) }
│ ├── output: [ party_account_direct_mv.party_id, party_account_direct_mv.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id ]
│ ├── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code ]
│ └── MergeExecutor
│ ├── output: [ party_account_direct_mv.party_id, party_account_direct_mv.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id ]
│ └── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code ]
└── MergeExecutor { output: [ now ], stream key: [] }
Fragment 24330 (Actor 112426,112425)
StreamProject { exprs: [party_account_direct_mv.party_id, party_account_direct_mv.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id] }
├── output: [ party_account_direct_mv.party_id, party_account_direct_mv.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id ]
├── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code ]
└── StreamDynamicFilter { predicate: ($expr1 <= now), output: [party_account_direct_mv.party_id, party_account_direct_mv.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, $expr1, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id], cleaned_by_watermark: true }
├── output: [ party_account_direct_mv.party_id, party_account_direct_mv.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, $expr1, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id ]
├── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code ]
├── StreamProject { exprs: [party_account_direct_mv.party_id, party_account_direct_mv.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, AtTimeZone(party_account_direct_mv.effective_start_date::Timestamp, 'UTC':Varchar) as $expr1, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id] }
│ ├── output: [ party_account_direct_mv.party_id, party_account_direct_mv.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, $expr1, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id ]
│ ├── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code ]
│ └── StreamTemporalJoin { type: LeftOuter, append_only: false, predicate: accounts_plain_mv.base_currency_code = currencies_dm.code, nested_loop: false }
│ ├── output: [ party_account_direct_mv.party_id, party_account_direct_mv.effective_start_date, party_account_direct_mv.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, currencies_dm.code ]
│ ├── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code ]
│ ├── MergeExecutor
│ │ ├── output: [ party_account_direct_mv.party_id, party_account_direct_mv.effective_start_date, party_account_direct_mv.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id ]
│ │ └── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id ]
│ └── MergeExecutor { output: [ currencies_dm.code, currencies_dm.symbol ], stream key: [ currencies_dm.code ] }
└── MergeExecutor { output: [ now ], stream key: [] }
Fragment 24331 (Actor 112493,112494)
StreamSyncLogStore
├── output: [ party_account_direct_mv.party_id, party_account_direct_mv.effective_start_date, party_account_direct_mv.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id ]
├── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id ]
└── MergeExecutor
├── output: [ party_account_direct_mv.party_id, party_account_direct_mv.effective_start_date, party_account_direct_mv.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id ]
└── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id ]
Fragment 24332 (Actor 112495,112496)
StreamSyncLogStore
├── output: [ party_account_direct_mv.party_id, party_account_direct_mv.effective_start_date, party_account_direct_mv.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id ]
├── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id ]
└── StreamHashJoin { type: Inner, predicate: party_account_direct_mv.account_id = accounts_plain_mv.account_id }
├── output: [ party_account_direct_mv.party_id, party_account_direct_mv.effective_start_date, party_account_direct_mv.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id ]
├── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id ]
├── MergeExecutor { output: [ party_account_direct_mv.party_id, party_account_direct_mv.account_id, party_account_direct_mv.effective_start_date, party_account_direct_mv.effective_end_date, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src ], stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src ] }
└── MergeExecutor { output: [ accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code ], stream key: [ accounts_plain_mv.account_id ] }
Fragment 24333 (Actor 112497,112498)
StreamProject { exprs: [party_account_direct_mv.party_id, party_account_direct_mv.account_id, party_account_direct_mv.effective_start_date, party_account_direct_mv.effective_end_date, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src] }
├── output: [ party_account_direct_mv.party_id, party_account_direct_mv.account_id, party_account_direct_mv.effective_start_date, party_account_direct_mv.effective_end_date, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src ]
├── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src ]
└── StreamFilter { predicate: (IsNotTrue(IsNull(party_account_direct_mv.effective_end_date)) OR IsNull(party_account_direct_mv.effective_end_date)) AND (party_account_direct_mv.type = 'all':Varchar) }
├── output: [ party_account_direct_mv.party_id, party_account_direct_mv.account_id, party_account_direct_mv.effective_start_date, party_account_direct_mv.effective_end_date, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.type ]
├── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src ]
└── StreamTableScan { table: party_account_direct_mv, columns: [party_id, account_id, effective_start_date, effective_end_date, party_involvements_dm.id, party_involvements_dm.entity_id, party_involvements_dm.customer_relationship_id, $src, type] }
├── output: [ party_account_direct_mv.party_id, party_account_direct_mv.account_id, party_account_direct_mv.effective_start_date, party_account_direct_mv.effective_end_date, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.type ]
├── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src ]
├── Upstream { output: [ party_id, account_id, effective_start_date, effective_end_date, party_involvements_dm.id, party_involvements_dm.entity_id, party_involvements_dm.customer_relationship_id, $src, type ], stream key: [] }
└── BatchPlanNode { output: [ party_id, account_id, effective_start_date, effective_end_date, party_involvements_dm.id, party_involvements_dm.entity_id, party_involvements_dm.customer_relationship_id, $src, type ], stream key: [] }
Fragment 24334 (Actor 112500,112499)
StreamTableScan { table: accounts_plain_mv, columns: [account_id, name, number, product_type_id, base_currency_code] } { output: [ accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code ], stream key: [ accounts_plain_mv.account_id ] }
├── Upstream { output: [ account_id, name, number, product_type_id, base_currency_code ], stream key: [] }
└── BatchPlanNode { output: [ account_id, name, number, product_type_id, base_currency_code ], stream key: [] }
Fragment 24335 (Actor 112423,112424)
StreamTableScan { table: currencies_dm, columns: [code, symbol] } { output: [ currencies_dm.code, currencies_dm.symbol ], stream key: [ currencies_dm.code ] }
├── Upstream { output: [ code, symbol ], stream key: [] }
└── BatchPlanNode { output: [ code, symbol ], stream key: [] }
Fragment 24336 (Actor 112501)
StreamNow { output: [ now ], stream key: [] }
Fragment 24337 (Actor 112502)
StreamNow { output: [ now ], stream key: [] }
Fragment 24338 (Actor 112422,112421)
StreamProject { exprs: [party_account_direct_mv.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code, party_account_direct_mv.$src, 1:Int32] }
├── output: [ party_account_direct_mv.party_id, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code, party_account_direct_mv.$src, 1:Int32 ]
├── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code ]
└── StreamFilter { predicate: IsNull(party_account_direct_mv.effective_end_date) }
├── output: [ party_account_direct_mv.party_id, party_account_direct_mv.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id ]
├── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code ]
└── MergeExecutor
├── output: [ party_account_direct_mv.party_id, party_account_direct_mv.effective_end_date, accounts_plain_mv.account_id, accounts_plain_mv.name, accounts_plain_mv.number, accounts_plain_mv.product_type_id, accounts_plain_mv.base_currency_code, currencies_dm.symbol, party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id ]
└── stream key: [ party_account_direct_mv.party_involvements_dm.id, party_account_direct_mv.party_involvements_dm.entity_id, party_account_direct_mv.party_involvements_dm.customer_relationship_id, party_account_direct_mv.$src, party_account_direct_mv.account_id, accounts_plain_mv.base_currency_code ]
Fragment 24339 (Actor 112504,112503)
StreamTableScan { table: investment_account_balance_snapshot_journal_mv, columns: [account_id, available_balance] } { output: [ investment_account_balance_snapshot_journal_mv.account_id, investment_account_balance_snapshot_journal_mv.available_balance ], stream key: [ investment_account_balance_snapshot_journal_mv.account_id ] }
├── Upstream { output: [ account_id, available_balance ], stream key: [] }
└── BatchPlanNode { output: [ account_id, available_balance ], stream key: [] }
Fragment 24340 (Actor 112505,112506)
StreamTableScan { table: cash_product_type_ids_mv, columns: [product_type_id, _row_id, $src] } { output: [ cash_product_type_ids_mv.product_type_id, cash_product_type_ids_mv._row_id, cash_product_type_ids_mv.$src ], stream key: [ cash_product_type_ids_mv._row_id, cash_product_type_ids_mv.$src ] }
├── Upstream { output: [ product_type_id, _row_id, $src ], stream key: [] }
└── BatchPlanNode { output: [ product_type_id, _row_id, $src ], stream key: [] }