Interface KafkaMessageRequest


public interface KafkaMessageRequest
Represents a request that can manipulate a flow of messages that is mapped from a Kafka native ProduceRequest.
Author:
Jeoffrey HAEYAERT (jeoffrey.haeyaert at graviteesource.com), GraviteeSource Team
  • Method Summary

    Modifier and Type
    Method
    Description
    io.reactivex.rxjava3.core.Flowable<KafkaMessage>
    Get the flow of messages.
    void
    messages(io.reactivex.rxjava3.core.Flowable<KafkaMessage> messages)
    Set the request message flow.
    default io.reactivex.rxjava3.core.Completable
    onMessage(Function<KafkaMessage,io.reactivex.rxjava3.core.Maybe<KafkaMessage>> onMessage)
    Applies a given transformation on each message.
    io.reactivex.rxjava3.core.Completable
    onMessages(io.reactivex.rxjava3.core.FlowableTransformer<KafkaMessage,KafkaMessage> onMessages)
    Applies a given transformation on each message.
  • Method Details

    • messages

      io.reactivex.rxjava3.core.Flowable<KafkaMessage> messages()
      Get the flow of messages. Messages are extracted from the ProduceRequest. WARN: you should not keep a direct reference on the message flow as it could be overridden by others at anytime.
      Returns:
      a Flowable of KafkaMessage.
    • messages

      void messages(io.reactivex.rxjava3.core.Flowable<KafkaMessage> messages)
      Set the request message flow. WARN:
      • Replacing the message flow DOES NOT take care of the previous message flow in place.
      • You MUST ensure to consume the previous message flow when using it.
      • You SHOULD consider using onMessages(FlowableTransformer) or onMessage(Function) that may be more appropriate for message transformation.
      Parameters:
      messages - the flow of messages.
      See Also:
    • onMessages

      io.reactivex.rxjava3.core.Completable onMessages(io.reactivex.rxjava3.core.FlowableTransformer<KafkaMessage,KafkaMessage> onMessages)
      Applies a given transformation on each message. Messages are extracted from the ProduceRequest. Ex: request.onMessages(messages -> messages.flatMap(message -> transformMyMessage(message)));
      Parameters:
      onMessages - the transformer that will be applied on each message.
      Returns:
      a Completable that completes once the message transformation has been set up on the message flow (not executed).
    • onMessage

      default io.reactivex.rxjava3.core.Completable onMessage(Function<KafkaMessage,io.reactivex.rxjava3.core.Maybe<KafkaMessage>> onMessage)
      Applies a given transformation on each message. Messages are extracted from the ProduceRequest. Ex: Discard a message: request.onMessage(message -> Maybe.empty()); Update a message: request.onMessage(message -> transformMyMessage(message));
      Parameters:
      onMessage - the transformer that will be applied on each message.
      Returns:
      a Completable that completes once the message transformation has been set up on the message flow (not executed).