【问题标题】:Cannot query local state store in Kafka Streams Application无法在 Kafka Streams 应用程序中查询本地状态存储
【发布时间】:2018-11-14 19:06:07
【问题描述】:

我正在使用 spring-kafka 构建一个 kafka 流应用程序,以按键对记录进行分组并应用一些业务逻辑。我遵循spring-kafka-streams doc 上的配置,但问题是当我想从本地存储中检索一个值时,我收到以下错误:

org.apache.kafka.streams.errors.InvalidStateStoreException: The state store, user-data-response-count, may have migrated to another instance.
  at org.apache.kafka.streams.state.internals.QueryableStoreProvider.getStore(QueryableStoreProvider.java:60)
  at org.apache.kafka.streams.KafkaStreams.store(KafkaStreams.java:1053)
  at com.umantis.management.service.UserDataManagementService.broadcastUserDataRequest(UserDataManagementService.java:121)

这是我的 KafkaStreamsConfiguration:

@Configuration
@EnableConfigurationProperties(EventsKafkaProperties.class)
@EnableKafka
@EnableKafkaStreams
public class KafkaConfiguration {

@Value("${app.kafka.streams.application-id}")
private String applicationId;

// This contains both the bootstrap servers and the schema registry url
@Autowired
private EventsKafkaProperties eventsKafkaProperties;

@Bean(name = KafkaStreamsDefaultConfiguration.DEFAULT_STREAMS_CONFIG_BEAN_NAME)
public StreamsConfig streamsConfig() {
    Map<String, Object> props = new HashMap<>();
    props.put(StreamsConfig.APPLICATION_ID_CONFIG, applicationId);
    props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, this.eventsKafkaProperties.getBrokers());
    props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName());
    props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, SpecificAvroSerde.class);
    props.put(AbstractKafkaAvroSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG, this.eventsKafkaProperties.getSchemaRegistryUrl());
    props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");

    return new StreamsConfig(props);
}

@Bean
public KGroupedStream<String, UserDataResponse> responseKStream(StreamsBuilder streamsBuilder, TopicUtils topicUtils) {
    final Map<String, String> serdeConfig = Collections.singletonMap("schema.registry.url", this.eventsKafkaProperties.getSchemaRegistryUrl());

    final Serde<UserDataResponse> valueSpecificAvroSerde = new SpecificAvroSerde<>();
    valueSpecificAvroSerde.configure(serdeConfig, false);

    return streamsBuilder
            .stream("myTopic", Consumed.with(Serdes.String(), valueSpecificAvroSerde))
            .groupByKey();
}

这是我在getKafkaStreams().store 上失败的服务代码:

@Slf4j
@Service
public class UserDataManagementService {

    private static final String RESPONSE_COUNT_STORE = "user-data-response-count";

    @Autowired
    private StreamsBuilderFactoryBean streamsBuilderFactory;

    public UserDataResponse broadcastUserDataRequest() {
        this.responseGroupStream.count(Materialized.as(RESPONSE_COUNT_STORE));

        if (!this.streamsBuilderFactory.isRunning()) {
            throw new KafkaStoreNotAvailableException();
        }

        // here we should have a single running kafka instance
        ReadOnlyKeyValueStore<String, Long> countStore =
                this.streamsBuilderFactory.getKafkaStreams().store(RESPONSE_COUNT_STORE, QueryableStoreTypes.keyValueStore());

        ...
    }

上下文:我在 Spring Boot 测试中在单个实例上运行应用程序,并确保 kafka 实例处于运行状态。我在 apache 上搜索了 this 问题的文档,但我的情况似乎不匹配。

谁能指出我做错了什么以及可能的解决方案?

我是 Kafka Streams 的新手,因此我们将不胜感激。

【问题讨论】:

    标签: java spring apache-kafka apache-kafka-streams spring-kafka


    【解决方案1】:

    好的,刚刚看到我在询问流工厂是否正在运行,但我并没有询问 kakfa 流实例是否真的在运行。

    轮询streamsBuilderFactory.getKafkaStreams().state 解决了这个问题。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-07-03
      • 1970-01-01
      • 1970-01-01
      • 2019-07-14
      • 2022-01-19
      • 1970-01-01
      相关资源
      最近更新 更多