【问题标题】:Adding data to state store for stateful processing and fault tolerance将数据添加到状态存储以进行状态处理和容错
【发布时间】:2020-12-10 12:00:48
【问题描述】:

我有一个执行一些有状态处理的微服务。应用程序从输入主题构造一个 KStream,进行一些有状态处理,然后将数据写入输出主题。

我将在同一个组中运行 3 个这样的应用程序。当微服务宕机时,我需要存储 3 个参数,接管的微服务可以查询共享的 statestore 并从崩溃的服务停止的地方继续。

我正在考虑将这 3 个参数推送到 statestore 中,并在其他微服务接管时查询数据。从我的研究中,我看到了很多人们使用状态存储执行事件计数的例子,但这并不是我想要的,有谁知道一个例子或解决这个问题的正确方法是什么?

【问题讨论】:

    标签: apache-kafka apache-kafka-streams stateful rocksdb


    【解决方案1】:

    所以你想做两件事:

    一个。关闭的服务必须存储参数:
    如果您想以一种直接的方式进行操作,那么您所要做的就是在与状态存储相关联的主题中写一条消息(您正在使用KTable 阅读的主题)。使用Kafka Producer APIKStream(可能是kTable.toStream())来完成它,就是这样。

    否则您可以手动创建状态存储:

    // take these serde as just an example
    Serde<String> keySerde = Serdes.String();
    Serde<String> valueSerde = Serdes.String();
    KeyValueBytesStoreSupplier storeSupplier = inMemoryKeyValueStore(stateStoreName);
    streamsBuilder.addStateStore(Stores.keyValueStoreBuilder(storeSupplier, keySerde, valueSerde));
    

    然后在转换器或处理器中使用它来添加项目;你必须在转换器/处理器中声明这个:

    // depending on the serde above you might have something else then String
    private KeyValueStore<String, String> stateStore;
    

    并初始化stateStore变量:

    @Override
    public void init(ProcessorContext context) {
      stateStore = (KeyValueStore<String, String>) context.getStateStore(stateStoreName);
    }
    

    稍后使用stateStore 变量:

    @Override
    public KeyValue<String, String> transform(String key, String value) {
      // using stateStore among other actions you might take here
      stateStore.put(key, processedValue);
    }
    

    b.读取服务接管中的参数:
    您可以使用 Kafka 消费者来实现,但使用 Kafka Streams,您首先必须使存储可用;最简单的方法是创建一个 KTable;那么您必须获取使用 KTable 自动创建的可查询商店名称;然后你必须真正进入商店;然后你从存储中提取一个记录值(即通过它的键获得一个参数值)。

    // this example is a modified copy of KTable javadocs example
    final StreamsBuilder streamsBuilder = new StreamsBuilder();
    
    // Creating a KTable over the topic containing your parameters a store shall automatically be created.
    //
    // The serde for your MyParametersClassType could be 
    // new org.springframework.kafka.support.serializer.JsonSerde(MyParametersClassType.class) 
    // though further configurations might be necessary here - e.g. setting the trusted packages for the ObjectMapper behind JsonSerde.
    //
    // If the parameter-value class is a String then you could use Serdes.String() instead of a MyParametersClassType serde.
    final KTable paramsTable = streamsBuilder.table("parametersTopicName", Consumed.with(Serdes.String(), <<your InstanceOfMyParametersClassType serde>>));
    
    ...
    // see the example from KafkaStreams javadocs for more KafkaStreams related details
    final KafkaStreams streams = ...;
    streams.start()
    ...
    
    // get the queryable store name that is automatically created with the KTable
    final String queryableStoreName = paramsTable.queryableStoreName();
    // get access to the store
    ReadOnlyKeyValueStore view = streams.store(queryableStoreName, QueryableStoreTypes.timestampedKeyValueStore());
    // extract a record value from the store
    InstanceOfMyParametersClassType parameter = view.get(key);
    

    【讨论】:

    • 感谢 adrhc 的回复,所以只是为了澄清。通过执行选项 a) 将这些参数写入单独的 Kafka 主题,比如“param-stored”,这个选项不会创建状态存储或任何东西。但是在过程 b)中,当我从“参数存储”主题生成 KTable 时,这将创建一个状态存储。然后我可以从 ReadOnlyKeyValueStore 对象中正常访问存储的参数...我最初认为您不需要创建单独的主题来存储参数,而是创建一个状态存储然后将值写入其中。然后访问它以补充微服务接管的状态。
    • 从主题创建 KTable 的行为也意味着您必须从头开始阅读整个主题以重建状态。与在同一消费者组中使用状态存储进行内存访问相比,这将产生更多开销,对吧?
    • 这将是一个启动延迟,这取决于您在压缩主题中拥有多少数据:至少您将只有 1 条记录(包含一个保留这 3 个参数的对象),而最多您可能有例如最近的 1000 次更新,仍然几乎没有。
    猜你喜欢
    • 1970-01-01
    • 2016-09-11
    • 2023-02-17
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多