From a4a769acb1e1138a439d3e5dbc013f7ea060f41f Mon Sep 17 00:00:00 2001 From: tryangul <11639460+tryangul@users.noreply.github.com> Date: Tue, 1 Sep 2026 11:24:27 -0700 Subject: [PATCH] feat(metrics): DLQ produce failures no increment a counter and log instead of panicking. Fix partition header bug. --- rust-arroyo/src/processing/dlq.rs | 9 +++++---- 1 file changed, 5 insertions(+), 4 deletions(-) diff --git a/rust-arroyo/src/processing/dlq.rs b/rust-arroyo/src/processing/dlq.rs index 886668e4..eaa90f14 100644 --- a/rust-arroyo/src/processing/dlq.rs +++ b/rust-arroyo/src/processing/dlq.rs @@ -84,7 +84,7 @@ impl DlqProducer for KafkaDlqProducer { headers = headers.insert( "original_partition", - Some(message.offset.to_string().into_bytes()), + Some(message.partition.index.to_string().into_bytes()), ); headers = headers.insert( "original_offset", @@ -98,9 +98,10 @@ impl DlqProducer for KafkaDlqProducer { ); Box::pin(async move { - producer - .produce(&topic, payload) - .expect("Message was produced"); + if let Err(err) = producer.produce(&topic, payload) { + counter!("arroyo.consumer.dlq.produce_error", 1); + tracing::error!("Failed to produce to DLQ: {:?}", err); + } message })