...
| Code Block |
|---|
| language | java |
|---|
| title | StreamsConfig |
|---|
|
public class StreamsConfig {
public static final String PROCESSOR_WRAPPER_CLASS_CONFIG = "processor.wrapper.class";
private static final String PROCESSOR_WRAPPER_CLASS_DOC = "A processor wrapper class or class name that implements the <code>org.apache.kafka.streams.state.ProcessorWrapper</code> interface. Must be passed in to the StreamsBuilder or Topology constructor in order to take effect";
} |
The ProcessorWrapper itself is defined as follows:
| Code Block |
|---|
| language | java |
|---|
| title | ProcessorWrapper |
|---|
|
package org.apache.kafka.streams.processor.api;
/**
* Wrapper class that can be used to inject custom wrappers around the processors of their application topology.
* The returned instance MUSTshould wrap the supplied {@code ProcessorSupplier} and the {@code Processor} it supplies
* to avoid disrupting the regular processing of the application, although this is not required and any processor
* implementation can be substituted in to replace the original processor entirely (which may be useful for example
* while testing or debugging an application topology).
* <p>
* NOTE: in order to use this feature, you must set the {@link StreamsConfig#PROCESSOR_WRAPPER} config and pass it
* in Returningas a new or completely different instance can have unexpected and undesirable affects. {@link TopologyConfig} when creating the {@link StreamsBuilder} or {@link Topology} by using the
* appropriate constructor (ie {@link StreamsBuilder#StreamsBuilder(TopologyConfig)} or {@link Topology#Topology(TopologyConfig)})
* <p>
* Can be configured, if desired, by implementing the {@link #configure(Map)} method,. whichThis will be invoked when
* the {@code ProcessorWrapper} is instantiated, and will provide it with the Streams application configs onceTopologyConfigs that were passed in
* to the {@code@link ProcessorWrapperStreamsBuilder} is instantiatedor {@link Topology} constructor.
*/
public interface ProcessorWrapper extends Configurable {
@Override
default void configure(final Map<String, ?> configs) {
// do nothing
}
<KIn, VIn, KOut, VOut> ProcessorSupplier<KInWrappedProcessorSupplier<KIn, VIn, KOut, VOut> wrapProcessorSupplier(final String processorName,
final ProcessorSupplier<KIn, VIn, KOut, VOut> processorSupplier);
<KIn, VIn, VOut> FixedKeyProcessorSupplier<KInWrappedFixedKeyProcessorSupplier<KIn, VIn, VOut> wrapFixedKeyProcessorSupplier(final String processorName,
final FixedKeyProcessorSupplier<KIn, VIn, VOut> processorSupplier);
} |
The return types are new interfaces introduced in this KIP for future compatibility (in case we want to add methods to the wrapped processor suppliers) and type checking (since we may want the ability to distinguish between an unwrapped processor supplier and one that has already been wrapped).
For now, these will simply be marker interfaces that extend the corresponding processor supplier class with no additional methods. They are defined below:
| Code Block |
|---|
| language | java |
|---|
| title | WrappedProcessorSupplier |
|---|
|
package org.apache.kafka.streams.processor.api;
/**
* Marker interface for classes implementing {@link ProcessorSupplier}
* that have been wrapped via a {@link ProcessorWrapper}
*/
public interface WrappedProcessorSupplier<KIn, VIn, KOut, VOut> extends ProcessorSupplier<KIn, VIn, KOut, VOut> {
} |
| Code Block |
|---|
| language | java |
|---|
| title | WrappedFixedKeyProcessorSupplier |
|---|
|
package org.apache.kafka.streams.processor.api;
/**
* Marker interface for classes implementing {@link FixedKeyProcessorSupplier}
* that have been wrapped via a {@link ProcessorWrapper}
*/
public interface WrappedFixedKeyProcessorSupplier<KIn, VIn, VOut> extends FixedKeyProcessorSupplier<KIn, VIn, VOut> {
} |
Proposed Changes
As noted above, we will introduce a new ProcessorWrapper class and its associated config. Unfortunately we have to add this to the TopologyConfig, rather than the Streams config, because we need access to the ProcessorWrapper during the topology construction and by the time the Topology is built and handed in to the KafkaStreams app, it is too late. We would like to do a cleanup of the config handling at some point as well, but will save that for a followup KIP at this time (eg StreamsConfig vs TopologyConfig with their overlapping configs, also KafkaStreams::new vs StreamsBuilder::new vs StreamsBuilder::build vs Topology::new all of which take in various forms of configs)
...