Skip to main content

OffsetCommitCallback

Trait OffsetCommitCallback 

Source
pub trait OffsetCommitCallback:
    Send
    + Sync
    + 'static {
    // Required method
    fn on_complete<'life0, 'life1, 'life2, 'async_trait>(
        &'life0 self,
        offsets: &'life1 HashMap<TopicPartition, OffsetAndMetadata>,
        error: Option<&'life2 Error>,
    ) -> Pin<Box<dyn Future<Output = ()> + Send + 'async_trait>>
       where Self: 'async_trait,
             'life0: 'async_trait,
             'life1: 'async_trait,
             'life2: 'async_trait;
}
Expand description

A callback interface that the user can implement to trigger custom actions when a commit request completes.

Corresponds to Java’s org.apache.kafka.clients.consumer.OffsetCommitCallback.

§Invocation model

Per consumer-threading.md §31, the callback executes on the caller’s task that is currently inside poll() / commit_*() / close() — never on the background task. Matches Java’s contract: “the callback may be executed in any thread calling poll()”.

§Bounds: Send + Sync + 'static

The callback travels through events on the background task and is invoked on the app task after the background task has already moved on. Arc<dyn> storage is required for the symmetric “send through channel, invoke later on app side” pipeline used by process_background_events. Box<dyn> would forbid the clone and the design collapses.

Required Methods§

Source

fn on_complete<'life0, 'life1, 'life2, 'async_trait>( &'life0 self, offsets: &'life1 HashMap<TopicPartition, OffsetAndMetadata>, error: Option<&'life2 Error>, ) -> Pin<Box<dyn Future<Output = ()> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait,

A callback method the user can implement to provide asynchronous handling of commit request completion. This method will be called when the commit request sent to the server has been acknowledged.

Corresponds to Java’s void onComplete(Map<TopicPartition, OffsetAndMetadata> offsets, Exception exception).

Returns () (not Result) because Java’s onComplete is void and does not allow the callback to fail the commit — the commit has already happened. The error parameter is informational and matches Java’s pattern of “exception == null means success”.

§Arguments
  • offsets - the offsets and associated metadata that this callback applies to
  • error - Some(&error) if the commit failed, None if it completed successfully

Implementors§