Package io.kroxylicious.proxy.filter
The protocol filter API.
An interface is provided for each kind of request and response in the Kafka protocol, e.g. ProduceRequestFilter.
Protocol Filter implementations inherit whichever of the per-RPC interfaces they need to intercept.
They can inherit multiple interfaces if necessary.
For filters which needs to intercept most or all of the protocol it is more convenient to inherit
RequestFilter and/or ResponseFilter.
Important facts about the Kafka protocol
Pipelining
The Kafka protocol supports pipelining (meaning a client can send multiple requests, before getting a response for any of them). Therefore when writing a filter implementation do not assume you won't see multiple requests before seeing any corresponding responses.
Ordering
A broker does not, in general, send responses in the same order as it receives requests. Therefore when writing a filter implementation do not assume ordering.
Local view
A client may obtain information from one broker in a cluster and use it to interact with other
brokers in the cluster (or the same broker, but on a different connection, and therefore a different
channel and filter chain). A classic example would
be a producer or consumer making a metadata connection and Metadata request to a broker and
then connecting to a partition leader to producer/consume records (Produce and Fetch requests).
So although your filter
implementation might intercept both Metadata and Produce request/response
(for example), those requests will not pass through the same instance of your filter
implementation. Therefore it is incorrect, in general, to assume your filter has a global view of
the communication between the client and broker.
Implementing Filters
Filter Results
Filter implementation must return a CompletionStage containing a
FilterResult object. It is the job of FilterResult to convey what
message is to be forwarded to the next filter in the chain (or client/broker if at the chain's beginning
or end). It is also used to carry instructions such as indicating that the connection must be closed,
or a message dropped.
If the filter returns a CompletionStage that is already completed normally, Kroxylicious will immediately perform the action described by the FilterResult.
If the CompletionStage completes exceptionally, the connection is closed. This also applies if the CompletionStage does not complete within a timeout (20000 milliseconds).
Deferring Forwards
The filter may return a CompletionStage that is not yet completed. When this happens, Kroxylicious will pause reading from the downstream (the Client writes will eventually block), and it begins to queue up in-flight requests/responses arriving at the filter. This is done so that message order is maintained. Once the CompletionStage completes, the action described by the FilterResult is performed, reading from the downstream resumes and any queued up requests/responses are processed.
IMPORTANT: The pausing of reads from the downstream is a relatively costly operation. To maintain optimal performance filter implementations should minimise the occasions on which an incomplete CompletionStage is returned.
Creating Filter Result objects
The FilterContext is the factory for the FilterResult objects.
There are two convenience methods that simply allow a filter to immediately forward a result:
FilterContext.forwardRequest(org.apache.kafka.common.message.RequestHeaderData, org.apache.kafka.common.protocol.ApiMessage), andFilterContext.forwardResponse(org.apache.kafka.common.message.ResponseHeaderData, org.apache.kafka.common.protocol.ApiMessage).
To access richer features, use the filter result builders:
Thread Safety
The Filter API provides the following thread-safety guarantees:
- There is a single thread associated with each connection and this association lasts for the lifetime of connection..
- Each filter instance is associated with exactly one connection.
- Construction of the filter instance and dispatch of the filter methods
onXxxRequestandonXxxResponsetakes place on that same thread. - Any computation stages chained to the
CompletionStagereturned byFilterContext.sendRequest(org.apache.kafka.common.message.RequestHeaderData, org.apache.kafka.common.protocol.ApiMessage)using the default execution methods (using methods without the suffix async) or default asynchronous execution (using methods with suffix async that employ the stage's default asynchronous execution facility) are guaranteed to be performed by that same thread. Computation stages chained using custom asynchronous execution (using methods with suffix async that take an Executor argument) do not get this guarantee.
Filter implementations are free to rely on these guarantees to safely maintain state within fields of the Filter without employing additional synchronization.
If a Filter needs to do some asynchronous work and mutate members of the Filter, they can access
the thread for the connection via FilterFactoryContext.filterDispatchExecutor(),
which is available in FilterFactory.createFilter(FilterFactoryContext, java.lang.Object)
when it is creating an instance of a Filter. Ensure that any work executes quickly as this is the IO thread for potentially
many connections.
-
InterfacesClassDescriptionA stateless filter for AddOffsetsToTxnRequests.A stateless filter for AddOffsetsToTxnResponses.A stateless filter for AddPartitionsToTxnRequests.A stateless filter for AddPartitionsToTxnResponses.A stateless filter for AddRaftVoterRequests.A stateless filter for AddRaftVoterResponses.A stateless filter for AllocateProducerIdsRequests.A stateless filter for AllocateProducerIdsResponses.A stateless filter for AlterClientQuotasRequests.A stateless filter for AlterClientQuotasResponses.A stateless filter for AlterConfigsRequests.A stateless filter for AlterConfigsResponses.A stateless filter for AlterPartitionReassignmentsRequests.A stateless filter for AlterPartitionReassignmentsResponses.A stateless filter for AlterPartitionRequests.A stateless filter for AlterPartitionResponses.A stateless filter for AlterReplicaLogDirsRequests.A stateless filter for AlterReplicaLogDirsResponses.A stateless filter for AlterShareGroupOffsetsRequests.A stateless filter for AlterShareGroupOffsetsResponses.A stateless filter for AlterUserScramCredentialsRequests.A stateless filter for AlterUserScramCredentialsResponses.A stateless filter for ApiVersionsRequests.A stateless filter for ApiVersionsResponses.A stateless filter for AssignReplicasToDirsRequests.A stateless filter for AssignReplicasToDirsResponses.A stateless filter for BeginQuorumEpochRequests.A stateless filter for BeginQuorumEpochResponses.A stateless filter for BrokerHeartbeatRequests.A stateless filter for BrokerHeartbeatResponses.A stateless filter for BrokerRegistrationRequests.A stateless filter for BrokerRegistrationResponses.A stateless filter for ConsumerGroupDescribeRequests.A stateless filter for ConsumerGroupDescribeResponses.A stateless filter for ConsumerGroupHeartbeatRequests.A stateless filter for ConsumerGroupHeartbeatResponses.A stateless filter for ControllerRegistrationRequests.A stateless filter for ControllerRegistrationResponses.A stateless filter for CreateAclsRequests.A stateless filter for CreateAclsResponses.A stateless filter for CreateDelegationTokenRequests.A stateless filter for CreateDelegationTokenResponses.A stateless filter for CreatePartitionsRequests.A stateless filter for CreatePartitionsResponses.A stateless filter for CreateTopicsRequests.A stateless filter for CreateTopicsResponses.A stateless filter for DeleteAclsRequests.A stateless filter for DeleteAclsResponses.A stateless filter for DeleteGroupsRequests.A stateless filter for DeleteGroupsResponses.A stateless filter for DeleteRecordsRequests.A stateless filter for DeleteRecordsResponses.A stateless filter for DeleteShareGroupOffsetsRequests.A stateless filter for DeleteShareGroupOffsetsResponses.A stateless filter for DeleteShareGroupStateRequests.A stateless filter for DeleteShareGroupStateResponses.A stateless filter for DeleteTopicsRequests.A stateless filter for DeleteTopicsResponses.A stateless filter for DescribeAclsRequests.A stateless filter for DescribeAclsResponses.A stateless filter for DescribeClientQuotasRequests.A stateless filter for DescribeClientQuotasResponses.A stateless filter for DescribeClusterRequests.A stateless filter for DescribeClusterResponses.A stateless filter for DescribeConfigsRequests.A stateless filter for DescribeConfigsResponses.A stateless filter for DescribeDelegationTokenRequests.A stateless filter for DescribeDelegationTokenResponses.A stateless filter for DescribeGroupsRequests.A stateless filter for DescribeGroupsResponses.A stateless filter for DescribeLogDirsRequests.A stateless filter for DescribeLogDirsResponses.A stateless filter for DescribeProducersRequests.A stateless filter for DescribeProducersResponses.A stateless filter for DescribeQuorumRequests.A stateless filter for DescribeQuorumResponses.A stateless filter for DescribeShareGroupOffsetsRequests.A stateless filter for DescribeShareGroupOffsetsResponses.A stateless filter for DescribeTopicPartitionsRequests.A stateless filter for DescribeTopicPartitionsResponses.A stateless filter for DescribeTransactionsRequests.A stateless filter for DescribeTransactionsResponses.A stateless filter for DescribeUserScramCredentialsRequests.A stateless filter for DescribeUserScramCredentialsResponses.A stateless filter for ElectLeadersRequests.A stateless filter for ElectLeadersResponses.A stateless filter for EndQuorumEpochRequests.A stateless filter for EndQuorumEpochResponses.A stateless filter for EndTxnRequests.A stateless filter for EndTxnResponses.A stateless filter for EnvelopeRequests.A stateless filter for EnvelopeResponses.A stateless filter for ExpireDelegationTokenRequests.A stateless filter for ExpireDelegationTokenResponses.A stateless filter for FetchRequests.A stateless filter for FetchResponses.A stateless filter for FetchSnapshotRequests.A stateless filter for FetchSnapshotResponses.Marker interface all Filter interfaces extend fromA context to allow filters to interact with other filters and the pipeline.An Executor backed by the Filter Dispatch Thread.FilterFactory<C,
I> A pluggable source ofFilterinstances.Construction context for Filters.The result of a filter request or response operation that encapsulates the request or response to be forwarded to the next filter in the chain.FilterResultBuilder<H extends org.apache.kafka.common.protocol.ApiMessage,R extends FilterResult> Fluent builder for filter results.A stateless filter for FindCoordinatorRequests.A stateless filter for FindCoordinatorResponses.A stateless filter for GetTelemetrySubscriptionsRequests.A stateless filter for GetTelemetrySubscriptionsResponses.A stateless filter for HeartbeatRequests.A stateless filter for HeartbeatResponses.A stateless filter for IncrementalAlterConfigsRequests.A stateless filter for IncrementalAlterConfigsResponses.A stateless filter for InitializeShareGroupStateRequests.A stateless filter for InitializeShareGroupStateResponses.A stateless filter for InitProducerIdRequests.A stateless filter for InitProducerIdResponses.A stateless filter for JoinGroupRequests.A stateless filter for JoinGroupResponses.A stateless filter for LeaveGroupRequests.A stateless filter for LeaveGroupResponses.A stateless filter for ListConfigResourcesRequests.A stateless filter for ListConfigResourcesResponses.A stateless filter for ListGroupsRequests.A stateless filter for ListGroupsResponses.A stateless filter for ListOffsetsRequests.A stateless filter for ListOffsetsResponses.A stateless filter for ListPartitionReassignmentsRequests.A stateless filter for ListPartitionReassignmentsResponses.A stateless filter for ListTransactionsRequests.A stateless filter for ListTransactionsResponses.A stateless filter for MetadataRequests.A stateless filter for MetadataResponses.A stateless filter for OffsetCommitRequests.A stateless filter for OffsetCommitResponses.A stateless filter for OffsetDeleteRequests.A stateless filter for OffsetDeleteResponses.A stateless filter for OffsetFetchRequests.A stateless filter for OffsetFetchResponses.A stateless filter for OffsetForLeaderEpochRequests.A stateless filter for OffsetForLeaderEpochResponses.A stateless filter for ProduceRequests.A stateless filter for ProduceResponses.A stateless filter for PushTelemetryRequests.A stateless filter for PushTelemetryResponses.A stateless filter for ReadShareGroupStateRequests.A stateless filter for ReadShareGroupStateResponses.A stateless filter for ReadShareGroupStateSummaryRequests.A stateless filter for ReadShareGroupStateSummaryResponses.A stateless filter for RemoveRaftVoterRequests.A stateless filter for RemoveRaftVoterResponses.A stateless filter for RenewDelegationTokenRequests.A stateless filter for RenewDelegationTokenResponses.A Filter that handles all request types, for example to modify the request headers.A specialization of theFilterResultfor request filters.Builder for request filter results.A Filter that handles all response types.A specialization of theFilterResultfor response filters.Builder for response filter results.A stateless filter for SaslAuthenticateRequests.A stateless filter for SaslAuthenticateResponses.A stateless filter for SaslHandshakeRequests.A stateless filter for SaslHandshakeResponses.A stateless filter for ShareAcknowledgeRequests.A stateless filter for ShareAcknowledgeResponses.A stateless filter for ShareFetchRequests.A stateless filter for ShareFetchResponses.A stateless filter for ShareGroupDescribeRequests.A stateless filter for ShareGroupDescribeResponses.A stateless filter for ShareGroupHeartbeatRequests.A stateless filter for ShareGroupHeartbeatResponses.A stateless filter for StreamsGroupDescribeRequests.A stateless filter for StreamsGroupDescribeResponses.A stateless filter for StreamsGroupHeartbeatRequests.A stateless filter for StreamsGroupHeartbeatResponses.A stateless filter for SyncGroupRequests.A stateless filter for SyncGroupResponses.A stateless filter for TxnOffsetCommitRequests.A stateless filter for TxnOffsetCommitResponses.A stateless filter for UnregisterBrokerRequests.A stateless filter for UnregisterBrokerResponses.A stateless filter for UpdateFeaturesRequests.A stateless filter for UpdateFeaturesResponses.A stateless filter for UpdateRaftVoterRequests.A stateless filter for UpdateRaftVoterResponses.A stateless filter for VoteRequests.A stateless filter for VoteResponses.A stateless filter for WriteShareGroupStateRequests.A stateless filter for WriteShareGroupStateResponses.A stateless filter for WriteTxnMarkersRequests.A stateless filter for WriteTxnMarkersResponses.