Job is idle — throughput ~0; structure shown.
Fragment 8857 (Actor 96951,96952)
StreamMaterialize { columns: [account_id, available_balance, fact_date], stream_key: [account_id], pk_columns: [account_id], pk_conflict: NoCheck } { output: [ holding_values_raw_ft.account_id, sum($expr3), max($expr4) ], stream key: [ holding_values_raw_ft.account_id ] }
└── StreamProject { exprs: [holding_values_raw_ft.account_id, sum($expr3), max($expr4)] } { output: [ holding_values_raw_ft.account_id, sum($expr3), max($expr4) ], stream key: [ holding_values_raw_ft.account_id ] }
└── StreamHashAgg { group_key: [holding_values_raw_ft.account_id], aggs: [sum($expr3), max($expr4), count] } { output: [ holding_values_raw_ft.account_id, sum($expr3), max($expr4), count ], stream key: [ holding_values_raw_ft.account_id ] }
└── StreamProject { exprs: [holding_values_raw_ft.account_id, Case((Not(IsNull(holding_values_intraday_ft.account_id)) AND ($expr2 >= holding_values_raw_ft.dim_value_date)), holding_values_intraday_ft.market_value, holding_values_raw_ft.market_value) as $expr3, Case((Not(IsNull(holding_values_intraday_ft.account_id)) AND ($expr2 >= holding_values_raw_ft.dim_value_date)), $expr2, holding_values_raw_ft.dim_value_date) as $expr4, holding_values_raw_ft.asset_id] }
├── output: [ holding_values_raw_ft.account_id, $expr3, $expr4, holding_values_raw_ft.asset_id ]
├── stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id ]
└── StreamHashJoin { type: Inner, predicate: holding_values_raw_ft.account_id = accounts_dm.account_id } { output: [ holding_values_raw_ft.account_id, holding_values_raw_ft.market_value, holding_values_raw_ft.dim_value_date, holding_values_intraday_ft.account_id, holding_values_intraday_ft.market_value, $expr2, accounts_dm.account_id, holding_values_raw_ft.asset_id ], stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id ] }
├── MergeExecutor { output: [ holding_values_raw_ft.account_id, holding_values_raw_ft.market_value, holding_values_raw_ft.dim_value_date, holding_values_intraday_ft.account_id, holding_values_intraday_ft.market_value, $expr2, holding_values_raw_ft.asset_id, holding_values_intraday_ft.asset_id ], stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id ] }
└── MergeExecutor { output: [ accounts_dm.account_id ], stream key: [ accounts_dm.account_id ] }
Fragment 8858 (Actor 96954,96953)
StreamHashJoin { type: LeftOuter, predicate: holding_values_raw_ft.account_id = holding_values_intraday_ft.account_id AND holding_values_raw_ft.asset_id = holding_values_intraday_ft.asset_id }
├── output: [ holding_values_raw_ft.account_id, holding_values_raw_ft.market_value, holding_values_raw_ft.dim_value_date, holding_values_intraday_ft.account_id, holding_values_intraday_ft.market_value, $expr2, holding_values_raw_ft.asset_id, holding_values_intraday_ft.asset_id ]
├── stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id ]
├── StreamProject { exprs: [holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.market_value, holding_values_raw_ft.type, $expr1] }
│ ├── output: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.market_value, holding_values_raw_ft.type, $expr1 ]
│ ├── stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id ]
│ └── StreamGroupTopN { order: [$expr1 ASC, holding_values_raw_ft.dim_value_date DESC, holding_values_raw_ft.type ASC], limit: 1, offset: 0, group_key: [holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id] }
│ ├── output: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.market_value, holding_values_raw_ft.type, $expr1 ]
│ ├── stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id ]
│ └── MergeExecutor { output: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.market_value, holding_values_raw_ft.type, $expr1 ], stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.type ] }
└── StreamProject { exprs: [holding_values_intraday_ft.account_id, holding_values_intraday_ft.asset_id, holding_values_intraday_ft.market_value, AtTimeZone(holding_values_intraday_ft.holding_timestamp, 'UTC':Varchar)::Date as $expr2] } { output: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.asset_id, holding_values_intraday_ft.market_value, $expr2 ], 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.holding_timestamp, holding_values_intraday_ft.market_value, holding_values_intraday_ft.id ]
├── stream key: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.asset_id ]
└── MergeExecutor { output: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.asset_id, 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 ] }
Fragment 8859 (Actor 96139,96140)
StreamProject { exprs: [holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.market_value, holding_values_raw_ft.type, IsTrue(holding_values_raw_ft.m_is_stub) as $expr1], output_watermarks: [[holding_values_raw_ft.dim_value_date]] }
├── output: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.market_value, holding_values_raw_ft.type, $expr1 ]
├── stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.type ]
└── StreamTableScan { table: holding_values_raw_ft, columns: [account_id, asset_id, dim_value_date, market_value, type, m_is_stub] }
├── output: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.market_value, holding_values_raw_ft.type, holding_values_raw_ft.m_is_stub ]
├── stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.type ]
├── Upstream { output: [ account_id, asset_id, dim_value_date, market_value, type, m_is_stub ], stream key: [] }
└── BatchPlanNode { output: [ account_id, asset_id, dim_value_date, market_value, type, m_is_stub ], stream key: [] }
Fragment 8860 (Actor 98250,98249)
StreamProject { exprs: [holding_values_intraday_ft.account_id, holding_values_intraday_ft.asset_id, 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.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.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, holding_timestamp, market_value, id, disabled_at] }
├── output: [ holding_values_intraday_ft.account_id, holding_values_intraday_ft.asset_id, 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, holding_timestamp, market_value, id, disabled_at ], stream key: [] }
└── BatchPlanNode { output: [ account_id, asset_id, holding_timestamp, market_value, id, disabled_at ], stream key: [] }
Fragment 8861 (Actor 97023,97024)
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: [] }