Job is idle — throughput ~0; structure shown.
Fragment 17825 (Actor 156642,156641)
StreamMaterialize { columns: [asset_id, date, currency_code, close, nav, last, reference_price, reference_price_source], stream_key: [asset_id, date], pk_columns: [asset_id, date], pk_conflict: Overwrite, watermark_columns: [date] }
├── output:
│ ┌── public.asset_prices_eod_ft.asset_id
│ ├── public.asset_prices_eod_ft.date
│ ├── public.asset_prices_eod_ft.currency_code
│ ├── public.asset_prices_eod_ft.close
│ ├── public.asset_prices_eod_ft.nav
│ ├── public.asset_prices_eod_ft.last
│ ├── public.asset_prices_eod_ft.reference_price
│ └── public.asset_prices_eod_ft.reference_price_source
├── stream key: [ public.asset_prices_eod_ft.asset_id, public.asset_prices_eod_ft.date ]
└── StreamWatermarkFilter [upsert] { watermark_descs: [Desc { column: public.asset_prices_eod_ft.date, expr: (public.asset_prices_eod_ft.date - '5 years':Interval)::Date }], output_watermarks: [[public.asset_prices_eod_ft.date]] }
├── output:
│ ┌── public.asset_prices_eod_ft.asset_id
│ ├── public.asset_prices_eod_ft.date
│ ├── public.asset_prices_eod_ft.currency_code
│ ├── public.asset_prices_eod_ft.close
│ ├── public.asset_prices_eod_ft.nav
│ ├── public.asset_prices_eod_ft.last
│ ├── public.asset_prices_eod_ft.reference_price
│ └── public.asset_prices_eod_ft.reference_price_source
├── stream key: []
└── StreamUnion { all: true }
├── output:
│ ┌── public.asset_prices_eod_ft.asset_id
│ ├── public.asset_prices_eod_ft.date
│ ├── public.asset_prices_eod_ft.currency_code
│ ├── public.asset_prices_eod_ft.close
│ ├── public.asset_prices_eod_ft.nav
│ ├── public.asset_prices_eod_ft.last
│ ├── public.asset_prices_eod_ft.reference_price
│ └── public.asset_prices_eod_ft.reference_price_source
├── stream key: []
├── MergeExecutor
│ ├── output:
│ │ ┌── public.asset_prices_eod_ft.asset_id
│ │ ├── public.asset_prices_eod_ft.date
│ │ ├── public.asset_prices_eod_ft.currency_code
│ │ ├── public.asset_prices_eod_ft.close
│ │ ├── public.asset_prices_eod_ft.nav
│ │ ├── public.asset_prices_eod_ft.last
│ │ ├── public.asset_prices_eod_ft.reference_price
│ │ └── public.asset_prices_eod_ft.reference_price_source
│ └── stream key: [ public.asset_prices_eod_ft.asset_id, public.asset_prices_eod_ft.date ]
├── MergeExecutor { output: [ asset_id, date, currency_code, close, nav, last, reference_price, reference_price_source ], stream key: [] }
└── StreamUpstreamSinkUnion { output: [ asset_id, date, currency_code, close, nav, last, reference_price, reference_price_source ], stream key: [] }
Fragment 17826 (Actor 156653)
StreamCdcTableScan { table: public.asset_prices_eod_ft, columns: [asset_id, date, currency_code, close, nav, last, reference_price, reference_price_source] }
├── output:
│ ┌── public.asset_prices_eod_ft.asset_id
│ ├── public.asset_prices_eod_ft.date
│ ├── public.asset_prices_eod_ft.currency_code
│ ├── public.asset_prices_eod_ft.close
│ ├── public.asset_prices_eod_ft.nav
│ ├── public.asset_prices_eod_ft.last
│ ├── public.asset_prices_eod_ft.reference_price
│ └── public.asset_prices_eod_ft.reference_price_source
├── stream key: [ public.asset_prices_eod_ft.asset_id, public.asset_prices_eod_ft.date ]
└── MergeExecutor { output: [ payload, _rw_offset, _rw_table_name ], stream key: [] }
Fragment 17827 (Actor 156855)
StreamCdcFilter { output: [ payload, _rw_offset, _rw_table_name ], stream key: [] }
└── Upstream { output: [ payload, _rw_offset, _rw_table_name ], stream key: [] }
Fragment 17828 (Actor 156654,156655)
StreamDml { columns: [asset_id, date, currency_code, close, nav, last, reference_price, reference_price_source] }
├── output: [ asset_id, date, currency_code, close, nav, last, reference_price, reference_price_source ]
├── stream key: []
└── StreamSource { output: [ asset_id, date, currency_code, close, nav, last, reference_price, reference_price_source ], stream key: [] }