Job is idle — throughput ~0; structure shown.
Fragment 62562 (Actor 743626,743625)
StreamMaterialize { columns: [benchmark_id, fact_date, currency_code, weight], stream_key: [benchmark_id, fact_date, currency_code], pk_columns: [benchmark_id, fact_date, currency_code], pk_conflict: NoCheck }
├── output: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, assets_dm_next.issue_currency_code, sum(benchmark_constituents_ft.weight) ]
├── stream key: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, assets_dm_next.issue_currency_code ]
└── StreamProject { exprs: [benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, assets_dm_next.issue_currency_code, sum(benchmark_constituents_ft.weight)] }
├── output: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, assets_dm_next.issue_currency_code, sum(benchmark_constituents_ft.weight) ]
├── stream key: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, assets_dm_next.issue_currency_code ]
└── StreamHashAgg { group_key: [benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, assets_dm_next.issue_currency_code], aggs: [sum(benchmark_constituents_ft.weight), count] }
├── output: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, assets_dm_next.issue_currency_code, sum(benchmark_constituents_ft.weight), count ]
├── stream key: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, assets_dm_next.issue_currency_code ]
└── StreamLocalityProvider { locality_columns: [benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, assets_dm_next.issue_currency_code] }
├── output:
│ ┌── benchmark_constituents_ft.benchmark_id
│ ├── benchmark_constituents_ft.date
│ ├── benchmark_constituents_ft.weight
│ ├── assets_dm_next.issue_currency_code
│ ├── benchmark_constituents_ft.asset_id
│ └── benchmarks_dm.id
├── stream key: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, assets_dm_next.issue_currency_code, benchmark_constituents_ft.asset_id ]
└── MergeExecutor
├── output:
│ ┌── benchmark_constituents_ft.benchmark_id
│ ├── benchmark_constituents_ft.date
│ ├── benchmark_constituents_ft.weight
│ ├── assets_dm_next.issue_currency_code
│ ├── benchmark_constituents_ft.asset_id
│ └── benchmarks_dm.id
└── stream key: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.date ]
Fragment 62563 (Actor 743628,743627)
StreamSyncLogStore
├── output:
│ ┌── benchmark_constituents_ft.benchmark_id
│ ├── benchmark_constituents_ft.date
│ ├── benchmark_constituents_ft.weight
│ ├── assets_dm_next.issue_currency_code
│ ├── benchmark_constituents_ft.asset_id
│ └── benchmarks_dm.id
├── stream key: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.date ]
└── StreamHashJoin { type: Inner, predicate: benchmark_constituents_ft.benchmark_id = benchmarks_dm.id }
├── output:
│ ┌── benchmark_constituents_ft.benchmark_id
│ ├── benchmark_constituents_ft.date
│ ├── benchmark_constituents_ft.weight
│ ├── assets_dm_next.issue_currency_code
│ ├── benchmark_constituents_ft.asset_id
│ └── benchmarks_dm.id
├── stream key: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.date ]
├── MergeExecutor
│ ├── output:
│ │ ┌── benchmark_constituents_ft.benchmark_id
│ │ ├── benchmark_constituents_ft.date
│ │ ├── benchmark_constituents_ft.weight
│ │ ├── assets_dm_next.issue_currency_code
│ │ ├── benchmark_constituents_ft.asset_id
│ │ └── assets_dm_next.id
│ └── stream key: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.date ]
└── MergeExecutor { output: [ benchmarks_dm.id ], stream key: [ benchmarks_dm.id ] }
Fragment 62564 (Actor 743643,743644)
StreamLocalityProvider { locality_columns: [benchmark_constituents_ft.benchmark_id] }
├── output:
│ ┌── benchmark_constituents_ft.benchmark_id
│ ├── benchmark_constituents_ft.date
│ ├── benchmark_constituents_ft.weight
│ ├── assets_dm_next.issue_currency_code
│ ├── benchmark_constituents_ft.asset_id
│ └── assets_dm_next.id
├── stream key: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.date ]
└── MergeExecutor
├── output:
│ ┌── benchmark_constituents_ft.benchmark_id
│ ├── benchmark_constituents_ft.date
│ ├── benchmark_constituents_ft.weight
│ ├── assets_dm_next.issue_currency_code
│ ├── benchmark_constituents_ft.asset_id
│ └── assets_dm_next.id
└── stream key: [ benchmark_constituents_ft.asset_id, benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date ]
Fragment 62565 (Actor 743646,743645)
StreamSyncLogStore
├── output:
│ ┌── benchmark_constituents_ft.benchmark_id
│ ├── benchmark_constituents_ft.date
│ ├── benchmark_constituents_ft.weight
│ ├── assets_dm_next.issue_currency_code
│ ├── benchmark_constituents_ft.asset_id
│ └── assets_dm_next.id
├── stream key: [ benchmark_constituents_ft.asset_id, benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date ]
└── StreamHashJoin { type: Inner, predicate: benchmark_constituents_ft.asset_id = assets_dm_next.id }
├── output:
│ ┌── benchmark_constituents_ft.benchmark_id
│ ├── benchmark_constituents_ft.date
│ ├── benchmark_constituents_ft.weight
│ ├── assets_dm_next.issue_currency_code
│ ├── benchmark_constituents_ft.asset_id
│ └── assets_dm_next.id
├── stream key: [ benchmark_constituents_ft.asset_id, benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date ]
├── MergeExecutor
│ ├── output: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.weight ]
│ └── stream key: [ benchmark_constituents_ft.asset_id, benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date ]
└── MergeExecutor { output: [ assets_dm_next.id, assets_dm_next.issue_currency_code ], stream key: [ assets_dm_next.id ] }
Fragment 62566 (Actor 743648,743647)
StreamLocalityProvider { locality_columns: [benchmark_constituents_ft.asset_id] }
├── output: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.weight ]
├── stream key: [ benchmark_constituents_ft.asset_id, benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date ]
└── MergeExecutor
├── output: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.weight ]
└── stream key: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, benchmark_constituents_ft.asset_id ]
Fragment 62567 (Actor 743042,743041)
StreamTableScan { table: benchmark_constituents_ft, columns: [benchmark_id, date, asset_id, weight] }
├── output: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.weight ]
├── stream key: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, benchmark_constituents_ft.asset_id ]
├── Upstream { output: [ benchmark_id, date, asset_id, weight ], stream key: [] }
└── BatchPlanNode { output: [ benchmark_id, date, asset_id, weight ], stream key: [] }
Fragment 62568 (Actor 743649,743650)
StreamFilter { predicate: Not(IsNull(assets_dm_next.issue_currency_code)) } { output: [ assets_dm_next.id, assets_dm_next.issue_currency_code ], stream key: [ assets_dm_next.id ] }
└── StreamTableScan { table: assets_dm_next, columns: [id, issue_currency_code] } { output: [ assets_dm_next.id, assets_dm_next.issue_currency_code ], stream key: [ assets_dm_next.id ] }
├── Upstream { output: [ id, issue_currency_code ], stream key: [] }
└── BatchPlanNode { output: [ id, issue_currency_code ], stream key: [] }
Fragment 62569 (Actor 743655,743656)
StreamProject { exprs: [benchmarks_dm.id] } { output: [ benchmarks_dm.id ], stream key: [ benchmarks_dm.id ] }
└── StreamFilter { predicate: IsNull(benchmarks_dm.disabled_at) } { output: [ benchmarks_dm.id, benchmarks_dm.disabled_at ], stream key: [ benchmarks_dm.id ] }
└── StreamTableScan { table: benchmarks_dm, columns: [id, disabled_at] } { output: [ benchmarks_dm.id, benchmarks_dm.disabled_at ], stream key: [ benchmarks_dm.id ] }
├── Upstream { output: [ id, disabled_at ], stream key: [] }
└── BatchPlanNode { output: [ id, disabled_at ], stream key: [] }