Spring Boot Starter Kafka Streams中TransformerSupplier自动注入配置问题
Here's a practical, Spring-friendly approach that meets all your requirements, keeping client code simple while enabling shared state and multiple store builders:
1. Refine CustomConfiguration for Store Builder Management
First, ensure your CustomConfiguration class properly collects and exposes state store builders using a fluent API for easy client setup:
import org.apache.kafka.streams.state.StoreBuilder; import java.util.ArrayList; import java.util.Collections; import java.util.List; public class CustomConfiguration { private final List<StoreBuilder<?>> transformationStoreBuilders = new ArrayList<>(); // Fluent method to add store builders (supports multiple entries) public CustomConfiguration addStoreBuilder(StoreBuilder<?> storeBuilder) { this.transformationStoreBuilders.add(storeBuilder); return this; } // Expose an unmodifiable list to prevent external modification public List<StoreBuilder<?>> getTransformationStoreBuilders() { return Collections.unmodifiableList(transformationStoreBuilders); } }
Clients create this as a Spring bean, adding their desired store builders:
import org.apache.kafka.streams.state.KeyValueStoreBuilder; import org.apache.kafka.streams.state.Stores; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @Configuration public class ClientKafkaConfig { @Bean public CustomConfiguration customConfiguration() { KeyValueStoreBuilder<String, String> kvStore = Stores.keyValueStoreBuilder( Stores.persistentKeyValueStore("user-profile-store"), Serdes.String(), Serdes.String() ); return new CustomConfiguration() .addStoreBuilder(kvStore) .addStoreBuilder(/* add more store builders here */); } }
2. Create a Spring-Managed Transformer Supplier Factory
Implement a factory bean that injects CustomConfiguration and handles the store builder injection into CustomTransformerSupplier instances. This factory is managed by Spring, so it has access to the configuration bean automatically:
import org.apache.kafka.streams.kstream.Transformer; import org.apache.kafka.streams.kstream.TransformerSupplier; import org.apache.kafka.streams.state.StoreBuilder; import org.springframework.stereotype.Component; import java.util.List; import java.util.Set; import java.util.function.Function; import java.util.stream.Collectors; @Component public class TransformerSupplierFactory { private final CustomConfiguration customConfiguration; // Constructor injection (Spring auto-wires CustomConfiguration) public TransformerSupplierFactory(CustomConfiguration customConfiguration) { this.customConfiguration = customConfiguration; } // Generic method to create a pre-configured TransformerSupplier public <K, V, R> TransformerSupplier<K, V, R> createSupplier( Function<List<StoreBuilder<?>>, Transformer<K, V, R>> transformerFactory) { return new TransformerSupplier<>() { @Override public Transformer<K, V, R> get() { // Pass the pre-collected store builders to the client's transformer return transformerFactory.apply(customConfiguration.getTransformationStoreBuilders()); } // Expose store names to Kafka Streams (critical for shared state) @Override public Set<String> getStoreNames() { return customConfiguration.getTransformationStoreBuilders() .stream() .map(StoreBuilder::name) .collect(Collectors.toSet()); } }; } }
3. Client Usage: Simple Supplier Creation
Clients use the factory to create their transformer suppliers—no manual injection logic required. The factory handles passing store builders to the transformer:
import org.apache.kafka.streams.StreamsBuilder; import org.apache.kafka.streams.kstream.KStream; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @Configuration public class ClientStreamConfig { private final TransformerSupplierFactory supplierFactory; // Inject the factory via constructor public ClientStreamConfig(TransformerSupplierFactory supplierFactory) { this.supplierFactory = supplierFactory; } @Bean public KStream<String, String> kStream(StreamsBuilder streamsBuilder) { KStream<String, String> inputStream = streamsBuilder.stream("input-topic"); // Create transformer supplier using the factory var transformerSupplier = supplierFactory.createSupplier( storeBuilders -> new UserProfileTransformer(storeBuilders) ); // Apply transformer (store names are automatically exposed via the supplier) inputStream.transform(transformerSupplier) .to("output-topic"); // For ProcessorSupplier, use the same pattern (create a ProcessorSupplierFactory) // inputStream.process(() -> new UserProfileProcessor(), storeNames); return inputStream; } }
4. Shared State Between Processors and Transformers
Since the store names are exposed via the supplier's getStoreNames() method, Kafka Streams registers these stores globally. Both processors and transformers can access the same store instances using the processor context:
import org.apache.kafka.streams.KeyValue; import org.apache.kafka.streams.kstream.Transformer; import org.apache.kafka.streams.processor.ProcessorContext; import org.apache.kafka.streams.state.KeyValueStore; import java.util.List; public class UserProfileTransformer implements Transformer<String, String, KeyValue<String, String>> { private final List<StoreBuilder<?>> storeBuilders; private KeyValueStore<String, String> userProfileStore; public UserProfileTransformer(List<StoreBuilder<?>> storeBuilders) { this.storeBuilders = storeBuilders; } @Override public void init(ProcessorContext context) { // Access the shared store by name userProfileStore = context.getStateStore("user-profile-store"); } @Override public KeyValue<String, String> transform(String key, String value) { // Use the store for stateful operations String existingProfile = userProfileStore.get(key); userProfileStore.put(key, value); return KeyValue.pair(key, existingProfile != null ? "Updated: " + value : "Created: " + value); } @Override public void close() { // Cleanup if needed } }
Key Benefits
- Multiple Store Builders: Clients can add any number of store builders to
CustomConfiguration. - Shared State: Kafka Streams manages a single instance of each store, accessible by both processors and transformers.
- Client Simplicity: Clients only need to define their store builders and use the factory—no manual injection or setup logic.
内容的提问来源于stack exchange,提问作者Rohit Agrawal

