Expand description
Kafka sources and sinks for Datum streams.
datum-mq is an additive Datum satellite crate with Tokio-native Kafka
consumer and producer backends. The rustls TLS path uses ring (Rust plus
assembly), with no AWS-LC C library or CMake build. Native produce and fetch
support uncompressed, gzip, Snappy, LZ4, and Zstd record batches. SASL
supports PLAIN, SCRAM-SHA-256, and SCRAM-SHA-512; GSSAPI and OAUTHBEARER are
not supported.
The consumer’s shipped offset-commit guarantee is at-least-once.
KafkaSource::committable emits a datum::SourceWithContext whose
context is a KafkaOffset; downstream code commits that offset only after
its own checkpoint/side effect succeeds.
KafkaSink uses the native producer by default, and its materialized
KafkaProducerControl drains final acknowledgements before shutdown. The
native producer enables idempotence by default: it negotiates a producer
id/epoch and uses per-partition sequences so the broker deduplicates a
retried batch. Set enable.idempotence=false to opt out explicitly. The
pipeline remains bounded and partition ordered.
Kafka transactions and consume-transform-produce EOS are not supported.
Structs§
- Consumer
Record - Owned Kafka consumer record.
- Kafka
Batch Offset - Committable watermark for a batch of Kafka records.
- Kafka
Config - Raw Kafka property map with typed helpers for common Datum defaults.
- Kafka
Consumer Settings - Consumer settings for
KafkaSource. - Kafka
Control - Materialized control for a Kafka source.
- Kafka
Header - Kafka record header.
- Kafka
Metrics - Cheap cloneable metrics handle for Kafka sources/sinks.
- Kafka
Metrics Snapshot - Point-in-time metrics snapshot.
- Kafka
Offset - Context value emitted by
KafkaSource::committable. - Kafka
Offset Batch - Batch of offsets for external checkpointing or batched downstream commit sinks.
- Kafka
Payload Batch - A payload batch plus one committable Kafka watermark per touched partition.
- Kafka
Payload Record - Payload-only Kafka record emitted by
KafkaSource::committable_payload_batches. - Kafka
Producer Control - Materialized control for a Kafka producer sink.
- Kafka
Producer Settings - Producer settings for
KafkaSink. - Kafka
Sink - Kafka producer sink entry points.
- Kafka
Source - Kafka source entry points.
- Offset
Commit - A Kafka commit position.
offsetis the next offset to consume. - Producer
Record - Owned Kafka producer record.
- Topic
Partition - Kafka topic/partition identifier.
- Topic
Partition Offset - A concrete topic/partition plus the starting offset for assignment subscriptions.
Enums§
- Commit
Policy - Offset commit mode for
KafkaSource. - Kafka
Consumer Backend - Consumer implementation used by
KafkaSource. - Kafka
Producer Backend - Producer implementation used by
crate::KafkaSink. - Kafka
Timestamp - Kafka record timestamp.
- MqError
- Kafka connector errors.
- Offset
- Kafka assignment start offset.
- Subscription
- Topic subscription mode.
Constants§
- VERSION
- The
datum-mqcrate version.
Type Aliases§
- MqResult
- Result type used by
datum-mq.