You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spring Boot Starter Kafka Streams中TransformerSupplier自动注入配置问题

Solution for Auto-Injecting StoreBuilders into CustomTransformerSupplier

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.09 21:55:25