Job is idle — throughput ~0; structure shown.
Fragment 58793 (Actor 739225,739224)
StreamMaterialize { columns: [account_id, asset_id, dim_value_date, type, currency_code, market_value, average_cost_per_unit, purchased_quantity, average_cost_per_unit_system_currency, total_cost_system_currency, deposit_profit_accrued, holding_values_raw_ft.account_id(hidden), holding_values_raw_ft.dim_value_date(hidden), min(holding_values_raw_ft.asset_id)(hidden), 'ASSET':Varchar(hidden)], stream_key: [account_id, dim_value_date, asset_id, type], pk_columns: [account_id, dim_value_date, asset_id, type], pk_conflict: NoCheck, watermark_columns: [dim_value_date, holding_values_raw_ft.dim_value_date(hidden)] }
├── output: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.type, holding_values_raw_ft.currency_code, holding_values_raw_ft.market_value, holding_values_raw_ft.average_cost_per_unit, holding_values_raw_ft.purchased_quantity, holding_values_raw_ft.average_cost_per_unit_system_currency, holding_values_raw_ft.total_cost_system_currency, max(fixed_deposit_accounts_ft.profit_accrued), holding_values_raw_ft.account_id, holding_values_raw_ft.dim_value_date, min(holding_values_raw_ft.asset_id), 'ASSET':Varchar ]
├── stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.asset_id, holding_values_raw_ft.type ]
└── 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.type, holding_values_raw_ft.currency_code, holding_values_raw_ft.market_value, holding_values_raw_ft.average_cost_per_unit, holding_values_raw_ft.purchased_quantity, holding_values_raw_ft.average_cost_per_unit_system_currency, holding_values_raw_ft.total_cost_system_currency, max(fixed_deposit_accounts_ft.profit_accrued), holding_values_raw_ft.account_id, holding_values_raw_ft.dim_value_date, min(holding_values_raw_ft.asset_id), 'ASSET':Varchar ]
└── stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.asset_id, holding_values_raw_ft.type ]
Fragment 58794 (Actor 739226,739227)
StreamSyncLogStore
├── output: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.type, holding_values_raw_ft.currency_code, holding_values_raw_ft.market_value, holding_values_raw_ft.average_cost_per_unit, holding_values_raw_ft.purchased_quantity, holding_values_raw_ft.average_cost_per_unit_system_currency, holding_values_raw_ft.total_cost_system_currency, max(fixed_deposit_accounts_ft.profit_accrued), holding_values_raw_ft.account_id, holding_values_raw_ft.dim_value_date, min(holding_values_raw_ft.asset_id), 'ASSET':Varchar ]
├── stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.asset_id, holding_values_raw_ft.type ]
└── StreamHashJoin [window] { type: LeftOuter, predicate: holding_values_raw_ft.dim_value_date = holding_values_raw_ft.dim_value_date AND holding_values_raw_ft.account_id = holding_values_raw_ft.account_id AND holding_values_raw_ft.asset_id = min(holding_values_raw_ft.asset_id) AND holding_values_raw_ft.type = 'ASSET':Varchar, conditions_to_clean_state_in_join_key: [(holding_values_raw_ft.dim_value_date = holding_values_raw_ft.dim_value_date)], output_watermarks: [[holding_values_raw_ft.dim_value_date], [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.type, holding_values_raw_ft.currency_code, holding_values_raw_ft.market_value, holding_values_raw_ft.average_cost_per_unit, holding_values_raw_ft.purchased_quantity, holding_values_raw_ft.average_cost_per_unit_system_currency, holding_values_raw_ft.total_cost_system_currency, max(fixed_deposit_accounts_ft.profit_accrued), holding_values_raw_ft.account_id, holding_values_raw_ft.dim_value_date, min(holding_values_raw_ft.asset_id), 'ASSET':Varchar ]
├── stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.asset_id, holding_values_raw_ft.type ]
├── 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.currency_code, holding_values_raw_ft.market_value, holding_values_raw_ft.average_cost_per_unit, holding_values_raw_ft.purchased_quantity, holding_values_raw_ft.type, holding_values_raw_ft.average_cost_per_unit_system_currency, holding_values_raw_ft.total_cost_system_currency ], stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.asset_id, holding_values_raw_ft.type ] }
└── MergeExecutor { output: [ holding_values_raw_ft.account_id, holding_values_raw_ft.dim_value_date, min(holding_values_raw_ft.asset_id), max(fixed_deposit_accounts_ft.profit_accrued), 'ASSET':Varchar ], stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.dim_value_date, min(holding_values_raw_ft.asset_id), 'ASSET':Varchar ] }
Fragment 58795 (Actor 739228,739229)
StreamLocalityProvider { locality_columns: [holding_values_raw_ft.account_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.asset_id, holding_values_raw_ft.type] }
├── output: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.currency_code, holding_values_raw_ft.market_value, holding_values_raw_ft.average_cost_per_unit, holding_values_raw_ft.purchased_quantity, holding_values_raw_ft.type, holding_values_raw_ft.average_cost_per_unit_system_currency, holding_values_raw_ft.total_cost_system_currency ]
├── stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.asset_id, holding_values_raw_ft.type ]
└── 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.currency_code, holding_values_raw_ft.market_value, holding_values_raw_ft.average_cost_per_unit, holding_values_raw_ft.purchased_quantity, holding_values_raw_ft.type, holding_values_raw_ft.average_cost_per_unit_system_currency, holding_values_raw_ft.total_cost_system_currency ], 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 ] }
Fragment 58796 (Actor 739231,739230)
StreamTableScan { table: holding_values_raw_ft, columns: [account_id, asset_id, dim_value_date, currency_code, market_value, average_cost_per_unit, purchased_quantity, type, average_cost_per_unit_system_currency, total_cost_system_currency] }
├── output: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.currency_code, holding_values_raw_ft.market_value, holding_values_raw_ft.average_cost_per_unit, holding_values_raw_ft.purchased_quantity, holding_values_raw_ft.type, holding_values_raw_ft.average_cost_per_unit_system_currency, holding_values_raw_ft.total_cost_system_currency ]
├── 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, currency_code, market_value, average_cost_per_unit, purchased_quantity, type, average_cost_per_unit_system_currency, total_cost_system_currency ], stream key: [] }
└── BatchPlanNode { output: [ account_id, asset_id, dim_value_date, currency_code, market_value, average_cost_per_unit, purchased_quantity, type, average_cost_per_unit_system_currency, total_cost_system_currency ], stream key: [] }
Fragment 58797 (Actor 739232,739233)
StreamLocalityProvider { locality_columns: [holding_values_raw_ft.account_id, holding_values_raw_ft.dim_value_date, min(holding_values_raw_ft.asset_id), 'ASSET':Varchar] } { output: [ holding_values_raw_ft.account_id, holding_values_raw_ft.dim_value_date, min(holding_values_raw_ft.asset_id), max(fixed_deposit_accounts_ft.profit_accrued), 'ASSET':Varchar ], stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.dim_value_date, min(holding_values_raw_ft.asset_id), 'ASSET':Varchar ] }
└── MergeExecutor { output: [ holding_values_raw_ft.account_id, holding_values_raw_ft.dim_value_date, min(holding_values_raw_ft.asset_id), max(fixed_deposit_accounts_ft.profit_accrued), 'ASSET':Varchar ], stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.dim_value_date ] }
Fragment 58798 (Actor 739254,739255)
StreamProject { exprs: [holding_values_raw_ft.account_id, holding_values_raw_ft.dim_value_date, min(holding_values_raw_ft.asset_id), max(fixed_deposit_accounts_ft.profit_accrued), 'ASSET':Varchar], output_watermarks: [[holding_values_raw_ft.dim_value_date]] } { output: [ holding_values_raw_ft.account_id, holding_values_raw_ft.dim_value_date, min(holding_values_raw_ft.asset_id), max(fixed_deposit_accounts_ft.profit_accrued), 'ASSET':Varchar ], stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.dim_value_date ] }
└── StreamHashAgg { group_key: [holding_values_raw_ft.account_id, holding_values_raw_ft.dim_value_date], aggs: [min(holding_values_raw_ft.asset_id), max(fixed_deposit_accounts_ft.profit_accrued), count], output_watermarks: [[holding_values_raw_ft.dim_value_date]] } { output: [ holding_values_raw_ft.account_id, holding_values_raw_ft.dim_value_date, min(holding_values_raw_ft.asset_id), max(fixed_deposit_accounts_ft.profit_accrued), count ], stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.dim_value_date ] }
└── StreamLocalityProvider { locality_columns: [holding_values_raw_ft.account_id, 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, fixed_deposit_accounts_ft.profit_accrued, holding_values_raw_ft.type, fixed_deposit_accounts_ft.account_id, fixed_deposit_accounts_ft.fact_date ], stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.asset_id, holding_values_raw_ft.type ] }
└── MergeExecutor { output: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, fixed_deposit_accounts_ft.profit_accrued, holding_values_raw_ft.type, fixed_deposit_accounts_ft.account_id, fixed_deposit_accounts_ft.fact_date ], stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.asset_id, holding_values_raw_ft.type ] }
Fragment 58799 (Actor 739253,739252)
StreamSyncLogStore { output: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, fixed_deposit_accounts_ft.profit_accrued, holding_values_raw_ft.type, fixed_deposit_accounts_ft.account_id, fixed_deposit_accounts_ft.fact_date ], stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.asset_id, holding_values_raw_ft.type ] }
└── StreamHashJoin [window] { type: Inner, predicate: holding_values_raw_ft.dim_value_date = fixed_deposit_accounts_ft.fact_date AND holding_values_raw_ft.account_id = fixed_deposit_accounts_ft.account_id, conditions_to_clean_state_in_join_key: [(holding_values_raw_ft.dim_value_date = fixed_deposit_accounts_ft.fact_date)], output_watermarks: [[holding_values_raw_ft.dim_value_date], [fixed_deposit_accounts_ft.fact_date]] }
├── output: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, fixed_deposit_accounts_ft.profit_accrued, holding_values_raw_ft.type, fixed_deposit_accounts_ft.account_id, fixed_deposit_accounts_ft.fact_date ]
├── stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.asset_id, holding_values_raw_ft.type ]
├── 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.type ], stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.asset_id, holding_values_raw_ft.type ] }
└── MergeExecutor { output: [ fixed_deposit_accounts_ft.account_id, fixed_deposit_accounts_ft.fact_date, fixed_deposit_accounts_ft.profit_accrued ], stream key: [ fixed_deposit_accounts_ft.account_id, fixed_deposit_accounts_ft.fact_date ] }
Fragment 58800 (Actor 739256,739257)
StreamLocalityProvider { locality_columns: [holding_values_raw_ft.account_id, 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.type ], stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.asset_id, holding_values_raw_ft.type ] }
└── 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.type ], 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 ] }
Fragment 58801 (Actor 739258,739259)
StreamFilter { predicate: (holding_values_raw_ft.type = 'ASSET':Varchar) } { output: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.type ], 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, type] } { output: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.type ], 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, type ], stream key: [] }
└── BatchPlanNode { output: [ account_id, asset_id, dim_value_date, type ], stream key: [] }
Fragment 58802 (Actor 739260,739261)
StreamProject { exprs: [fixed_deposit_accounts_ft.account_id, fixed_deposit_accounts_ft.fact_date, fixed_deposit_accounts_ft.profit_accrued], output_watermarks: [[fixed_deposit_accounts_ft.fact_date]] } { output: [ fixed_deposit_accounts_ft.account_id, fixed_deposit_accounts_ft.fact_date, fixed_deposit_accounts_ft.profit_accrued ], stream key: [ fixed_deposit_accounts_ft.account_id, fixed_deposit_accounts_ft.fact_date ] }
└── StreamFilter { predicate: Not(IsNull(fixed_deposit_accounts_ft.profit_accrued)) AND IsNull(fixed_deposit_accounts_ft.disabled_at) } { output: [ fixed_deposit_accounts_ft.account_id, fixed_deposit_accounts_ft.fact_date, fixed_deposit_accounts_ft.profit_accrued, fixed_deposit_accounts_ft.disabled_at ], stream key: [ fixed_deposit_accounts_ft.account_id, fixed_deposit_accounts_ft.fact_date ] }
└── StreamTableScan { table: fixed_deposit_accounts_ft, columns: [account_id, fact_date, profit_accrued, disabled_at] } { output: [ fixed_deposit_accounts_ft.account_id, fixed_deposit_accounts_ft.fact_date, fixed_deposit_accounts_ft.profit_accrued, fixed_deposit_accounts_ft.disabled_at ], stream key: [ fixed_deposit_accounts_ft.account_id, fixed_deposit_accounts_ft.fact_date ] }
├── Upstream { output: [ account_id, fact_date, profit_accrued, disabled_at ], stream key: [] }
└── BatchPlanNode { output: [ account_id, fact_date, profit_accrued, disabled_at ], stream key: [] }