Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions crates/cli/src/observability.rs
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@ pub struct Metrics {
pub cdc_events_processed: Counter<u64>,
pub cdc_events_failed: Counter<u64>,
pub cdc_batch_duration: Histogram<f64>,
pub cdc_transform_duration: Histogram<f64>,
pub backfill_rows_processed: Counter<u64>,
pub backfill_batches_failed: Counter<u64>,
pub backfill_batch_duration: Histogram<f64>,
Expand Down Expand Up @@ -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")
Expand Down
5 changes: 5 additions & 0 deletions crates/cli/src/pipeline/streaming.rs
Original file line number Diff line number Diff line change
Expand Up @@ -87,7 +87,12 @@ async fn process_config_events(
events_processed: &mut HashMap<String, u64>,
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) => {
Expand Down