Interface KafkaMessageResponse
public interface KafkaMessageResponse
Represents a response that can manipulate a flow of messages that is mapped from a Kafka native FetchResponse.
- Author:
- Jeoffrey HAEYAERT (jeoffrey.haeyaert at graviteesource.com), GraviteeSource Team
-
Method Summary
Modifier and TypeMethodDescriptionio.reactivex.rxjava3.core.Flowable<KafkaMessage>messages()Get the flow of messages.voidmessages(io.reactivex.rxjava3.core.Flowable<KafkaMessage> messages) Set the request message flow.default io.reactivex.rxjava3.core.CompletableonMessage(Function<KafkaMessage, io.reactivex.rxjava3.core.Maybe<KafkaMessage>> onMessage) Applies a given transformation on each message.io.reactivex.rxjava3.core.CompletableonMessages(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
FlowableofKafkaMessage.
-
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)oronMessage(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
Completablethat 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
Completablethat completes once the message transformation has been set up on the message flow (not executed).
-