【问题标题】:Why Kafka KTable is missing entries?为什么 Kafka KTable 缺少条目?
【发布时间】:2017-04-06 14:22:27
【问题描述】:

我有一个使用来自 Kafka Streams 的 KTable 的单实例 Java 应用程序。直到最近,当某些消息突然消失时,我才可以使用 KTable 检索所有数据。那里应该有大约 33k 条带有唯一键的消息。

当我想按键检索消息时,我没有收到一些消息。我使用 ReadOnlyKeyValueStore 来检索消息:

final ReadOnlyKeyValueStore<GenericRecord, GenericRecord> store = ((KafkaStreams)streams).store(storeName, QueryableStoreTypes.keyValueStore());
store.get(key);

这些是我为 KafkaStreams 设置的配置设置。

final Properties config = new Properties();
config.put(StreamsConfig.APPLICATION_SERVER_CONFIG, serverId);
config.put(StreamsConfig.APPLICATION_ID_CONFIG, applicationId);
config.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
config.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
config.put(AbstractKafkaAvroSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG, schemaRegistryUrl);
config.put(StreamsConfig.KEY_SERDE_CLASS_CONFIG, GenericAvroSerde.class);
config.put(StreamsConfig.VALUE_SERDE_CLASS_CONFIG, GenericAvroSerde.class);
config.put(StreamsConfig.CACHE_MAX_BYTES_BUFFERING_CONFIG, 0);

Kafka:0.10.2.0-cp1
融合:3.2.0

调查给我带来了一些非常令人担忧的见解。使用 REST 代理我手动读取分区并发现一些偏移量返回错误。

请求: /topics/{topic}/partitions/{partition}/messages?offset={offset}

{
    "error_code": 50002,
    "message": "Kafka error: Fetch response contains an error code: 1"
}

没有客户端,java 和命令行都没有返回任何错误。他们只是跳过 faulty 丢失的消息,导致 KTables 中的数据丢失。一切都很好,似乎有些消息不知何故损坏了。

我有两个代理,所有主题的复制因子都是 2 并且完全复制。两个经纪人分别返回相同的。重启代理没有任何区别。

  • 可能是什么原因?
  • 如何在客户端检测到这种情况?

【问题讨论】:

  • 我不知道StoreManager 是什么——这不是 Kafka Streams 的一部分。您使用窗口化的还是非窗口化的 KTable?你使用什么版本的 Kafka Streams?
  • @MatthiasJ.Sax 抱歉我的错误我把问题做得更精确了。
  • 感谢您的更新。这听起来真的很奇怪。 “他们只是跳过导致数据丢失的错误消息”——这听起来也很奇怪——AFAIK,消费者没有“跳过”消息的内置机制。也许你应该在 Kafka 用户列表中询问 kafka.apache.org/contact(这甚至可能是一个错误......)——但这似乎不是 Kafka Streams 的问题,因为 Kafka Streams 内部只是使用了 Kafka Consumer——这个,如果消费者的行为很奇怪,Kafka Streams 无法解决这个问题。
  • 是否会因为日志压缩或超过配置的保留期而删除已消失的消息?
  • @HansJespersen 可耻地我必须承认问题就是这么简单。谢谢。

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


【解决方案1】:

通过default Kafka Broker 配置键cleanup.policy 设置为delete。将其设置为compact 以保留每个键的最新消息。 See compaction.

删除旧消息不会更改最小偏移量,因此尝试检索低于它的消息会导致错误。错误非常模糊。 Kafka Streams 客户端将从最小偏移量开始读取消息,因此没有错误。唯一可见的影响是 KTables 中缺少数据。

由于caches,应用程序正在运行时,即使从 Kafka 本身删除消息后,所有数据仍可能可用。它们会在清理后消失。

【讨论】:

    猜你喜欢
    • 2021-04-15
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-09-07
    • 1970-01-01
    • 2019-10-31
    • 2014-06-12
    • 1970-01-01
    相关资源
    最近更新 更多