【发布时间】:2018-10-27 22:08:43
【问题描述】:
我反复阅读 spring-kafka/kafka 文档,但仍然找不到方法,如何通过错误恢复来执行正确的事务行为。我相信这不是一个微不足道的问题,所以请阅读到最后。我相信整个这个问题都围绕着寻找如何重新定位失败记录或如何确认错误处理程序的方法。但也许有更好的方法,我不知道。
所以记录在流入,其中一些是无效的。我想作为一个最小的解决方案是(然后我将解决你可能会看到的几个问题):
1) 如果发生一些小事故,例如一个或几个无效记录,我们无法负担停止生产的奢侈。因此,如果kafka主题中存在无效记录,我想将其记录下来,或者将其重新发送到不同的队列,然后继续处理以下记录。
2) 存在永久和临时故障。永久失败是记录无法反序列化,记录数据验证失败。在这种情况下,我想跳过无效记录,如 1) 中所述。临时失败可能是一些特定的异常或状态,例如数据库连接错误、网络问题等。在这种情况下,我们不想跳过失败的记录,我们想在延迟一段时间后重试。
这个问题的主题只是实现跳过/不跳过行为。
可以说,这是我们的出发点:
private Map<String, Object> createKafkaConsumerFactoryProperties(String bootstrapServers, String groupId, Class<?> valueDeserializerClass) {
Map<String, Object> props = new HashMap<>();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, valueDeserializerClass);
props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
return props;
}
@Bean(name="SomeFactory")
public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory(
@Value("${…}") String bootstrapServers,
@Value("${…}") String groupId) {
ConcurrentKafkaListenerContainerFactory<String, String> factory =
new ConcurrentKafkaListenerContainerFactory<>();
ConsumerFactory<String, String> consumerFactory = new DefaultKafkaConsumerFactory<>(
createKafkaConsumerFactoryProperties(bootstrapServers, groupId, AvroDeserializer.class),
new StringDeserializer(),
new AvroDeserializer(SomeClass.class));
factory.setConsumerFactory(consumerFactory);
// factory.setConcurrency(2);
// factory.setBatchListener(true);
return factory;
}
我们有这样的听众:
@KafkaListener(topics = "${…}", containerFactory = "SomeFactory")
public void receive(@Valid List<SomeClass> messageList) {/*logic*/}
如果我理解正确,现在它的表现如何:
当监听器收到消息时,~当我们到达receive方法内部时,kafka消息已经被确认,如果receive方法抛出异常,下一次轮询将返回以下记录。因为 ack 发生了,并且我们没有定义错误处理程序,因此记录错误处理程序将启动。这不一定是我们想要的。我们可以使用 SeekToCurrentErrorHandler 重新处理消息。或者可以指定 TransactionManager,如果从侦听器中“泄漏”异常,也会发生重新定位。如果有人知道这两种方法的性能比较,请告诉我。
当消息无法反序列化时,反序列化器将失败,消息将不会被确认,并且将再次轮询相同的记录。这是某种“毒包”,因为 kafka 将无限期地旋转此消息。我们确实有 retry.backoff.ms 至少可以减慢它的速度,但我看不到任何最大重试次数或其他东西。所以我们能做的最好的事情就是在这种情况下停止/暂停容器。这太苛刻了。顺便提一句。我是 kafka/spring-kafka 的新手,我没有看到任何提及,如何从应用程序外部手动重新定位偏移量,这意味着好的,侦听器已关闭,但现在呢?另一种解决方案是不使反序列化器失败,并返回一些东西。但是什么?? KafkaNull,很好,但是我们的监听器会因为 SomeClass ClassCastException 而失败。我们可以发送一些 SomeClass 的人工值,这又是可怕的,因为这不是我们实际得到的数据。这在架构上也是不正确的。
或者我们可以使用重新定位错误处理程序,如果我们知道如何做到这一点,那就太好了。我需要寻找下一个记录。但是,虽然文档说,ErrorHandler 应该传达导致失败的记录,但它似乎没有这样做。因此,即使在非批处理侦听器中,我也有记录列表(1 个失败 + 一堆未处理),并且不知道将偏移设置到哪里。
那么解决这种疯狂的方法是什么? 好吧,我现在能想到的最好的方法是非常丑陋:不要在反序列化器中失败(坏),不接受侦听器中的特定类型(坏),手动过滤掉 KafkaNulls(坏),最后手动触发 bean 验证(坏) .有没有更好的办法?感谢您的示例,我将不胜感激给出如何实现这一目标的每一个提示或指导。
【问题讨论】:
标签: spring-kafka