Skip to main content

KafkaFuture

Struct KafkaFuture 

Source
pub struct KafkaFuture<T: Send + 'static> { /* private fields */ }
Expand description

A flexible future which supports async result retrieval.

Translated from org.apache.kafka.common.KafkaFuture.

In Java, KafkaFuture<T> implements Future<T> and provides methods like get(), get(timeout, unit), and isDone(). In Rust, this struct wraps an internal implementation and provides equivalent async methods.

This type is used as the return type for Producer.send(), hiding the internal FutureRecordMetadata behind a public interface — matching how Java’s Producer.send() returns Future<RecordMetadata> rather than the internal implementation type.

Implementations§

Source§

impl<T: Send + 'static> KafkaFuture<T>

Source

pub fn completed(result: Result<T, Error>) -> Self
where T: Clone + Sync,

Create a KafkaFuture that is already resolved with the given result.

Useful when the result is known at the time the future is constructed — for example when implementing the Producer trait by awaiting a remote call and wrapping the response in a future to satisfy the trait’s return type. Analogous to std::future::ready.

Both get and get_with_timeout resolve immediately with a clone of the result; is_done returns true.

Source

pub async fn get(&self) -> Result<T, Error>

Await the result of this future.

This is the Rust equivalent of Java’s Future.get().

§Errors

Returns the error from the underlying operation if it failed.

Source

pub async fn get_with_timeout(&self, timeout: Duration) -> Result<T, Error>

Await the result of this future with a timeout.

This is the Rust equivalent of Java’s Future.get(timeout, unit).

§Arguments
  • timeout - The maximum time to wait
§Errors

Returns Error::LocalTimeout if the timeout elapses before the result is available — Java’s java.util.concurrent.TimeoutException, which Future.get(timeout, unit) declares, not the retriable org.apache.kafka.common.errors.TimeoutException. Returns the error from the underlying operation if it failed.

Source

pub fn is_done(&self) -> bool

Whether this future is complete.

This is the Rust equivalent of Java’s Future.isDone().

Source

pub fn all_of(futures: Vec<KafkaFuture<T>>) -> KafkaFuture<()>
where T: Clone + Sync,

Returns a new KafkaFuture that is completed when all the given futures have completed. If any future completes exceptionally, the returned future returns that error. If multiple futures fail, which error gets returned is arbitrarily chosen (the first encountered while awaiting in order).

Translated from org.apache.kafka.common.KafkaFuture.allOf. Unlike Java’s variadic allOf(KafkaFuture<?>...), this Rust version is homogeneous (Vec<KafkaFuture<T>>), which is all the admin *Result types require (they combine per-key futures of a single type).

Source

pub fn then_apply<R, F>(&self, function: F) -> KafkaFuture<R>
where T: Clone + Sync, R: Clone + Send + Sync + 'static, F: Fn(T) -> R + Send + Sync + 'static,

Returns a new KafkaFuture that, when this future completes normally, is completed with the result of applying function to this future’s value. If this future completes exceptionally, the returned future completes with the same exception.

Translated from org.apache.kafka.common.KafkaFuture.thenApply, for the common case where the transform cannot fail.

Source

pub fn join_map<K>( entries: Vec<(K, KafkaFuture<T>)>, ) -> KafkaFuture<HashMap<K, T>>
where T: Clone + Sync, K: Hash + Eq + Clone + Send + Sync + 'static,

Returns a future that completes when all the given keyed futures complete, yielding a map from each key to its resolved value. If any future completes exceptionally, the returned future yields that error.

This is the combinator behind the admin *Result aggregators (DescribeTopicsResult::all_topic_names, etc.). Java expresses the same thing inline as KafkaFuture.allOf(...).thenApply(v -> collect each future.get()); the Rust port names it because the get-driven model cannot call .get() synchronously inside a then_apply closure.

Source

pub fn join_map_results<K>( entries: Vec<(K, KafkaFuture<T>)>, ) -> KafkaFuture<HashMap<K, Result<T, Error>>>
where T: Clone + Sync, K: Hash + Eq + Clone + Send + Sync + 'static,

Returns a future that completes when all the given keyed futures complete, yielding a map from each key to that future’s outcome.

This is the collect-all counterpart of join_map: a failed input future does not short-circuit, so every key’s Result is preserved. Java has no direct equivalent because callers there hold the per-key KafkaFutures themselves and inspect each one; the C FFI has no KafkaFuture type, so the admin bindings flatten a multi-key *Result into one handle carrying a value and an error per key (see PLAN-bindings.md §2 / D2 and admin-client.md §5, which requires per-key granularity to survive the boundary). join_map cannot be reused for that: its future.get().await? abandons the remaining keys on the first error.

Source

pub fn then_apply_try<R, F>(&self, function: F) -> KafkaFuture<R>
where T: Clone + Sync, R: Clone + Send + Sync + 'static, F: Fn(T) -> Result<R, Error> + Send + Sync + 'static,

Like then_apply but the transform may fail. If function returns Err, the returned future completes with that error.

This models the Java thenApply cases whose BaseFunction throws — for example CreateTopicsResult.TopicMetadataAndConfig accessors that call ensureSuccess() and rethrow a stored exception.

Trait Implementations§

Source§

impl<T: Send + 'static> Clone for KafkaFuture<T>

Source§

fn clone(&self) -> Self

Returns a duplicate of the value. Read more
1.0.0 · Source§

fn clone_from(&mut self, source: &Self)

Performs copy-assignment from source. Read more
Source§

impl<T: Send + 'static> Debug for KafkaFuture<T>

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more

Auto Trait Implementations§

§

impl<T> Freeze for KafkaFuture<T>

§

impl<T> !RefUnwindSafe for KafkaFuture<T>

§

impl<T> Send for KafkaFuture<T>

§

impl<T> Sync for KafkaFuture<T>

§

impl<T> Unpin for KafkaFuture<T>

§

impl<T> UnsafeUnpin for KafkaFuture<T>

§

impl<T> !UnwindSafe for KafkaFuture<T>

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> CloneToUninit for T
where T: Clone,

Source§

unsafe fn clone_to_uninit(&self, dest: *mut u8)

🔬This is a nightly-only experimental API. (clone_to_uninit)
Performs copy-assignment from self to dest. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> ToOwned for T
where T: Clone,

Source§

type Owned = T

The resulting type after obtaining ownership.
Source§

fn to_owned(&self) -> T

Creates owned data from borrowed data, usually by cloning. Read more
Source§

fn clone_into(&self, target: &mut T)

Uses borrowed data to replace owned data, usually by cloning. Read more
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = Infallible

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

§

fn vzip(self) -> V