Skip to content
Merged
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
1 change: 1 addition & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ ssl = ["rdkafka/ssl"]
[dependencies]
chrono = "0.4.26"
coarsetime = "0.1.33"
metrics = "0.24"
once_cell = "1.18.0"
rand = "0.8.5"
rdkafka = { version = ">=0.37.0,<0.40", features = ["cmake-build", "tracing"] }
Expand Down
14 changes: 5 additions & 9 deletions rust-arroyo/src/backends/kafka/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,6 @@ use super::CommitOffsets;
use super::Consumer as ArroyoConsumer;
use super::ConsumerError;
use crate::backends::kafka::types::KafkaPayload;
use crate::gauge;
use crate::types::{BrokerMessage, Partition, Topic};
use chrono::{DateTime, Utc};
use parking_lot::Mutex;
Expand Down Expand Up @@ -218,18 +217,15 @@ impl<C: AssignmentCallbacks + Send + Sync> ClientContext for CustomContext<C> {
}

fn stats(&self, stats: Statistics) {
gauge!(
"arroyo.consumer.librdkafka.total_queue_size",
stats.replyq as u64,
);
metrics::gauge!("arroyo.consumer.librdkafka.total_queue_size").set(stats.replyq as f64);
for (topic_name, topic) in stats.topics.iter() {
for (partition_num, partition) in topic.partitions.iter() {
gauge!(
metrics::gauge!(
"arroyo.consumer.librdkafka.fetch_queue_count",
partition.fetchq_cnt as u64,
"topic" => topic_name,
"topic" => topic_name.clone(),
"partition" => partition_num.to_string()
);
)
.set(partition.fetchq_cnt as f64);
}
}
}
Expand Down
105 changes: 60 additions & 45 deletions rust-arroyo/src/backends/kafka/producer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,8 +5,6 @@ use crate::backends::ProducerError;
use crate::backends::{
AsyncProducer as ArroyoAsyncProducer, Producer as ArroyoProducer, ProducerFuture,
};
use crate::counter;
use crate::gauge;
use crate::types::TopicOrPartition;
use rdkafka::client::ClientContext;
use rdkafka::config::ClientConfig;
Expand Down Expand Up @@ -40,80 +38,80 @@ impl ClientContext for ProducerContext {
// Record broker latency metrics
if let Some(int_latency) = &broker_stats.int_latency {
let p99_latency_ms = int_latency.p99 as f64 / 1000.0;
gauge!(
metrics::gauge!(
"arroyo.producer.librdkafka.p99_int_latency",
p99_latency_ms as u64,
"broker_id" => broker_id_str.clone(),
"producer_name" => producer_name
);
"producer_name" => producer_name.to_owned()
)
.set(p99_latency_ms as u64 as f64);
}

if let Some(outbuf_latency) = &broker_stats.outbuf_latency {
let p99_latency_ms = outbuf_latency.p99 as f64 / 1000.0;
gauge!(
metrics::gauge!(
"arroyo.producer.librdkafka.p99_outbuf_latency",
p99_latency_ms as u64,
"broker_id" => broker_id_str.clone(),
"producer_name" => producer_name
);
"producer_name" => producer_name.to_owned()
)
.set(p99_latency_ms as u64 as f64);
}

if let Some(rtt) = &broker_stats.rtt {
let p99_rtt_ms = rtt.p99 as f64 / 1000.0;
gauge!(
metrics::gauge!(
"arroyo.producer.librdkafka.p99_rtt",
p99_rtt_ms as u64,
"broker_id" => broker_id_str.clone(),
"producer_name" => producer_name
);
"producer_name" => producer_name.to_owned()
)
.set(p99_rtt_ms as u64 as f64);
}

// Record broker transmission error metrics
gauge!(
metrics::gauge!(
"arroyo.producer.librdkafka.broker_txerrs",
broker_stats.txerrs as i64,
"broker_id" => broker_id_str.clone(),
"producer_name" => producer_name
);
"producer_name" => producer_name.to_owned()
)
.set(broker_stats.txerrs as f64);

gauge!(
metrics::gauge!(
"arroyo.producer.librdkafka.broker_txretries",
broker_stats.txretries as i64,
"broker_id" => broker_id_str,
"producer_name" => producer_name
);
"producer_name" => producer_name.to_owned()
)
.set(broker_stats.txretries as f64);
}

// Record global producer metrics
gauge!(
metrics::gauge!(
"arroyo.producer.librdkafka.message_count",
stats.msg_cnt as i64,
"producer_name" => producer_name
);
"producer_name" => producer_name.to_owned()
)
.set(stats.msg_cnt as f64);

gauge!(
metrics::gauge!(
"arroyo.producer.librdkafka.message_count_max",
stats.msg_max as i64,
"producer_name" => producer_name
);
"producer_name" => producer_name.to_owned()
)
.set(stats.msg_max as f64);

gauge!(
metrics::gauge!(
"arroyo.producer.librdkafka.message_size",
stats.msg_size as i64,
"producer_name" => producer_name
);
"producer_name" => producer_name.to_owned()
)
.set(stats.msg_size as f64);

gauge!(
metrics::gauge!(
"arroyo.producer.librdkafka.message_size_max",
stats.msg_size_max as i64,
"producer_name" => producer_name
);
"producer_name" => producer_name.to_owned()
)
.set(stats.msg_size_max as f64);

gauge!(
metrics::gauge!(
"arroyo.producer.librdkafka.reply_queue_size",
stats.replyq as i64,
"producer_name" => producer_name
);
"producer_name" => producer_name.to_owned()
)
.set(stats.replyq as f64);
}
}

Expand All @@ -130,7 +128,12 @@ impl RdkafkaProducerContext for ProducerContext {
Err((err, _)) => get_error_name(err),
};
let producer_name = self.get_producer_name();
counter!("arroyo.producer.produce_status", 1, "status" => result, "producer_name" => producer_name);
metrics::counter!(
"arroyo.producer.produce_status",
"status" => result,
"producer_name" => producer_name.to_owned()
)
.increment(1);
}
}

Expand Down Expand Up @@ -204,13 +207,25 @@ fn record_producer_error(
let producer_error = ProducerError::ProducerFailure {
error: error_name.clone(),
};
counter!("arroyo.producer.produce_status", 1, "status" => "error", "code" => error_name, "producer_name" => producer_name);
metrics::counter!(
"arroyo.producer.produce_status",
"status" => "error",
"code" => error_name,
"producer_name" => producer_name.to_owned()
)
.increment(1);
return producer_error;
}
let producer_error = ProducerError::ProducerFailure {
error: default_error.to_string(),
};
counter!("arroyo.producer.produce_status", 1, "status" => "error", "code" => default_error, "producer_name" => producer_name);
metrics::counter!(
"arroyo.producer.produce_status",
"status" => "error",
"code" => default_error.to_owned(),
"producer_name" => producer_name.to_owned()
)
.increment(1);
producer_error
}

Expand Down
1 change: 0 additions & 1 deletion rust-arroyo/src/lib.rs
Original file line number Diff line number Diff line change
@@ -1,5 +1,4 @@
pub mod backends;
pub mod metrics;
pub mod processing;
pub mod testutils;
pub mod types;
Expand Down
44 changes: 0 additions & 44 deletions rust-arroyo/src/metrics/globals.rs

This file was deleted.

62 changes: 0 additions & 62 deletions rust-arroyo/src/metrics/macros.rs

This file was deleted.

8 changes: 0 additions & 8 deletions rust-arroyo/src/metrics/mod.rs

This file was deleted.

Loading
Loading