Fragment 23653 (Actor 104896,104895)
StreamMaterialize { columns: [client_id, cash_tiles], stream_key: [client_id], pk_columns: [client_id], pk_conflict: NoCheck }
├── output: [ clients_dm.id, jsonb_agg($expr3) ]
├── stream key: [ clients_dm.id ]
└── StreamProject { exprs: [clients_dm.id, jsonb_agg($expr3)] }
├── output: [ clients_dm.id, jsonb_agg($expr3) ]
├── stream key: [ clients_dm.id ]
└── StreamHashAgg { group_key: [clients_dm.id], aggs: [jsonb_agg($expr3), count] }
├── output: [ clients_dm.id, jsonb_agg($expr3), count ]
├── stream key: [ clients_dm.id ]
└── MergeExecutor
├── output:
│ ┌── clients_dm.id
│ ├── $expr3
│ ├── accounts_to_clients_dm.account_id
│ ├── accounts_to_clients_dm.effective_start_date
│ ├── accounts_plain_mv.base_currency_code
│ ├── accounts_plain_mv.account_id
│ └── accounts_plain_mv.product_type_id
└── stream key:
┌── clients_dm.id
├── accounts_to_clients_dm.account_id
├── accounts_to_clients_dm.effective_start_date
├── accounts_plain_mv.base_currency_code
├── accounts_plain_mv.account_id
└── accounts_plain_mv.product_type_id
Fragment 23654 (Actor 104897,104898)
StreamProject { exprs: [clients_dm.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_next.available_balance::Varchar, 'currency':Varchar, $expr2), 'baseCurrency':Varchar, $expr2) as $expr3, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date, accounts_plain_mv.base_currency_code, accounts_plain_mv.account_id, accounts_plain_mv.product_type_id] }
├── output: [ clients_dm.id, $expr3, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date, accounts_plain_mv.base_currency_code, accounts_plain_mv.account_id, accounts_plain_mv.product_type_id ]
├── stream key: [ clients_dm.id, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date, accounts_plain_mv.base_currency_code, accounts_plain_mv.account_id, accounts_plain_mv.product_type_id ]
└── StreamProject { exprs: [clients_dm.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_next.available_balance, JsonbBuildObject('code':Varchar, accounts_plain_mv.base_currency_code, 'symbol':Varchar, currencies_dm.symbol) as $expr2, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date, accounts_plain_mv.product_type_id] }
├── output: [ clients_dm.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_next.available_balance, $expr2, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date, accounts_plain_mv.product_type_id ]
├── stream key: [ clients_dm.id, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date, accounts_plain_mv.base_currency_code, accounts_plain_mv.account_id, accounts_plain_mv.product_type_id ]
└── MergeExecutor
├── output: [ clients_dm.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_next.available_balance, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date, accounts_plain_mv.product_type_id ]
└── stream key: [ clients_dm.id, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date, accounts_plain_mv.base_currency_code, accounts_plain_mv.account_id, accounts_plain_mv.product_type_id ]
Fragment 23655 (Actor 104899,104900)
StreamSyncLogStore
├── output: [ clients_dm.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_next.available_balance, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date, accounts_plain_mv.product_type_id ]
├── stream key: [ clients_dm.id, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date, accounts_plain_mv.base_currency_code, 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: [ clients_dm.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_next.available_balance, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date, accounts_plain_mv.product_type_id ]
├── stream key: [ clients_dm.id, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date, accounts_plain_mv.base_currency_code, accounts_plain_mv.account_id, accounts_plain_mv.product_type_id ]
├── MergeExecutor
│ ├── output: [ clients_dm.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_next.available_balance, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date, investment_account_balance_snapshot_journal_mv_next.account_id ]
│ └── stream key: [ clients_dm.id, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date, accounts_plain_mv.base_currency_code, 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 23656 (Actor 104902,104901)
StreamSyncLogStore
├── output: [ clients_dm.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_next.available_balance, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date, investment_account_balance_snapshot_journal_mv_next.account_id ]
├── stream key: [ clients_dm.id, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date, accounts_plain_mv.base_currency_code, accounts_plain_mv.account_id ]
└── StreamHashJoin { type: LeftOuter, predicate: accounts_plain_mv.account_id = investment_account_balance_snapshot_journal_mv_next.account_id }
├── output: [ clients_dm.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_next.available_balance, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date, investment_account_balance_snapshot_journal_mv_next.account_id ]
├── stream key: [ clients_dm.id, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date, accounts_plain_mv.base_currency_code, accounts_plain_mv.account_id ]
├── MergeExecutor { output: [ clients_dm.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, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date, currencies_dm.code ], stream key: [ clients_dm.id, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date, accounts_plain_mv.base_currency_code ] }
└── MergeExecutor { output: [ investment_account_balance_snapshot_journal_mv_next.account_id, investment_account_balance_snapshot_journal_mv_next.available_balance ], stream key: [ investment_account_balance_snapshot_journal_mv_next.account_id ] }
Fragment 23657 (Actor 104827,104826)
StreamTemporalJoin { type: LeftOuter, append_only: false, predicate: accounts_plain_mv.base_currency_code = currencies_dm.code, nested_loop: false }
├── output: [ clients_dm.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, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date, currencies_dm.code ]
├── stream key: [ clients_dm.id, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date, accounts_plain_mv.base_currency_code ]
├── MergeExecutor { output: [ clients_dm.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, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date ], stream key: [ clients_dm.id, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date ] }
└── MergeExecutor { output: [ currencies_dm.code, currencies_dm.symbol ], stream key: [ currencies_dm.code ] }
Fragment 23658 (Actor 104904,104903)
StreamSyncLogStore { output: [ clients_dm.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, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date ], stream key: [ clients_dm.id, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date ] }
└── StreamHashJoin { type: Inner, predicate: accounts_to_clients_dm.account_id = accounts_plain_mv.account_id } { output: [ clients_dm.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, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date ], stream key: [ clients_dm.id, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date ] }
├── MergeExecutor { output: [ clients_dm.id, accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ], stream key: [ clients_dm.id, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date ] }
└── 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 23659 (Actor 104905,104906)
StreamSyncLogStore { output: [ clients_dm.id, accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ], stream key: [ clients_dm.id, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date ] }
└── StreamHashJoin { type: Inner, predicate: clients_dm.id = accounts_to_clients_dm.client_id } { output: [ clients_dm.id, accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ], stream key: [ clients_dm.id, accounts_to_clients_dm.account_id, accounts_to_clients_dm.effective_start_date ] }
├── MergeExecutor { output: [ clients_dm.id ], stream key: [ clients_dm.id ] }
└── MergeExecutor { output: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ] }
Fragment 23660 (Actor 104907,104908)
StreamProject { exprs: [clients_dm.id] } { output: [ clients_dm.id ], stream key: [ clients_dm.id ] }
└── StreamFilter { predicate: IsNull(clients_dm.disabled_at) } { output: [ clients_dm.id, clients_dm.disabled_at ], stream key: [ clients_dm.id ] }
└── StreamTableScan { table: clients_dm, columns: [id, disabled_at] } { output: [ clients_dm.id, clients_dm.disabled_at ], stream key: [ clients_dm.id ] }
├── Upstream { output: [ id, disabled_at ], stream key: [] }
└── BatchPlanNode { output: [ id, disabled_at ], stream key: [] }
Fragment 23661 (Actor 104910,104909)
StreamProject { exprs: [accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date] } { output: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ] }
└── StreamDynamicFilter { predicate: ($expr1 <= now), output: [accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, $expr1, accounts_to_clients_dm.effective_start_date], cleaned_by_watermark: true } { output: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, $expr1, accounts_to_clients_dm.effective_start_date ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ] }
├── StreamProject { exprs: [accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, AtTimeZone(accounts_to_clients_dm.effective_start_date::Timestamp, 'UTC':Varchar) as $expr1, accounts_to_clients_dm.effective_start_date] } { output: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, $expr1, accounts_to_clients_dm.effective_start_date ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ] }
│ └── StreamFilter { predicate: IsNull(accounts_to_clients_dm.disabled_at) AND IsNull(accounts_to_clients_dm.effective_end_date) } { output: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, accounts_to_clients_dm.disabled_at, accounts_to_clients_dm.effective_end_date ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ] }
│ └── StreamTableScan { table: accounts_to_clients_dm, columns: [account_id, client_id, effective_start_date, disabled_at, effective_end_date] } { output: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date, accounts_to_clients_dm.disabled_at, accounts_to_clients_dm.effective_end_date ], stream key: [ accounts_to_clients_dm.account_id, accounts_to_clients_dm.client_id, accounts_to_clients_dm.effective_start_date ] }
│ ├── Upstream { output: [ account_id, client_id, effective_start_date, disabled_at, effective_end_date ], stream key: [] }
│ └── BatchPlanNode { output: [ account_id, client_id, effective_start_date, disabled_at, effective_end_date ], stream key: [] }
└── MergeExecutor { output: [ now ], stream key: [] }
Fragment 23662 (Actor 104911)
StreamNow { output: [ now ], stream key: [] }
Fragment 23663 (Actor 104913,104912)
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 23664 (Actor 104825,104824)
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 23665 (Actor 104914,104915)
StreamTableScan { table: investment_account_balance_snapshot_journal_mv_next, columns: [account_id, available_balance] } { output: [ investment_account_balance_snapshot_journal_mv_next.account_id, investment_account_balance_snapshot_journal_mv_next.available_balance ], stream key: [ investment_account_balance_snapshot_journal_mv_next.account_id ] }
├── Upstream { output: [ account_id, available_balance ], stream key: [] }
└── BatchPlanNode { output: [ account_id, available_balance ], stream key: [] }
Fragment 23666 (Actor 104916,104917)
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: [] }