Skip to content

Commit f0fe169

Browse files
committed
Kafka: Tweaks
1 parent 63c0cac commit f0fe169

File tree

1 file changed

+3
-4
lines changed

1 file changed

+3
-4
lines changed

sea-streamer-kafka/src/producer.rs

+3-4
Original file line numberDiff line numberDiff line change
@@ -4,12 +4,11 @@ use crate::{
44
cluster::cluster_uri, impl_into_string, stream_err, BaseOptionKey, KafkaConnectOptions,
55
KafkaErr, KafkaResult, DEFAULT_TIMEOUT,
66
};
7-
pub use rdkafka::producer::FutureRecord;
87
use rdkafka::{
98
config::ClientConfig,
10-
producer::{DeliveryFuture, Producer as ProducerTrait},
9+
producer::{DeliveryFuture, FutureRecord as RawPayload, Producer as ProducerTrait},
1110
};
12-
pub use rdkafka::{consumer::ConsumerGroupMetadata, TopicPartitionList};
11+
pub use rdkafka::{consumer::ConsumerGroupMetadata, producer::FutureRecord, TopicPartitionList};
1312
use sea_streamer_runtime::spawn_blocking;
1413
use sea_streamer_types::{
1514
export::{async_trait, futures::FutureExt},
@@ -89,7 +88,7 @@ impl Producer for KafkaProducer {
8988
fn send_to<S: Buffer>(&self, stream: &StreamKey, payload: S) -> KafkaResult<Self::SendFuture> {
9089
let fut = self
9190
.get()
92-
.send_result(FutureRecord::<str, [u8]>::to(stream.name()).payload(payload.as_bytes()))
91+
.send_result(RawPayload::<str, [u8]>::to(stream.name()).payload(payload.as_bytes()))
9392
.map_err(|(err, _raw)| stream_err(err))?;
9493

9594
Ok(SendFuture {

0 commit comments

Comments
 (0)