diff --git a/crates/cli/src/observability.rs b/crates/cli/src/observability.rs index 6dab891..80fc119 100644 --- a/crates/cli/src/observability.rs +++ b/crates/cli/src/observability.rs @@ -43,6 +43,7 @@ pub struct Metrics { pub cdc_events_processed: Counter, pub cdc_events_failed: Counter, pub cdc_batch_duration: Histogram, + pub cdc_transform_duration: Histogram, pub backfill_rows_processed: Counter, pub backfill_batches_failed: Counter, pub backfill_batch_duration: Histogram, @@ -218,6 +219,10 @@ fn build_metrics(meter: &Meter) -> Metrics { .f64_histogram("puffgres.cdc.batch_duration_ms") .with_description("Time to process one CDC batch") .build(), + cdc_transform_duration: meter + .f64_histogram("puffgres.cdc.transform_duration_ms") + .with_description("Time spent in transform_batch for one CDC config batch") + .build(), backfill_rows_processed: meter .u64_counter("puffgres.backfill.rows_processed") .with_description("Backfill rows processed") diff --git a/crates/cli/src/pipeline/streaming.rs b/crates/cli/src/pipeline/streaming.rs index 1912c24..0fa62bc 100644 --- a/crates/cli/src/pipeline/streaming.rs +++ b/crates/cli/src/pipeline/streaming.rs @@ -87,7 +87,12 @@ async fn process_config_events( events_processed: &mut HashMap, dlq_lsn: u64, ) -> Result<(), CliError> { + let transform_start = std::time::Instant::now(); let transform_result = transformer.transform_batch(events).await; + if let Some(m) = metrics { + m.cdc_transform_duration + .record(transform_start.elapsed().as_millis() as f64, &[]); + } match transform_result { Err(e) => {