【问题标题】:Why does kafka streams reprocess the messages produced after broker restart为什么kafka流重新处理broker重启后产生的消息
【发布时间】:2018-11-17 09:38:21
【问题描述】:

我有一个单节点 kafka 代理和简单的流应用程序。我创建了 2 个主题(主题 1 和主题 2)。

Produced on topic1 - processed message - write to topic2

注意:对于每条生成的消息,只有一条消息写入目标主题

我生成了一条消息。写到 topic2 后,我停止了 kafka 代理。过了一段时间,我重新启动了代理并在 topic1 上产生了另一条消息。现在流应用程序处理该消息 3 次。现在,在没有停止代理的情况下,我向 topic1 生成了消息,并等待流应用程序写入 topic2,然后再次生成。

Streams 应用的行为异常。有时对于一条生成的消息,有 2 条消息写入目标主题,有时 3 条。我不明白为什么会这样。我的意思是即使在代理重启后产生的消息也会被复制。

更新 1:

我正在使用 Kafka 1.0.0 版和 Kafka-Streams 1.1.0 版

下面是代码。

Main.java

String credentials = env.get("CREDENTIALS");

props.put(StreamsConfig.APPLICATION_ID_CONFIG, "activity-collection");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());
props.put(StreamsConfig.RECONNECT_BACKOFF_MS_CONFIG, 100000);
props.put(StreamsConfig.RECONNECT_BACKOFF_MAX_MS_CONFIG, 200000);
props.put(StreamsConfig.REQUEST_TIMEOUT_MS_CONFIG, 60000);
props.put(StreamsConfig.RETRY_BACKOFF_MS_CONFIG, 60000);
props.put(StreamsConfig.producerPrefix(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG), true);
props.put(StreamsConfig.producerPrefix(ProducerConfig.ACKS_CONFIG), "all");

final StreamsBuilder builder = new StreamsBuilder();

KStream<String, String> activityStream = builder.stream("activity_contenturl");
KStream<String, String> activityResultStream = AppUtil.hitContentUrls(credentials , activityStream);
activityResultStream.to("o365_user_activity");

AppUtil.java

public static KStream<String, String> hitContentUrls(String credentials, KStream<String, String> activityStream) {

        KStream<String, String> activityResultStream = activityStream
                .flatMapValues(new ValueMapper<String, Iterable<String>>() {
                    @Override
                    public Iterable<String> apply(String value) {

                        ArrayList<String> log = new ArrayList<String>();
                        JSONObject received = new JSONObject(value);
                        String url = received.get("url").toString();

                        String accessToken = ServiceUtil.getAccessToken(credentials);
                        JSONObject activityLog = ServiceUtil.getActivityLogs(url, accessToken);

                        log.add(activityLog.toString());
                    }
                    return log;
                }                   
            });

                return activityResultStream;
    }

更新 2:

在具有上述配置的单代理和单分区环境中,我启动了 Kafka 代理和流应用程序。在源主题上产生了 6 条消息,当我在目标主题上启动消费者时,有 36 条消息并且还在计数。他们一直在来。

所以我运行这个来查看consumer-groups

kafka_2.11-1.1.0/bin/kafka-consumer-groups.sh --new-consumer --bootstrap-server localhost:9092 --list

输出:

streams-collection-app-0

接下来我运行了这个:

kafka_2.11-1.1.0/bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group streams-collection-app-0

输出:

TOPIC                    PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG             CONSUMER-ID                                                                                                                HOST            CLIENT-ID
o365_activity_contenturl 0          1               1               0               streams-collection-app-0-244b6f55-b6be-40c4-9160-00ea45bba645-StreamThread-1-consumer-3a2940c2-47ab-49a0-ba72-4e49d341daee /127.0.0.1      streams-collection-app-0-244b6f55-b6be-40c4-9160-00ea45bba645-StreamThread-1-consumer

一段时间后,输出显示如下:

TOPIC                    PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG             CONSUMER-ID     HOST            CLIENT-ID
o365_activity_contenturl 0          1               6               5               -               -               -

然后:

TOPIC                    PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG             CONSUMER-ID     HOST            CLIENT-ID
o365_activity_contenturl 0          1               7               6               -               -               -

【问题讨论】:

  • 你用的是什么版本?你能展示实际的程序吗?您是否在停止/重新启动代理之前/之后检查了提交的偏移量?
  • @MatthiasJ.Sax 请查看更新。
  • @MatthiasJ.Sax 有没有办法使用控制台消费者从__consumer_offsets 主题消费?
  • 代码看起来不错,乍一看。只是想知道如果您只想为每个输入记录发出一个输出记录,为什么要使用 flatMapValues() 而不是 mapValues()__consumer_offsets 主题是一个特殊的内部主题,不能作为其他主题访问;您可以使用命令行工具bin/kafka-consumer-groups.sh 获取偏移信息。
  • 谢谢,我已将flatMapValues() 更改为mapValues()。我将尝试使用 consumer-group 命令行工具并更新您。再次感谢所有帮助

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


【解决方案1】:

您似乎面临着已知的限制。 Kafka 主题默认存储消息至少 7 天,但提交的偏移量存储 1 天(默认配置值 offsets.retention.minutes = 1440)。因此,如果在超过 1 天的时间内没有向您的源主题生成任何消息,则在应用重新启动后,来自主题的所有消息将被再次重新处理(实际上是多次,取决于重新启动的次数,每个此类主题每天最多 1 次,很少有传入消息)。

您可以找到关于到期提交的偏移量的描述How does an offset expire for consumer group

在 kafka 版本 2.0 中,提交的偏移量的保留增加了 KIP-186: Increase offsets retention default to 7 days

为防止重新处理,您可以添加消费者属性auto.offset.reset: latest(默认值为earliest)。 latest 存在一个小风险:如果当天没有人在源主题中生成消息,然后您重新启动应用程序,您可能会丢失一些消息(仅在重新启动期间到达的消息)。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-08-31
    • 1970-01-01
    • 2018-02-21
    • 2022-01-17
    相关资源
    最近更新 更多