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>
impl<T: Send + 'static> KafkaFuture<T>
Sourcepub fn completed(result: Result<T, Error>) -> Self
pub fn completed(result: Result<T, Error>) -> Self
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.
Sourcepub async fn get(&self) -> Result<T, Error>
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.
Sourcepub async fn get_with_timeout(&self, timeout: Duration) -> Result<T, Error>
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.
Sourcepub fn is_done(&self) -> bool
pub fn is_done(&self) -> bool
Whether this future is complete.
This is the Rust equivalent of Java’s Future.isDone().
Sourcepub fn all_of(futures: Vec<KafkaFuture<T>>) -> KafkaFuture<()>
pub fn all_of(futures: Vec<KafkaFuture<T>>) -> KafkaFuture<()>
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).
Sourcepub fn then_apply<R, F>(&self, function: F) -> KafkaFuture<R>
pub fn then_apply<R, F>(&self, function: F) -> KafkaFuture<R>
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.
Sourcepub fn join_map<K>(
entries: Vec<(K, KafkaFuture<T>)>,
) -> KafkaFuture<HashMap<K, T>>
pub fn join_map<K>( entries: Vec<(K, KafkaFuture<T>)>, ) -> KafkaFuture<HashMap<K, T>>
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.
Sourcepub fn join_map_results<K>(
entries: Vec<(K, KafkaFuture<T>)>,
) -> KafkaFuture<HashMap<K, Result<T, Error>>>
pub fn join_map_results<K>( entries: Vec<(K, KafkaFuture<T>)>, ) -> KafkaFuture<HashMap<K, Result<T, Error>>>
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.
Sourcepub fn then_apply_try<R, F>(&self, function: F) -> KafkaFuture<R>
pub fn then_apply_try<R, F>(&self, function: F) -> KafkaFuture<R>
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.