【问题标题】:transactional behavior in spring-kafkaspring-kafka 中的事务行为
【发布时间】: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


    【解决方案1】:

    the documentation for the upcoming 2.2 release (due tomorrow)

    DefaultAfterRollbackProcessor(使用事务时)和SeekToCurrentErrorHandler(不使用事务时)现在可以恢复(跳过)一直失败的记录,默认情况下会在 10 次失败后恢复。它们可以配置为将失败的记录发布到死信主题。

    另请参阅Error Handling Deserializer,它捕获反序列化问题并将它们传递给容器,以便将它们发送到错误处理程序。

    【讨论】:

    • 感谢我正在尝试阅读此内容。只是第一句话:为什么:“ErrorHandlingDeserializer 实现 ExtendedDeserializer”——这有点限制了它的可用性。如果一个人有类型正确的反序列化器和类型正确的工厂,他就不能使用 ErrorHandlingDeserializer,因为它没有类型参数,这是被阻塞的。 IIUC 您无法修复以下代码。 ConsumerFactory consumerFactory = new DefaultKafkaConsumerFactory(createKafkaConsumerFactoryProperties(bootstrapServers, groupId), new StringDeserializer(), new AvroDeserializer(SomeClass.class));
    • 是的——应该是&lt;? extends Object&gt;。但这很简单;您可以使用正确的类型信息制作自己的副本。
    • 当然,但我会重新考虑更新。即使有人依赖 ...,添加此更改也不会破坏任何内容
    • ... 并且使用这个处理程序与我之前返回的 KafkaNull 类似,只是现在我有特定的数据(这确实是改进),但是在消费者方法中,我仍然不能使用泛型,因为我得到预期的数据或 DeserializationException。当然,我可以使用过滤器过滤掉这些消息,但我不确定我是否可以像将它们发送到错误主题一样处理它们(对不起,我还没有读到那部分)
    • ?如果您使用基于记录的侦听器,DE 将直接进入错误处理程序 - 参见 github.com/spring-projects/spring-kafka/blob/master/… - 对于基于批处理的侦听器,我们别无选择,只能将其交给侦听器。
    猜你喜欢
    • 2018-05-01
    • 1970-01-01
    • 2018-10-25
    • 1970-01-01
    • 2021-10-31
    • 1970-01-01
    • 2019-11-30
    • 1970-01-01
    • 2020-03-07
    相关资源
    最近更新 更多