#![cfg_attr(not(feature = "rdkafka"), allow(dead_code))]
use std::sync::{
Arc,
atomic::{AtomicBool, AtomicI64, AtomicU64, Ordering},
};
#[derive(Debug, Clone, Default)]
pub struct KafkaMetrics {
inner: Arc<KafkaMetricsInner>,
}
#[derive(Debug, Default)]
struct KafkaMetricsInner {
assigned_partitions: AtomicU64,
revoked_partitions: AtomicU64,
lost_partitions: AtomicU64,
rebalances: AtomicU64,
emitted: AtomicU64,
committed_offsets: AtomicU64,
commit_failures: AtomicU64,
paused: AtomicBool,
outstanding: AtomicU64,
producer_in_flight: AtomicU64,
produced: AtomicU64,
delivery_failures: AtomicU64,
queue_full: AtomicU64,
high_watermark: AtomicI64,
committed_watermark: AtomicI64,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct KafkaMetricsSnapshot {
pub assigned_partitions: u64,
pub revoked_partitions: u64,
pub lost_partitions: u64,
pub rebalances: u64,
pub emitted: u64,
pub committed_offsets: u64,
pub commit_failures: u64,
pub paused: bool,
pub outstanding: u64,
pub producer_in_flight: u64,
pub produced: u64,
pub delivery_failures: u64,
pub queue_full: u64,
pub high_watermark: i64,
pub committed_watermark: i64,
}
impl KafkaMetrics {
#[must_use]
pub fn snapshot(&self) -> KafkaMetricsSnapshot {
KafkaMetricsSnapshot {
assigned_partitions: self.inner.assigned_partitions.load(Ordering::Relaxed),
revoked_partitions: self.inner.revoked_partitions.load(Ordering::Relaxed),
lost_partitions: self.inner.lost_partitions.load(Ordering::Relaxed),
rebalances: self.inner.rebalances.load(Ordering::Relaxed),
emitted: self.inner.emitted.load(Ordering::Relaxed),
committed_offsets: self.inner.committed_offsets.load(Ordering::Relaxed),
commit_failures: self.inner.commit_failures.load(Ordering::Relaxed),
paused: self.inner.paused.load(Ordering::Relaxed),
outstanding: self.inner.outstanding.load(Ordering::Relaxed),
producer_in_flight: self.inner.producer_in_flight.load(Ordering::Relaxed),
produced: self.inner.produced.load(Ordering::Relaxed),
delivery_failures: self.inner.delivery_failures.load(Ordering::Relaxed),
queue_full: self.inner.queue_full.load(Ordering::Relaxed),
high_watermark: self.inner.high_watermark.load(Ordering::Relaxed),
committed_watermark: self.inner.committed_watermark.load(Ordering::Relaxed),
}
}
pub(crate) fn add_assigned(&self, count: u64) {
self.inner
.assigned_partitions
.fetch_add(count, Ordering::Relaxed);
self.inner.rebalances.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn add_revoked(&self, count: u64) {
self.inner
.revoked_partitions
.fetch_add(count, Ordering::Relaxed);
self.inner.rebalances.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn add_lost(&self, count: u64) {
self.inner
.lost_partitions
.fetch_add(count, Ordering::Relaxed);
self.inner.rebalances.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn emitted(&self, outstanding: u64, high_watermark: i64) {
self.inner.emitted.fetch_add(1, Ordering::Relaxed);
self.inner.outstanding.store(outstanding, Ordering::Relaxed);
self.inner
.high_watermark
.store(high_watermark, Ordering::Relaxed);
}
pub(crate) fn emitted_batch(&self, count: u64, outstanding: u64, high_watermark: i64) {
self.inner.emitted.fetch_add(count, Ordering::Relaxed);
self.inner.outstanding.store(outstanding, Ordering::Relaxed);
self.inner
.high_watermark
.store(high_watermark, Ordering::Relaxed);
}
pub(crate) fn committed(&self, outstanding: u64, committed_watermark: i64) {
self.inner.committed_offsets.fetch_add(1, Ordering::Relaxed);
self.inner.outstanding.store(outstanding, Ordering::Relaxed);
self.inner
.committed_watermark
.store(committed_watermark, Ordering::Relaxed);
}
pub(crate) fn set_outstanding(&self, outstanding: u64) {
self.inner.outstanding.store(outstanding, Ordering::Relaxed);
}
pub(crate) fn commit_failed(&self) {
self.inner.commit_failures.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn set_paused(&self, paused: bool) {
self.inner.paused.store(paused, Ordering::Relaxed);
}
pub(crate) fn producer_in_flight(&self, in_flight: u64) {
self.inner
.producer_in_flight
.store(in_flight, Ordering::Relaxed);
}
pub(crate) fn produced(&self) {
self.inner.produced.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn produced_batch(&self, count: u64) {
self.inner.produced.fetch_add(count, Ordering::Relaxed);
}
pub(crate) fn delivery_failed(&self) {
self.inner.delivery_failures.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn delivery_failed_batch(&self, count: u64) {
self.inner
.delivery_failures
.fetch_add(count, Ordering::Relaxed);
}
pub(crate) fn queue_full(&self) {
self.inner.queue_full.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn replace_from(&self, snapshot: KafkaMetricsSnapshot) {
self.inner
.assigned_partitions
.store(snapshot.assigned_partitions, Ordering::Relaxed);
self.inner
.revoked_partitions
.store(snapshot.revoked_partitions, Ordering::Relaxed);
self.inner
.lost_partitions
.store(snapshot.lost_partitions, Ordering::Relaxed);
self.inner
.rebalances
.store(snapshot.rebalances, Ordering::Relaxed);
self.inner
.emitted
.store(snapshot.emitted, Ordering::Relaxed);
self.inner
.committed_offsets
.store(snapshot.committed_offsets, Ordering::Relaxed);
self.inner
.commit_failures
.store(snapshot.commit_failures, Ordering::Relaxed);
self.inner.paused.store(snapshot.paused, Ordering::Relaxed);
self.inner
.outstanding
.store(snapshot.outstanding, Ordering::Relaxed);
self.inner
.producer_in_flight
.store(snapshot.producer_in_flight, Ordering::Relaxed);
self.inner
.produced
.store(snapshot.produced, Ordering::Relaxed);
self.inner
.delivery_failures
.store(snapshot.delivery_failures, Ordering::Relaxed);
self.inner
.queue_full
.store(snapshot.queue_full, Ordering::Relaxed);
self.inner
.high_watermark
.store(snapshot.high_watermark, Ordering::Relaxed);
self.inner
.committed_watermark
.store(snapshot.committed_watermark, Ordering::Relaxed);
}
}