Fragment 23569 (Actor 103984,103985)
StreamMaterialize { columns: [account_id, available_balance, fact_date], stream_key: [account_id], pk_columns: [account_id], pk_conflict: NoCheck } { output: [ settled_cash_balances_mv_next.account_id, sum($expr3), max($expr4) ], stream key: [ settled_cash_balances_mv_next.account_id ] }
└── StreamProject { exprs: [settled_cash_balances_mv_next.account_id, sum($expr3), max($expr4)] } { output: [ settled_cash_balances_mv_next.account_id, sum($expr3), max($expr4) ], stream key: [ settled_cash_balances_mv_next.account_id ] }
└── StreamHashAgg { group_key: [settled_cash_balances_mv_next.account_id], aggs: [sum($expr3), max($expr4), count] } { output: [ settled_cash_balances_mv_next.account_id, sum($expr3), max($expr4), count ], stream key: [ settled_cash_balances_mv_next.account_id ] }
└── StreamProject { exprs: [settled_cash_balances_mv_next.account_id, Case((Not(IsNull(holding_values_intraday_ft.account_id)) AND (IsNull(settled_cash_balances_mv_next.account_id) OR (max($expr2) >= settled_cash_balances_mv_next.dim_settlement_date))), sum(holding_values_intraday_ft.market_value), settled_cash_balances_mv_next.settled_cash_balance) as $expr3, Case((Not(IsNull(holding_values_intraday_ft.account_id)) AND (IsNull(settled_cash_balances_mv_next.account_id) OR (max($expr2) >= settled_cash_balances_mv_next.dim_settlement_date))), max($expr2), settled_cash_balances_mv_next.dim_settlement_date) as $expr4, settled_cash_balances_mv_next.currency_code] }
├── output: [ settled_cash_balances_mv_next.account_id, $expr3, $expr4, settled_cash_balances_mv_next.currency_code ]
├── stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ]
└── MergeExecutor { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.settled_cash_balance, settled_cash_balances_mv_next.dim_settlement_date, holding_values_intraday_ft.account_id, sum(holding_values_intraday_ft.market_value), max($expr2), accounts_dm.account_id, settled_cash_balances_mv_next.currency_code ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
Fragment 23570 (Actor 103982,103983)
StreamSyncLogStore { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.settled_cash_balance, settled_cash_balances_mv_next.dim_settlement_date, holding_values_intraday_ft.account_id, sum(holding_values_intraday_ft.market_value), max($expr2), accounts_dm.account_id, settled_cash_balances_mv_next.currency_code ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
└── StreamHashJoin { type: Inner, predicate: settled_cash_balances_mv_next.account_id = accounts_dm.account_id } { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.settled_cash_balance, settled_cash_balances_mv_next.dim_settlement_date, holding_values_intraday_ft.account_id, sum(holding_values_intraday_ft.market_value), max($expr2), accounts_dm.account_id, settled_cash_balances_mv_next.currency_code ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
├── MergeExecutor { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.settled_cash_balance, settled_cash_balances_mv_next.dim_settlement_date, holding_values_intraday_ft.account_id, sum(holding_values_intraday_ft.market_value), max($expr2), settled_cash_balances_mv_next.currency_code, holding_values_intraday_ft.currency_code ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
└── MergeExecutor { output: [ accounts_dm.account_id ], stream key: [ accounts_dm.account_id ] }
Fragment 23571 (Actor 103992,103993)
StreamSyncLogStore { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.settled_cash_balance, settled_cash_balances_mv_next.dim_settlement_date, holding_values_intraday_ft.account_id, sum(holding_values_intraday_ft.market_value), max($expr2), settled_cash_balances_mv_next.currency_code, holding_values_intraday_ft.currency_code ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
└── StreamHashJoin { type: LeftOuter, predicate: settled_cash_balances_mv_next.account_id = holding_values_intraday_ft.account_id AND settled_cash_balances_mv_next.currency_code = holding_values_intraday_ft.currency_code }
├── output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.settled_cash_balance, settled_cash_balances_mv_next.dim_settlement_date, holding_values_intraday_ft.account_id, sum(holding_values_intraday_ft.market_value), max($expr2), settled_cash_balances_mv_next.currency_code, holding_values_intraday_ft.currency_code ]
├── stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ]
├── MergeExecutor { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.settled_cash_balance, settled_cash_balances_mv_next.dim_settlement_date, settled_cash_balances_mv_next.currency_code ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
└── MergeExecutor { output: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code, sum(holding_values_intraday_ft.market_value), max($expr2) ], stream key: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code ] }
Fragment 23572 (Actor 103988,103989)
StreamSyncLogStore { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.settled_cash_balance, settled_cash_balances_mv_next.dim_settlement_date, settled_cash_balances_mv_next.currency_code ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
└── StreamHashJoin { type: LeftOuter, predicate: settled_cash_balances_mv_next.account_id = settled_cash_balances_mv_next.account_id AND settled_cash_balances_mv_next.currency_code = settled_cash_balances_mv_next.currency_code } { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.settled_cash_balance, settled_cash_balances_mv_next.dim_settlement_date, settled_cash_balances_mv_next.currency_code ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
├── StreamProject { exprs: [settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code] } { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
│ └── StreamHashAgg { group_key: [settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code], aggs: [count] } { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, count ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
│ └── MergeExecutor { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, $src ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, $src ] }
└── MergeExecutor { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.settled_cash_balance, settled_cash_balances_mv_next.dim_settlement_date ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
Fragment 23573 (Actor 103995,103994)
StreamUnion { all: true } { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, $src ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, $src ] }
├── MergeExecutor { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, 0:Int32 ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
└── MergeExecutor { output: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code, 1:Int32 ], stream key: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code ] }
Fragment 23574 (Actor 104001,104000)
StreamProject { exprs: [settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, 0:Int32] } { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, 0:Int32 ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
└── MergeExecutor { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.settled_cash_balance, settled_cash_balances_mv_next.dim_settlement_date ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
Fragment 23575 (Actor 103999,103998)
StreamProject { exprs: [settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.settled_cash_balance, settled_cash_balances_mv_next.dim_settlement_date] } { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.settled_cash_balance, settled_cash_balances_mv_next.dim_settlement_date ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
└── StreamGroupTopN { order: [settled_cash_balances_mv_next.dim_settlement_date DESC], limit: 1, offset: 0, group_key: [settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code] } { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.dim_settlement_date, settled_cash_balances_mv_next.settled_cash_balance ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
└── StreamProject { exprs: [settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.dim_settlement_date, settled_cash_balances_mv_next.settled_cash_balance] } { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.dim_settlement_date, settled_cash_balances_mv_next.settled_cash_balance ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.dim_settlement_date ] }
└── StreamDynamicFilter { predicate: ($expr1 <= now), output: [settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.dim_settlement_date, settled_cash_balances_mv_next.settled_cash_balance, $expr1], cleaned_by_watermark: true } { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.dim_settlement_date, settled_cash_balances_mv_next.settled_cash_balance, $expr1 ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.dim_settlement_date ] }
├── StreamProject { exprs: [settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.dim_settlement_date, settled_cash_balances_mv_next.settled_cash_balance, AtTimeZone(settled_cash_balances_mv_next.dim_settlement_date::Timestamp, 'UTC':Varchar) as $expr1] }
│ ├── output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.dim_settlement_date, settled_cash_balances_mv_next.settled_cash_balance, $expr1 ]
│ ├── stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.dim_settlement_date ]
│ └── StreamTableScan { table: settled_cash_balances_mv_next, columns: [account_id, currency_code, dim_settlement_date, settled_cash_balance] } { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.dim_settlement_date, settled_cash_balances_mv_next.settled_cash_balance ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.dim_settlement_date ] }
│ ├── Upstream { output: [ account_id, currency_code, dim_settlement_date, settled_cash_balance ], stream key: [] }
│ └── BatchPlanNode { output: [ account_id, currency_code, dim_settlement_date, settled_cash_balance ], stream key: [] }
└── MergeExecutor { output: [ now ], stream key: [] }
Fragment 23576 (Actor 104002)
StreamNow { output: [ now ], stream key: [] }
Fragment 23577 (Actor 103987,103986)
StreamProject { exprs: [holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code, 1:Int32] } { output: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code, 1:Int32 ], stream key: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code ] }
└── MergeExecutor { output: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code, sum(holding_values_intraday_ft.market_value), max($expr2) ], stream key: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code ] }
Fragment 23578 (Actor 103990,103991)
StreamProject { exprs: [holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code, sum(holding_values_intraday_ft.market_value), max($expr2)] } { output: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code, sum(holding_values_intraday_ft.market_value), max($expr2) ], stream key: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code ] }
└── StreamHashAgg { group_key: [holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code], aggs: [sum(holding_values_intraday_ft.market_value), max($expr2), count] } { output: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code, sum(holding_values_intraday_ft.market_value), max($expr2), count ], stream key: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code ] }
└── MergeExecutor { output: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code, holding_values_intraday_ft.market_value, $expr2, holding_values_intraday_ft.asset_id ], stream key: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.asset_id ] }
Fragment 23579 (Actor 103932,103931)
StreamProject { exprs: [holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code, holding_values_intraday_ft.market_value, AtTimeZone(holding_values_intraday_ft.holding_timestamp, 'UTC':Varchar)::Date as $expr2, holding_values_intraday_ft.asset_id] } { output: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.currency_code, holding_values_intraday_ft.market_value, $expr2, holding_values_intraday_ft.asset_id ], stream key: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.asset_id ] }
└── StreamGroupTopN { order: [holding_values_intraday_ft.holding_timestamp DESC], limit: 1, offset: 0, group_key: [holding_values_intraday_ft.account_id, holding_values_intraday_ft.asset_id] } { output: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.asset_id, holding_values_intraday_ft.currency_code, holding_values_intraday_ft.holding_timestamp, holding_values_intraday_ft.market_value, holding_values_intraday_ft.id, assets_dm.id ], stream key: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.asset_id ] }
└── StreamTemporalJoin { type: Inner, append_only: false, predicate: holding_values_intraday_ft.asset_id = assets_dm.id AND (assets_dm.type = 'CASH':Varchar), nested_loop: false, output_watermarks: [[holding_values_intraday_ft.holding_timestamp]] } { output: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.asset_id, holding_values_intraday_ft.currency_code, holding_values_intraday_ft.holding_timestamp, holding_values_intraday_ft.market_value, holding_values_intraday_ft.id, assets_dm.id ], stream key: [ holding_values_intraday_ft.id, holding_values_intraday_ft.holding_timestamp, holding_values_intraday_ft.asset_id ] }
├── MergeExecutor { output: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.asset_id, holding_values_intraday_ft.currency_code, holding_values_intraday_ft.holding_timestamp, holding_values_intraday_ft.market_value, holding_values_intraday_ft.id ], stream key: [ holding_values_intraday_ft.id, holding_values_intraday_ft.holding_timestamp ] }
└── MergeExecutor { output: [ assets_dm.id, assets_dm.type ], stream key: [ assets_dm.id ] }
Fragment 23580 (Actor 104003,104004)
StreamProject { exprs: [holding_values_intraday_ft.account_id, holding_values_intraday_ft.asset_id, holding_values_intraday_ft.currency_code, holding_values_intraday_ft.holding_timestamp, holding_values_intraday_ft.market_value, holding_values_intraday_ft.id], output_watermarks: [[holding_values_intraday_ft.holding_timestamp]] } { output: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.asset_id, holding_values_intraday_ft.currency_code, holding_values_intraday_ft.holding_timestamp, holding_values_intraday_ft.market_value, holding_values_intraday_ft.id ], stream key: [ holding_values_intraday_ft.id, holding_values_intraday_ft.holding_timestamp ] }
└── StreamFilter { predicate: IsNull(holding_values_intraday_ft.disabled_at) } { output: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.asset_id, holding_values_intraday_ft.currency_code, holding_values_intraday_ft.holding_timestamp, holding_values_intraday_ft.market_value, holding_values_intraday_ft.id, holding_values_intraday_ft.disabled_at ], stream key: [ holding_values_intraday_ft.id, holding_values_intraday_ft.holding_timestamp ] }
└── StreamTableScan { table: holding_values_intraday_ft, columns: [account_id, asset_id, currency_code, holding_timestamp, market_value, id, disabled_at] } { output: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.asset_id, holding_values_intraday_ft.currency_code, holding_values_intraday_ft.holding_timestamp, holding_values_intraday_ft.market_value, holding_values_intraday_ft.id, holding_values_intraday_ft.disabled_at ], stream key: [ holding_values_intraday_ft.id, holding_values_intraday_ft.holding_timestamp ] }
├── Upstream { output: [ account_id, asset_id, currency_code, holding_timestamp, market_value, id, disabled_at ], stream key: [] }
└── BatchPlanNode { output: [ account_id, asset_id, currency_code, holding_timestamp, market_value, id, disabled_at ], stream key: [] }
Fragment 23581 (Actor 103929,103930)
StreamTableScan { table: assets_dm, columns: [id, type] } { output: [ assets_dm.id, assets_dm.type ], stream key: [ assets_dm.id ] }
├── Upstream { output: [ id, type ], stream key: [] }
└── BatchPlanNode { output: [ id, type ], stream key: [] }
Fragment 23582 (Actor 103996,103997)
StreamNoOp { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.settled_cash_balance, settled_cash_balances_mv_next.dim_settlement_date ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
└── MergeExecutor { output: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code, settled_cash_balances_mv_next.settled_cash_balance, settled_cash_balances_mv_next.dim_settlement_date ], stream key: [ settled_cash_balances_mv_next.account_id, settled_cash_balances_mv_next.currency_code ] }
Fragment 23583 (Actor 104005,104006)
StreamProject { exprs: [accounts_dm.account_id] } { output: [ accounts_dm.account_id ], stream key: [ accounts_dm.account_id ] }
└── StreamFilter { predicate: IsNull(accounts_dm.disabled_at) } { output: [ accounts_dm.account_id, accounts_dm.disabled_at ], stream key: [ accounts_dm.account_id ] }
└── StreamTableScan { table: accounts_dm, columns: [account_id, disabled_at] } { output: [ accounts_dm.account_id, accounts_dm.disabled_at ], stream key: [ accounts_dm.account_id ] }
├── Upstream { output: [ account_id, disabled_at ], stream key: [] }
└── BatchPlanNode { output: [ account_id, disabled_at ], stream key: [] }