【发布时间】: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