use std::{
collections::{BTreeMap, BTreeSet},
sync::{
Arc, Condvar, Mutex,
atomic::{AtomicU64, Ordering},
},
};
use bytes::Bytes;
use crate::native::{
KafkaClientError, KafkaClientResult,
profile::{self, ProfileBucket},
};
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct TopicPartition {
pub topic: String,
pub partition: i32,
}
impl TopicPartition {
#[must_use]
pub fn new(topic: impl Into<String>, partition: i32) -> Self {
Self {
topic: topic.into(),
partition,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum StartOffset {
Beginning,
End,
Timestamp(i64),
Committed,
Offset(i64),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TopicPartitionAssignment {
pub topic: String,
pub partition: i32,
pub start_offset: StartOffset,
}
impl TopicPartitionAssignment {
#[must_use]
pub fn beginning(topic: impl Into<String>, partition: i32) -> Self {
Self {
topic: topic.into(),
partition,
start_offset: StartOffset::Beginning,
}
}
#[must_use]
pub fn end(topic: impl Into<String>, partition: i32) -> Self {
Self {
topic: topic.into(),
partition,
start_offset: StartOffset::End,
}
}
#[must_use]
pub fn timestamp(topic: impl Into<String>, partition: i32, timestamp_ms: i64) -> Self {
Self {
topic: topic.into(),
partition,
start_offset: StartOffset::Timestamp(timestamp_ms),
}
}
#[must_use]
pub fn committed(topic: impl Into<String>, partition: i32) -> Self {
Self {
topic: topic.into(),
partition,
start_offset: StartOffset::Committed,
}
}
#[must_use]
pub fn absolute(topic: impl Into<String>, partition: i32, offset: i64) -> Self {
Self {
topic: topic.into(),
partition,
start_offset: StartOffset::Offset(offset),
}
}
#[must_use]
pub fn topic_partition(&self) -> TopicPartition {
TopicPartition::new(self.topic.clone(), self.partition)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum KafkaTimestamp {
NotAvailable,
CreateTime(i64),
LogAppendTime(i64),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct KafkaPayloadRecord {
topic_index: usize,
pub partition: i32,
pub offset: i64,
pub timestamp: KafkaTimestamp,
payload_start: usize,
payload_end: usize,
}
impl KafkaPayloadRecord {
pub(crate) fn new(
topic_index: usize,
partition: i32,
offset: i64,
timestamp: KafkaTimestamp,
payload_start: usize,
payload_end: usize,
) -> Self {
Self {
topic_index,
partition,
offset,
timestamp,
payload_start,
payload_end,
}
}
#[must_use]
pub fn payload<'a>(&self, batch: &'a KafkaPayloadBatch) -> &'a [u8] {
&batch.payloads[self.payload_start..self.payload_end]
}
#[must_use]
pub fn payload_range(&self) -> std::ops::Range<usize> {
self.payload_start..self.payload_end
}
#[must_use]
pub fn topic<'a>(&self, batch: &'a KafkaPayloadBatch) -> &'a str {
&batch.topics[self.topic_index]
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct OffsetCommit {
pub topic: String,
pub partition: i32,
pub offset: i64,
}
impl OffsetCommit {
#[must_use]
pub fn topic_partition(&self) -> TopicPartition {
TopicPartition::new(self.topic.clone(), self.partition)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) struct OffsetRange {
pub(crate) partition: i32,
pub(crate) first_offset: i64,
pub(crate) next_offset: i64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct TopicPartitionOffsetRange {
pub(crate) topic: String,
pub(crate) partition: i32,
pub(crate) first_offset: i64,
pub(crate) next_offset: i64,
}
impl TopicPartitionOffsetRange {
fn topic_partition(&self) -> TopicPartition {
TopicPartition::new(self.topic.clone(), self.partition)
}
}
#[derive(Debug, Clone, Default)]
pub struct NativeKafkaMetrics {
inner: Arc<NativeKafkaMetricsInner>,
}
#[derive(Debug, Default)]
struct NativeKafkaMetricsInner {
assigned_partitions: AtomicU64,
revoked_partitions: AtomicU64,
lost_partitions: AtomicU64,
rebalances: AtomicU64,
emitted: AtomicU64,
committed_offsets: AtomicU64,
outstanding: AtomicU64,
high_watermark: AtomicU64,
committed_watermark: AtomicU64,
offset_out_of_range_resets: AtomicU64,
broker_commit_requests: AtomicU64,
broker_commit_failures: AtomicU64,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct NativeKafkaMetricsSnapshot {
pub assigned_partitions: u64,
pub revoked_partitions: u64,
pub lost_partitions: u64,
pub rebalances: u64,
pub emitted: u64,
pub committed_offsets: u64,
pub outstanding: u64,
pub high_watermark: u64,
pub committed_watermark: u64,
pub offset_out_of_range_resets: u64,
pub broker_commit_requests: u64,
pub broker_commit_failures: u64,
}
impl NativeKafkaMetrics {
pub(crate) fn add_assigned(&self, count: u64) {
if count == 0 {
return;
}
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) {
if count == 0 {
return;
}
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) {
if count == 0 {
return;
}
self.inner
.lost_partitions
.fetch_add(count, Ordering::Relaxed);
self.inner.rebalances.fetch_add(1, 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.max(0) as u64, Ordering::Relaxed);
}
pub(crate) fn committed(&self, count: u64, outstanding: u64, committed_watermark: i64) {
self.inner
.committed_offsets
.fetch_add(count, Ordering::Relaxed);
self.inner.outstanding.store(outstanding, Ordering::Relaxed);
self.inner
.committed_watermark
.store(committed_watermark.max(0) as u64, Ordering::Relaxed);
}
pub(crate) fn offset_out_of_range_reset(&self) {
self.inner
.offset_out_of_range_resets
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn broker_commit_request(&self) {
self.inner
.broker_commit_requests
.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn broker_commit_failed(&self) {
self.inner
.broker_commit_failures
.fetch_add(1, Ordering::Relaxed);
}
#[must_use]
pub fn snapshot(&self) -> NativeKafkaMetricsSnapshot {
NativeKafkaMetricsSnapshot {
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),
outstanding: self.inner.outstanding.load(Ordering::Relaxed),
high_watermark: self.inner.high_watermark.load(Ordering::Relaxed),
committed_watermark: self.inner.committed_watermark.load(Ordering::Relaxed),
offset_out_of_range_resets: self
.inner
.offset_out_of_range_resets
.load(Ordering::Relaxed),
broker_commit_requests: self.inner.broker_commit_requests.load(Ordering::Relaxed),
broker_commit_failures: self.inner.broker_commit_failures.load(Ordering::Relaxed),
}
}
}
#[derive(Debug, Clone)]
pub struct KafkaPayloadBatch {
topics: Vec<String>,
records: Vec<KafkaPayloadRecord>,
payloads: Bytes,
watermarks: Vec<OffsetCommit>,
ranges: Vec<TopicPartitionOffsetRange>,
active_partitions: Vec<TopicPartition>,
commit_state: Option<Arc<LocalCommitState>>,
}
#[derive(Debug, Clone)]
pub struct KafkaPayloadBatchParts {
pub topics: Vec<String>,
pub records: Vec<KafkaPayloadRecord>,
pub payloads: Bytes,
pub watermarks: Vec<OffsetCommit>,
pub active_partitions: Vec<TopicPartition>,
}
impl KafkaPayloadBatch {
pub(crate) fn new(
topics: Vec<String>,
records: Vec<KafkaPayloadRecord>,
payloads: Vec<u8>,
watermarks: Vec<OffsetCommit>,
ranges: Vec<TopicPartitionOffsetRange>,
active_partitions: Vec<TopicPartition>,
) -> Self {
Self {
topics,
records,
payloads: Bytes::from(payloads),
watermarks,
ranges,
active_partitions,
commit_state: None,
}
}
pub(crate) fn with_commit_state(mut self, commit_state: Arc<LocalCommitState>) -> Self {
profile::measure(ProfileBucket::CommitBookkeeping, || {
commit_state.observe(&self.ranges);
});
self.commit_state = Some(commit_state);
self
}
#[must_use]
pub fn records(&self) -> &[KafkaPayloadRecord] {
&self.records
}
#[must_use]
pub fn topics(&self) -> &[String] {
&self.topics
}
#[must_use]
pub fn len(&self) -> usize {
self.records.len()
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.records.is_empty()
}
#[must_use]
pub fn payload(&self, record: &KafkaPayloadRecord) -> &[u8] {
record.payload(self)
}
#[must_use]
pub fn watermarks(&self) -> &[OffsetCommit] {
&self.watermarks
}
#[must_use]
pub fn active_partitions(&self) -> &[TopicPartition] {
&self.active_partitions
}
#[must_use]
pub fn into_parts(self) -> KafkaPayloadBatchParts {
KafkaPayloadBatchParts {
topics: self.topics,
records: self.records,
payloads: self.payloads,
watermarks: self.watermarks,
active_partitions: self.active_partitions,
}
}
pub fn commit(&self) -> KafkaClientResult<()> {
let Some(commit_state) = &self.commit_state else {
return Err(KafkaClientError::InvalidConfig(
"native Kafka payload batch has no local commit state".to_owned(),
));
};
profile::measure(ProfileBucket::CommitBookkeeping, || {
commit_state.process(&self.ranges)
})
}
pub fn commit_offset(
&self,
topic: impl Into<String>,
partition: i32,
offset: i64,
) -> KafkaClientResult<()> {
let Some(commit_state) = &self.commit_state else {
return Err(KafkaClientError::InvalidConfig(
"native Kafka payload batch has no local commit state".to_owned(),
));
};
profile::measure(ProfileBucket::CommitBookkeeping, || {
commit_state.process(&[TopicPartitionOffsetRange {
topic: topic.into(),
partition,
first_offset: offset,
next_offset: offset + 1,
}])
})
}
}
#[derive(Debug, Default)]
pub(crate) struct PayloadBatchBuilder {
topics: Vec<String>,
topic_indexes: BTreeMap<String, usize>,
records: Vec<KafkaPayloadRecord>,
payloads: Vec<u8>,
watermarks: Vec<OffsetCommit>,
ranges: Vec<TopicPartitionOffsetRange>,
partition_index_topic: Option<String>,
partition_index: Vec<Option<usize>>,
include_timestamps: bool,
}
impl PayloadBatchBuilder {
pub(crate) fn with_capacity(
records: usize,
payload_bytes: usize,
include_timestamps: bool,
) -> Self {
Self {
records: Vec::with_capacity(records),
topics: Vec::new(),
topic_indexes: BTreeMap::new(),
payloads: Vec::with_capacity(payload_bytes),
watermarks: Vec::new(),
ranges: Vec::new(),
partition_index_topic: None,
partition_index: Vec::new(),
include_timestamps,
}
}
pub(crate) fn len(&self) -> usize {
self.records.len()
}
pub(crate) fn include_timestamps(&self) -> bool {
self.include_timestamps
}
pub(crate) fn push_record(
&mut self,
topic: &str,
partition: i32,
offset: i64,
timestamp: KafkaTimestamp,
payload: &[u8],
) {
let topic_index = self.topic_index(topic);
let payload_start = self.payloads.len();
self.payloads.extend_from_slice(payload);
let payload_end = self.payloads.len();
self.records.push(KafkaPayloadRecord::new(
topic_index,
partition,
offset,
timestamp,
payload_start,
payload_end,
));
}
pub(crate) fn observe_partition(&mut self, topic: &str, range: OffsetRange) {
profile::measure(ProfileBucket::OffsetBookkeeping, || {
self.observe(topic, range);
});
}
pub(crate) fn finish(
self,
active_partitions: Vec<TopicPartition>,
) -> Option<KafkaPayloadBatch> {
if self.records.is_empty() {
return None;
}
Some(KafkaPayloadBatch::new(
self.topics,
self.records,
self.payloads,
self.watermarks,
self.ranges,
active_partitions,
))
}
fn observe(&mut self, topic: &str, range: OffsetRange) {
if let Some(index) = self.cached_partition_index(topic, range.partition) {
let watermark = &mut self.watermarks[index];
watermark.offset = watermark.offset.max(range.next_offset);
let existing = &mut self.ranges[index];
existing.first_offset = existing.first_offset.min(range.first_offset);
existing.next_offset = existing.next_offset.max(range.next_offset);
return;
}
if let Some(index) = self.watermarks.iter().position(|watermark| {
watermark.topic == topic && watermark.partition == range.partition
}) {
let watermark = &mut self.watermarks[index];
watermark.offset = watermark.offset.max(range.next_offset);
let existing = &mut self.ranges[index];
existing.first_offset = existing.first_offset.min(range.first_offset);
existing.next_offset = existing.next_offset.max(range.next_offset);
self.store_partition_index(topic, range.partition, index);
return;
}
let index = self.watermarks.len();
self.watermarks.push(OffsetCommit {
topic: topic.to_owned(),
partition: range.partition,
offset: range.next_offset,
});
self.ranges.push(TopicPartitionOffsetRange {
topic: topic.to_owned(),
partition: range.partition,
first_offset: range.first_offset,
next_offset: range.next_offset,
});
self.store_partition_index(topic, range.partition, index);
}
fn topic_index(&mut self, topic: &str) -> usize {
if let Some(index) = self.topic_indexes.get(topic) {
return *index;
}
let index = self.topics.len();
self.topics.push(topic.to_owned());
self.topic_indexes.insert(topic.to_owned(), index);
index
}
fn cached_partition_index(&self, topic: &str, partition: i32) -> Option<usize> {
let partition = usize::try_from(partition).ok()?;
if self.partition_index_topic.as_deref() != Some(topic) {
return None;
}
self.partition_index.get(partition).and_then(|index| *index)
}
fn store_partition_index(&mut self, topic: &str, partition: i32, index: usize) {
let Ok(partition) = usize::try_from(partition) else {
return;
};
match self.partition_index_topic.as_deref() {
None => self.partition_index_topic = Some(topic.to_owned()),
Some(existing) if existing == topic => {}
Some(_) => return,
}
if self.partition_index.len() <= partition {
self.partition_index.resize(partition + 1, None);
}
self.partition_index[partition] = Some(index);
}
}
#[derive(Debug, Default)]
pub(crate) struct LocalCommitState {
inner: Mutex<BTreeMap<TopicPartition, PartitionCommitState>>,
lost: Mutex<BTreeSet<TopicPartition>>,
progress: Condvar,
}
#[derive(Debug, Default)]
struct PartitionCommitState {
emitted_next: Option<i64>,
processed_next: Option<i64>,
committed_next: Option<i64>,
processed_ranges: BTreeMap<i64, i64>,
}
impl LocalCommitState {
pub(crate) fn observe(&self, ranges: &[TopicPartitionOffsetRange]) {
let mut guard = self
.inner
.lock()
.expect("native Kafka commit state poisoned");
for range in ranges {
let state = guard.entry(range.topic_partition()).or_default();
state.observe(range.first_offset, range.next_offset);
}
self.progress.notify_all();
}
pub(crate) fn process(&self, ranges: &[TopicPartitionOffsetRange]) -> KafkaClientResult<()> {
self.reject_lost(ranges)?;
let mut guard = self
.inner
.lock()
.expect("native Kafka commit state poisoned");
for range in ranges {
let state = guard.entry(range.topic_partition()).or_default();
state.observe(range.first_offset, range.next_offset);
state.process(range.first_offset, range.next_offset);
}
self.progress.notify_all();
Ok(())
}
pub(crate) fn outstanding(&self) -> u64 {
let guard = self
.inner
.lock()
.expect("native Kafka commit state poisoned");
outstanding_locked(&guard)
}
pub(crate) fn outstanding_by_partition(&self) -> BTreeMap<TopicPartition, u64> {
let guard = self
.inner
.lock()
.expect("native Kafka commit state poisoned");
guard
.iter()
.map(|(partition, state)| (partition.clone(), state.outstanding()))
.collect()
}
pub(crate) fn committed_offsets_sum(&self) -> u64 {
let guard = self
.inner
.lock()
.expect("native Kafka commit state poisoned");
guard
.values()
.map(|state| state.committed_next.unwrap_or(0).max(0) as u64)
.sum()
}
pub(crate) fn uncommitted(&self) -> u64 {
let guard = self
.inner
.lock()
.expect("native Kafka commit state poisoned");
uncommitted_locked(&guard)
}
pub(crate) fn due_commits(&self) -> Vec<OffsetCommit> {
let guard = self
.inner
.lock()
.expect("native Kafka commit state poisoned");
due_commits_locked(&guard, None)
}
pub(crate) fn due_commits_for(
&self,
partitions: &BTreeSet<TopicPartition>,
) -> Vec<OffsetCommit> {
let guard = self
.inner
.lock()
.expect("native Kafka commit state poisoned");
due_commits_locked(&guard, Some(partitions))
}
pub(crate) fn mark_committed(&self, commits: &[OffsetCommit]) {
let mut guard = self
.inner
.lock()
.expect("native Kafka commit state poisoned");
for commit in commits {
if let Some(state) = guard.get_mut(&commit.topic_partition()) {
state.committed_next = Some(
state
.committed_next
.map_or(commit.offset, |current| current.max(commit.offset)),
);
}
}
self.progress.notify_all();
}
pub(crate) fn wait_for_all_committed(&self, timeout: std::time::Duration) -> bool {
let deadline = std::time::Instant::now() + timeout;
let mut guard = self
.inner
.lock()
.expect("native Kafka commit state poisoned");
while outstanding_locked(&guard) != 0 || uncommitted_locked(&guard) != 0 {
let now = std::time::Instant::now();
if now >= deadline {
return false;
}
let remaining = deadline.saturating_duration_since(now);
let (next_guard, result) = self
.progress
.wait_timeout(guard, remaining)
.expect("native Kafka commit condvar poisoned");
guard = next_guard;
if result.timed_out()
&& (outstanding_locked(&guard) != 0 || uncommitted_locked(&guard) != 0)
{
return false;
}
}
true
}
pub(crate) fn mark_lost(&self, partitions: &BTreeSet<TopicPartition>) {
if partitions.is_empty() {
return;
}
let mut guard = self.lost.lock().expect("native Kafka lost set poisoned");
guard.extend(partitions.iter().cloned());
self.progress.notify_all();
}
fn reject_lost(&self, ranges: &[TopicPartitionOffsetRange]) -> KafkaClientResult<()> {
let guard = self.lost.lock().expect("native Kafka lost set poisoned");
if let Some(range) = ranges
.iter()
.find(|range| guard.contains(&range.topic_partition()))
{
return Err(KafkaClientError::AssignmentLost {
topic: range.topic.clone(),
partition: range.partition,
});
}
Ok(())
}
}
fn outstanding_locked(guard: &BTreeMap<TopicPartition, PartitionCommitState>) -> u64 {
guard.values().map(PartitionCommitState::outstanding).sum()
}
fn uncommitted_locked(guard: &BTreeMap<TopicPartition, PartitionCommitState>) -> u64 {
guard.values().map(PartitionCommitState::uncommitted).sum()
}
fn due_commits_locked(
guard: &BTreeMap<TopicPartition, PartitionCommitState>,
partitions: Option<&BTreeSet<TopicPartition>>,
) -> Vec<OffsetCommit> {
guard
.iter()
.filter(|(partition, _)| {
partitions.is_none_or(|partitions| partitions.contains(*partition))
})
.filter_map(|(partition, state)| {
let processed = state.processed_next?;
if state
.committed_next
.is_some_and(|committed| committed >= processed)
{
return None;
}
Some(OffsetCommit {
topic: partition.topic.clone(),
partition: partition.partition,
offset: processed,
})
})
.collect()
}
impl PartitionCommitState {
fn observe(&mut self, first_offset: i64, next_offset: i64) {
self.emitted_next = Some(
self.emitted_next
.map_or(next_offset, |current| current.max(next_offset)),
);
self.processed_next
.get_or_insert(first_offset.min(next_offset));
self.committed_next
.get_or_insert(first_offset.min(next_offset));
}
fn process(&mut self, first_offset: i64, next_offset: i64) {
if first_offset >= next_offset {
return;
}
self.processed_ranges
.entry(first_offset)
.and_modify(|end| *end = (*end).max(next_offset))
.or_insert(next_offset);
self.advance_processed();
}
fn advance_processed(&mut self) {
let Some(mut cursor) = self.processed_next else {
return;
};
while let Some((&start, &end)) = self.processed_ranges.range(..=cursor).next_back() {
if end <= cursor {
self.processed_ranges.remove(&start);
continue;
}
if start > cursor {
break;
}
cursor = end;
self.processed_next = Some(cursor);
self.processed_ranges.remove(&start);
}
}
fn outstanding(&self) -> u64 {
let emitted = self.emitted_next.unwrap_or(0);
let processed = self.processed_next.unwrap_or(emitted);
emitted.saturating_sub(processed) as u64
}
fn uncommitted(&self) -> u64 {
let processed = self.processed_next.unwrap_or(0);
let committed = self.committed_next.unwrap_or(processed);
processed.saturating_sub(committed) as u64
}
}
#[cfg(test)]
mod tests {
use super::*;
fn range(first_offset: i64, next_offset: i64) -> TopicPartitionOffsetRange {
TopicPartitionOffsetRange {
topic: "topic".to_owned(),
partition: 0,
first_offset,
next_offset,
}
}
#[test]
fn local_commit_state_does_not_advance_past_out_of_order_gap() {
let state = LocalCommitState::default();
state.observe(&[range(0, 10)]);
state.observe(&[range(10, 20)]);
state.process(&[range(10, 20)]).expect("process second");
assert_eq!(state.outstanding(), 20);
assert!(state.due_commits().is_empty());
state.process(&[range(0, 10)]).expect("process first");
let commits = state.due_commits();
assert_eq!(commits.len(), 1);
assert_eq!(commits[0].offset, 20);
state.mark_committed(&commits);
assert_eq!(state.outstanding(), 0);
assert_eq!(state.uncommitted(), 0);
assert_eq!(state.committed_offsets_sum(), 20);
}
#[test]
fn local_commit_state_rejects_lost_assignment_commit() {
let state = LocalCommitState::default();
state.observe(&[range(0, 10)]);
let lost = [TopicPartition::new("topic", 0)]
.into_iter()
.collect::<BTreeSet<_>>();
state.mark_lost(&lost);
let error = state.process(&[range(0, 10)]).expect_err("lost commit");
assert!(matches!(
error,
KafkaClientError::AssignmentLost {
topic,
partition: 0
} if topic == "topic"
));
}
}