Skip to main content

ConsumerInterceptor

Trait ConsumerInterceptor 

Source
pub trait ConsumerInterceptor:
    Send
    + Sync
    + Debug {
    // Required method
    fn before_consume(&self, msg: &mut IncomingMessage);

    // Provided methods
    fn on_acknowledge(
        &self,
        _message_id: MessageId,
        _outcome: Result<(), &PulsarError>,
    ) { ... }
    fn on_acknowledge_cumulative(
        &self,
        _message_id: MessageId,
        _outcome: Result<(), &PulsarError>,
    ) { ... }
    fn on_negative_acks_send(&self, _message_ids: &[MessageId]) { ... }
}
Expand description

Java ConsumerInterceptor SPI. Plug receive-side hooks behind Consumer::receive to inspect / mutate incoming messages and observe ack outcomes. Mirrors org.apache.pulsar.client.api.interceptor.ConsumerInterceptor:

  • before_consume runs on every received message and may mutate it.
  • on_acknowledge fires on every individual / batch ack.
  • on_acknowledge_cumulative fires on every cumulative ack.
  • on_negative_acks_send fires when the runtime forwards a redeliver-unacknowledged command (negative ack with delay or immediate).

Each callback runs on the receive / ack path — keep them fast and non-blocking. Use receive_with_interceptors to chain a list against a magnetar_runtime_tokio::Consumer.

Required Methods§

Source

fn before_consume(&self, msg: &mut IncomingMessage)

Inspect and optionally mutate the incoming message before it is handed back to the user. Mirrors Java ConsumerInterceptor#beforeConsume.

Provided Methods§

Source

fn on_acknowledge( &self, _message_id: MessageId, _outcome: Result<(), &PulsarError>, )

Fired after an individual or batch ack completes (success or error). Mirrors Java ConsumerInterceptor#onAcknowledge.

Source

fn on_acknowledge_cumulative( &self, _message_id: MessageId, _outcome: Result<(), &PulsarError>, )

Fired after a cumulative ack completes. Mirrors Java ConsumerInterceptor#onAcknowledgeCumulative.

Source

fn on_negative_acks_send(&self, _message_ids: &[MessageId])

Fired when the runtime forwards a CommandRedeliverUnacknowledgedMessages for one or more message ids. Mirrors Java ConsumerInterceptor#onNegativeAcksSend.

Dyn Compatibility§

This trait is dyn compatible.

In older versions of Rust, dyn compatibility was called "object safety".

Implementors§