Fragment 25377 (Actor 124005,124004)
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 25378 (Actor 124006,124007)
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 25379 (Actor 124010,124011)
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 25380 (Actor 124013,124012)
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 25381 (Actor 124016,124017)
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 25382 (Actor 124022,124023)
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 25383 (Actor 124018,124019)
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 25384 (Actor 124024)
StreamNow { output: [ now ], stream key: [] }
Fragment 25385 (Actor 124014,124015)
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 25386 (Actor 124008,124009)
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 25387 (Actor 123929,123930)
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_next.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_next.id AND (assets_dm_next.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_next.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_next.id, assets_dm_next.type ], stream key: [ assets_dm_next.id ] }
Fragment 25388 (Actor 124026,124025)
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 25389 (Actor 123928,123927)
StreamTableScan { table: assets_dm_next, columns: [id, type] } { output: [ assets_dm_next.id, assets_dm_next.type ], stream key: [ assets_dm_next.id ] }
├── Upstream { output: [ id, type ], stream key: [] }
└── BatchPlanNode { output: [ id, type ], stream key: [] }
Fragment 25390 (Actor 124020,124021)
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 25391 (Actor 124027,124028)
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: [] }