Job is idle — throughput ~0; structure shown.
Fragment 23382 (Actor 102000,101999)
StreamMaterialize { columns: [asset_id, as_of_timestamp, last, value_timestamp], stream_key: [asset_id, as_of_timestamp, value_timestamp], pk_columns: [asset_id, as_of_timestamp], pk_conflict: Overwrite, watermark_columns: [value_timestamp] }
├── output: [ public.intraday_asset_prices_ft.asset_id, public.intraday_asset_prices_ft.as_of_timestamp, public.intraday_asset_prices_ft.last, public.intraday_asset_prices_ft.value_timestamp ]
├── stream key: [ public.intraday_asset_prices_ft.asset_id, public.intraday_asset_prices_ft.as_of_timestamp, public.intraday_asset_prices_ft.value_timestamp ]
└── StreamWatermarkFilter [upsert] { watermark_descs: [Desc { column: public.intraday_asset_prices_ft.value_timestamp, expr: SubtractWithTimeZone(public.intraday_asset_prices_ft.value_timestamp, '00:30:00':Interval, 'UTC':Varchar) }], output_watermarks: [[public.intraday_asset_prices_ft.value_timestamp]] }
├── output: [ public.intraday_asset_prices_ft.asset_id, public.intraday_asset_prices_ft.as_of_timestamp, public.intraday_asset_prices_ft.last, public.intraday_asset_prices_ft.value_timestamp ]
├── stream key: []
└── StreamUnion { all: true } { output: [ public.intraday_asset_prices_ft.asset_id, public.intraday_asset_prices_ft.as_of_timestamp, public.intraday_asset_prices_ft.last, public.intraday_asset_prices_ft.value_timestamp ], stream key: [] }
├── MergeExecutor
│ ├── output: [ public.intraday_asset_prices_ft.asset_id, public.intraday_asset_prices_ft.as_of_timestamp, public.intraday_asset_prices_ft.last, public.intraday_asset_prices_ft.value_timestamp ]
│ └── stream key: [ public.intraday_asset_prices_ft.asset_id, public.intraday_asset_prices_ft.as_of_timestamp ]
├── MergeExecutor { output: [ asset_id, as_of_timestamp, last, value_timestamp ], stream key: [] }
└── StreamUpstreamSinkUnion { output: [ asset_id, as_of_timestamp, last, value_timestamp ], stream key: [] }
Fragment 23383 (Actor 102001)
StreamCdcTableScan { table: public.intraday_asset_prices_ft, columns: [asset_id, as_of_timestamp, last, value_timestamp] }
├── output: [ public.intraday_asset_prices_ft.asset_id, public.intraday_asset_prices_ft.as_of_timestamp, public.intraday_asset_prices_ft.last, public.intraday_asset_prices_ft.value_timestamp ]
├── stream key: [ public.intraday_asset_prices_ft.asset_id, public.intraday_asset_prices_ft.as_of_timestamp ]
└── MergeExecutor { output: [ payload, _rw_offset, _rw_table_name ], stream key: [] }
Fragment 23384 (Actor 101998)
StreamCdcFilter { output: [ payload, _rw_offset, _rw_table_name ], stream key: [] }
└── Upstream { output: [ payload, _rw_offset, _rw_table_name ], stream key: [] }
Fragment 23385 (Actor 102003,102002)
StreamDml { columns: [asset_id, as_of_timestamp, last, value_timestamp] } { output: [ asset_id, as_of_timestamp, last, value_timestamp ], stream key: [] }
└── StreamSource { output: [ asset_id, as_of_timestamp, last, value_timestamp ], stream key: [] }