【问题标题】:How to log offset in KStreams Bean using spring-kafka and kafka-streams如何使用 spring-kafka 和 kafka-streams 在 KStreams Bean 中记录偏移量
【发布时间】:2020-09-07 09:40:30
【问题描述】:

我已经参考了几乎所有关于通过处理器 API 的 transform() 或 process() 方法在 KStreams 上记录偏移量的问题,就像这里的许多问题中提到的那样 -

How can I get the offset value in KStream

但我无法得到这些答案的解决方案,所以我问这个问题。

我想在每次消息被流消费时记录分区、消费者组 ID 和偏移量,我不知道如何将 process() 或 transform() 方法与 ProcessorContext API 集成?如果我在我的 CustomParser 类中实现处理器接口,那么我将不得不实现所有方法,但我不确定这是否会起作用,就像在记录元数据的汇合文档中提到的那样 - https://docs.confluent.io/current/streams/developer-guide/processor-api.html#streams-developer-guide-processor-api

我已经在下面的 spring-boot 应用程序中设置了 KStreams(更改变量名以供参考)

 @Bean
    public Set<KafkaStreams> myKStreamJson(StreamsBuilder profileBuilder) {
        Serde<JsonNode> jsonSerde = Serdes.serdeFrom(jsonSerializer, jsonDeserializer);

        final KStream<String, JsonNode> pStream = myBuilder.stream(inputTopic, Consumed.with(Serdes.String(), jsonSerde));

        Properties props = streamsConfig.kStreamsConfigs().asProperties();
       

        pstream
                .map((key, value) -> {
                            try {
                                return CustomParser.parse(key, value);
                            } catch (Exception e) {
                                LOGGER.error("Error occurred - " + e.getMessage());
                            }
                            return new KeyValue<>(null, null);
                        }
                )
                .filter((key, value) -> {
                    try {
                        return MessageFilter.filterNonNull(key, value);
                    } catch (Exception e) {
                        LOGGER.error("Error occurred - " + e.getMessage());
                    }
                    return false;
                })
                .through(
                        outputTopic,
                        Produced.with(Serdes.String(), new JsonPOJOSerde<>(TransformedMessage.class)));

        return Sets.newHashSet(
                new KafkaStreams(profileBuilder.build(), props)
        );
    }

【问题讨论】:

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


    【解决方案1】:

    实现Transformer;在init() 中保存ProcessorContext;然后,您可以访问transform() 中的记录元数据并简单地返回原始键/值。

    这是example of a Transformer。它由 Spring 提供,供 Apache Kafka 调用 Spring Integration 流来转换键/值。

    【讨论】:

    • 谢谢,添加了一个 CustomTransformer 记录详细信息,然后返回相同的键和值。
    猜你喜欢
    • 2019-03-16
    • 1970-01-01
    • 2021-02-14
    • 2017-01-10
    • 2014-08-08
    • 1970-01-01
    • 2019-09-04
    • 2018-06-28
    • 2022-01-06
    相关资源
    最近更新 更多