Job is idle — throughput ~0; structure shown.
Fragment 62402 (Actor 741997,741998)
StreamMaterialize { columns: [asset_id, as_of_date, provider_code, yield_to_maturity], stream_key: [asset_id, as_of_date, provider_code], pk_columns: [asset_id, as_of_date, provider_code], pk_conflict: Overwrite, watermark_columns: [as_of_date] }
├── output: [ public.asset_snapshots_ft.asset_id, public.asset_snapshots_ft.as_of_date, public.asset_snapshots_ft.provider_code, public.asset_snapshots_ft.yield_to_maturity ]
├── stream key: [ public.asset_snapshots_ft.asset_id, public.asset_snapshots_ft.as_of_date, public.asset_snapshots_ft.provider_code ]
└── StreamWatermarkFilter [upsert] { watermark_descs: [Desc { column: public.asset_snapshots_ft.as_of_date, expr: (public.asset_snapshots_ft.as_of_date - '5 years':Interval)::Date }], output_watermarks: [[public.asset_snapshots_ft.as_of_date]] }
├── output: [ public.asset_snapshots_ft.asset_id, public.asset_snapshots_ft.as_of_date, public.asset_snapshots_ft.provider_code, public.asset_snapshots_ft.yield_to_maturity ]
├── stream key: []
└── StreamUnion { all: true } { output: [ public.asset_snapshots_ft.asset_id, public.asset_snapshots_ft.as_of_date, public.asset_snapshots_ft.provider_code, public.asset_snapshots_ft.yield_to_maturity ], stream key: [] }
├── MergeExecutor
│ ├── output: [ public.asset_snapshots_ft.asset_id, public.asset_snapshots_ft.as_of_date, public.asset_snapshots_ft.provider_code, public.asset_snapshots_ft.yield_to_maturity ]
│ └── stream key: [ public.asset_snapshots_ft.asset_id, public.asset_snapshots_ft.as_of_date, public.asset_snapshots_ft.provider_code ]
├── MergeExecutor { output: [ asset_id, as_of_date, provider_code, yield_to_maturity ], stream key: [] }
└── StreamUpstreamSinkUnion { output: [ asset_id, as_of_date, provider_code, yield_to_maturity ], stream key: [] }
Fragment 62403 (Actor 741999)
StreamCdcTableScan { table: public.asset_snapshots_ft, columns: [asset_id, as_of_date, provider_code, yield_to_maturity] }
├── output: [ public.asset_snapshots_ft.asset_id, public.asset_snapshots_ft.as_of_date, public.asset_snapshots_ft.provider_code, public.asset_snapshots_ft.yield_to_maturity ]
├── stream key: [ public.asset_snapshots_ft.asset_id, public.asset_snapshots_ft.as_of_date, public.asset_snapshots_ft.provider_code ]
└── MergeExecutor { output: [ payload, _rw_offset, _rw_table_name ], stream key: [] }
Fragment 62404 (Actor 741820)
StreamCdcFilter { output: [ payload, _rw_offset, _rw_table_name ], stream key: [] }
└── Upstream { output: [ payload, _rw_offset, _rw_table_name ], stream key: [] }
Fragment 62405 (Actor 742003,742002)
StreamDml { columns: [asset_id, as_of_date, provider_code, yield_to_maturity] } { output: [ asset_id, as_of_date, provider_code, yield_to_maturity ], stream key: [] }
└── StreamSource { output: [ asset_id, as_of_date, provider_code, yield_to_maturity ], stream key: [] }