#![forbid(unsafe_code)]
use datum_mq::{
CommitPolicy, KafkaConsumerSettings, KafkaProducerSettings, KafkaSink, KafkaSource,
Subscription,
};
#[test]
fn kafka_committable_source_blueprint() {
use datum::SourceWithContext;
use datum_mq::{ConsumerRecord, KafkaControl, KafkaOffset};
let settings = KafkaConsumerSettings::new("127.0.0.1:9092", "orders-consumer")
.with("auto.offset.reset", "earliest")
.with_commit_policy(CommitPolicy::Manual)
.with_backpressure(2_048, 4_096);
let source: SourceWithContext<ConsumerRecord, KafkaOffset, KafkaControl> =
KafkaSource::committable(settings, Subscription::topics(["orders"]));
let _ = source;
}
#[test]
fn kafka_sink_blueprint() {
use datum::Sink;
use datum_mq::{KafkaProducerControl, ProducerRecord};
let sink: Sink<ProducerRecord, KafkaProducerControl> =
KafkaSink::plain(KafkaProducerSettings::new("127.0.0.1:9092"));
let _ = sink;
}
#[test]
fn kafka_native_producer_backend_blueprint() {
use datum::Sink;
use datum_mq::{
KafkaProducerBackend, KafkaProducerControl, KafkaProducerSettings, KafkaSink,
ProducerRecord,
};
let settings = KafkaProducerSettings::new("127.0.0.1:9092")
.with_producer_backend(KafkaProducerBackend::Native);
let sink: Sink<ProducerRecord, KafkaProducerControl> = KafkaSink::plain(settings);
let _ = sink;
}
#[test]
fn kafka_native_consumer_backend_blueprint() {
use datum_mq::{KafkaConsumerBackend, KafkaConsumerSettings, KafkaSource, Subscription};
let settings = KafkaConsumerSettings::new("127.0.0.1:9092", "orders-consumer")
.with_consumer_backend(KafkaConsumerBackend::Native)
.with("auto.offset.reset", "earliest");
let source =
KafkaSource::committable_payload_batches(settings, Subscription::topics(["orders"]));
let _ = source;
}
#[test]
fn cookbook_native_committable_at_least_once() {
use datum::{Sink, StreamError};
use datum_mq::{
CommitPolicy, ConsumerRecord, KafkaConsumerBackend, KafkaConsumerSettings, KafkaOffset,
KafkaSource, Subscription,
};
let Some(topic) = std::env::var("DATUM_DOCS_ORDERS_TOPIC").ok() else {
return;
};
let bootstrap =
std::env::var("MQ_BOOTSTRAP_SERVERS").unwrap_or_else(|_| "127.0.0.1:9092".to_owned());
let group = format!("datum-docs-native-{}", std::process::id());
let settings = KafkaConsumerSettings::new(bootstrap, group)
.with_consumer_backend(KafkaConsumerBackend::Native)
.with_commit_policy(CommitPolicy::Manual)
.with("auto.offset.reset", "earliest")
.with("partition.assignment.strategy", "cooperative-sticky")
.with_backpressure(512, 1_024)
.with_poll_batch_size(64);
let processed = KafkaSource::committable(settings, Subscription::topics([topic]))
.as_source()
.take(3)
.run_with(Sink::fold_result(
0_u64,
|count, (record, offset): (ConsumerRecord, KafkaOffset)| {
let _dedupe_key = (&record.topic, record.partition, record.offset);
offset.commit().map_err(StreamError::from)?;
Ok(count + 1)
},
))
.expect("native Kafka consumer materializes")
.wait()
.expect("native Kafka consumer completes");
println!("processed {processed} orders");
assert_eq!(processed, 3);
}
#[test]
fn cookbook_native_consumer_group_drain() {
use std::time::Duration;
use datum::{Keep, Sink, StreamCompletion, StreamError};
use datum_mq::{
KafkaConsumerBackend, KafkaConsumerSettings, KafkaControl, KafkaPayloadBatch, KafkaSource,
MqError, Subscription,
};
let Some(topic) = std::env::var("DATUM_DOCS_GROUP_TOPIC").ok() else {
return;
};
let bootstrap =
std::env::var("MQ_BOOTSTRAP_SERVERS").unwrap_or_else(|_| "127.0.0.1:9092".to_owned());
let settings = KafkaConsumerSettings::new(bootstrap, "orders-workers")
.with_consumer_backend(KafkaConsumerBackend::Native)
.with("auto.offset.reset", "earliest")
.with("partition.assignment.strategy", "cooperative-sticky")
.with_backpressure(256, 512)
.with_commit_batch_size(128)
.with_commit_interval(Duration::from_millis(100));
let (control, completion): (KafkaControl, StreamCompletion<datum::NotUsed>) =
KafkaSource::committable_payload_batches(settings, Subscription::topics([topic]))
.to_mat(
Sink::foreach_result(|batch: KafkaPayloadBatch| {
for record in batch.records() {
let _payload = batch.payload(record);
}
match batch.commit() {
Ok(()) => Ok(()),
Err(error @ MqError::AssignmentLost { .. }) => {
Err(StreamError::from(error))
}
Err(error) => Err(StreamError::from(error)),
}
}),
Keep::both,
)
.run()
.expect("native Kafka group consumer materializes");
control
.drain_and_shutdown(Duration::from_secs(30))
.expect("consumer drains outstanding commits");
completion.wait().expect("consumer completes after drain");
let metrics = control.metrics().snapshot();
println!(
"drained; committed={} rebalances={} lost={}",
metrics.committed_offsets, metrics.rebalances, metrics.lost_partitions
);
}