mod config;
mod group_coordinator;
mod handles;
mod reactor;
mod types;
mod util;
pub use config::{AutoOffsetReset, ConsumerConfig, PartitionAssignmentStrategy};
pub use handles::{GroupHandle, OffsetHandle};
pub use types::ConsumerRecord;
use std::collections::VecDeque;
use std::sync::Arc;
use std::time::Duration;
use tokio::sync::{mpsc, oneshot};
use crate::cluster::ClusterClient;
use crate::consumer::reactor::spawn_consumer_task;
use crate::error::Result;
pub struct ConsumerStream {
rx: mpsc::Receiver<Vec<ConsumerRecord>>,
cmd_tx: mpsc::UnboundedSender<types::ConsumerCommand>,
buffer: VecDeque<ConsumerRecord>,
}
impl ConsumerStream {
pub async fn recv(&mut self) -> Option<ConsumerRecord> {
if let Some(record) = self.buffer.pop_front() {
return Some(record);
}
let batch = self.rx.recv().await?;
self.buffer.extend(batch);
self.buffer.pop_front()
}
}
impl Drop for ConsumerStream {
fn drop(&mut self) {
let _ = self.cmd_tx.send(types::ConsumerCommand::Shutdown);
}
}
pub struct Consumer {
cmd_tx: mpsc::UnboundedSender<types::ConsumerCommand>,
record_rx: mpsc::Receiver<Vec<ConsumerRecord>>,
offset_handle: handles::OffsetHandle,
group_handle: handles::GroupHandle,
is_group: bool,
started: bool,
}
impl Consumer {
pub(crate) fn new(cluster: Arc<ClusterClient>, config: ConsumerConfig) -> Self {
let (cmd_tx, cmd_rx) = mpsc::unbounded_channel();
let (record_tx, record_rx) = mpsc::channel(64);
let is_group = config.group_id.is_some();
spawn_consumer_task(cluster, config, cmd_rx, record_tx);
let cmd_tx_clone = cmd_tx.clone();
Self {
cmd_tx,
record_rx,
offset_handle: handles::OffsetHandle {
cmd_tx: cmd_tx_clone.clone(),
},
group_handle: handles::GroupHandle {
cmd_tx: cmd_tx_clone,
},
is_group,
started: false,
}
}
pub async fn subscribe(&mut self, topics: Vec<String>) -> Result<()> {
let (tx, rx) = oneshot::channel();
self.cmd_tx
.send(types::ConsumerCommand::Subscribe {
topics,
reply: Some(tx),
})
.map_err(|_| crate::error::KafkaError::ConnectionClosed)?;
rx.await
.map_err(|_| crate::error::KafkaError::ConnectionClosed)?
}
pub async fn assign(&mut self, topic: impl Into<String>, partitions: Vec<i32>) -> Result<()> {
let (tx, rx) = oneshot::channel();
self.cmd_tx
.send(types::ConsumerCommand::Assign {
topic: topic.into(),
partitions,
reply: Some(tx),
})
.map_err(|_| crate::error::KafkaError::ConnectionClosed)?;
rx.await
.map_err(|_| crate::error::KafkaError::ConnectionClosed)?
}
pub async fn seek(&self, topic: impl Into<String>, partition: i32, offset: i64) -> Result<()> {
self.cmd_tx
.send(types::ConsumerCommand::Seek {
topic: topic.into(),
partition,
offset,
})
.map_err(|_| crate::error::KafkaError::ConnectionClosed)?;
Ok(())
}
pub async fn unsubscribe(&mut self) -> Result<()> {
let (tx, rx) = oneshot::channel();
self.cmd_tx
.send(types::ConsumerCommand::Unsubscribe { reply: tx })
.map_err(|_| crate::error::KafkaError::ConnectionClosed)?;
rx.await
.map_err(|_| crate::error::KafkaError::ConnectionClosed)?
}
pub async fn poll(&mut self) -> Result<Vec<ConsumerRecord>> {
if !self.started {
self.cmd_tx
.send(types::ConsumerCommand::StartPolling)
.map_err(|_| crate::error::KafkaError::ConnectionClosed)?;
self.started = true;
}
let batch = self
.record_rx
.recv()
.await
.ok_or(crate::error::KafkaError::ConnectionClosed)?;
let mut records = batch;
while let Ok(batch) = self.record_rx.try_recv() {
records.extend(batch);
}
Ok(records)
}
pub async fn poll_timeout(&mut self, timeout: Duration) -> Result<Vec<ConsumerRecord>> {
if !self.started {
self.cmd_tx
.send(types::ConsumerCommand::StartPolling)
.map_err(|_| crate::error::KafkaError::ConnectionClosed)?;
self.started = true;
}
let mut records = Vec::new();
while let Ok(batch) = self.record_rx.try_recv() {
records.extend(batch);
}
if !records.is_empty() {
return Ok(records);
}
match tokio::time::timeout(timeout, self.record_rx.recv()).await {
Ok(Some(batch)) => {
records = batch;
while let Ok(batch) = self.record_rx.try_recv() {
records.extend(batch);
}
Ok(records)
}
Ok(None) => Err(crate::error::KafkaError::ConnectionClosed),
Err(_) => Ok(Vec::new()),
}
}
pub async fn try_poll(&mut self) -> Result<Vec<ConsumerRecord>> {
if !self.started {
self.cmd_tx
.send(types::ConsumerCommand::StartPolling)
.map_err(|_| crate::error::KafkaError::ConnectionClosed)?;
self.started = true;
}
let mut records = Vec::new();
while let Ok(batch) = self.record_rx.try_recv() {
records.extend(batch);
}
Ok(records)
}
pub async fn close(&self) -> Result<()> {
let _ = self.cmd_tx.send(types::ConsumerCommand::Shutdown);
Ok(())
}
pub fn offsets(&self) -> &OffsetHandle {
&self.offset_handle
}
pub fn group(&self) -> &GroupHandle {
&self.group_handle
}
pub fn is_group(&self) -> bool {
self.is_group
}
pub fn into_stream(mut self) -> ConsumerStream {
if !self.started {
let _ = self.cmd_tx.send(types::ConsumerCommand::StartPolling);
self.started = true;
}
ConsumerStream {
rx: self.record_rx,
cmd_tx: self.cmd_tx.clone(),
buffer: VecDeque::new(),
}
}
}