【问题标题】:Transactional transfer messages between kafka clusters with spring使用spring的kafka集群之间的事务传输消息
【发布时间】:2019-09-15 20:28:24
【问题描述】:

我有两个 kafka 集群。我需要使用 kafka-spring 在它们之间实现某种同步。

[cluster A, topic A]  <-- [spring app] --> [cluster B, topic B]

我创建了注释为@Transactional 的侦听器,它使用 kafkaTemplate 发布消息。当两个集群都有连接时,这非常有效。当与目标集群的连接丢失时 - 侦听器似乎仍在确认新消息,但它们没有发布。我尝试对侦听器进行手动破解,禁用自动提交等,但它们似乎没有像我认为的那样工作。当连接恢复在线时,消息永远不会被传递。在这方面需要帮助。

    @KafkaListener(topics = "A", containerFactory = "syncLocalListenerFactory")
    public void consumeLocal(@Header(KafkaHeaders.RECEIVED_MESSAGE_KEY) String key, @Payload SyncEvent message, Acknowledgment ack) {
        kafkaSyncRemoteTemplate.send("B", key, message);
        ack.acknowledge();
    }

我正在获取日志:

2019-04-26 12:11:40.808  WARN 21304 --- [ad | producer-1] org.apache.kafka.clients.NetworkClient   : [Producer clientId=producer-1] Connection to node 1001 could not be established. Broker may not be available.
2019-04-26 12:11:40.828  WARN 21304 --- [ntainer#0-0-C-1] org.apache.kafka.clients.NetworkClient   : [Consumer clientId=consumer-2, groupId=app-sync] Connection to node 1001 could not be established. Broker may not be available.
2019-04-26 12:11:47.829 ERROR 21304 --- [ad | producer-1] o.s.k.support.LoggingProducerListener    : Exception thrown when sending a message with key='...' and payload='...' to topic B:

org.apache.kafka.common.errors.TimeoutException: Expiring 2 record(s) for sync-2: 30002 ms has passed since batch creation plus linger time

2019-04-26 12:11:47.829 ERROR 21304 --- [ad | producer-1] o.s.k.support.LoggingProducerListener    : Exception thrown when sending a message with key='...' and payload='...' to topic B:

org.apache.kafka.common.errors.TimeoutException: Expiring 2 record(s) for sync-2: 30002 ms has passed since batch creation plus linger time

--- 编辑---

这里的kafkaProperties是从application.properties文件中读取的默认kafka-spring属性,但是在这种情况下它们都是默认的

    @Bean
    public ConsumerFactory<String, SyncEvent> syncLocalConsumerFactory() {
        Map<String, Object> config = kafkaProperties.buildConsumerProperties();

        config.put(ConsumerConfig.GROUP_ID_CONFIG, kafkaProperties.getStreams().getApplicationId() + "-sync");
        config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class);
        config.put(JsonDeserializer.VALUE_DEFAULT_TYPE, SyncEvent.class);
        config.put(JsonDeserializer.TRUSTED_PACKAGES, "app.structures");

        config.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);

        DefaultKafkaConsumerFactory<String, SyncEvent> cf = new DefaultKafkaConsumerFactory<>(config);
        cf.setValueDeserializer(new JsonDeserializer<>(SyncEvent.class, objectMapper));
        return cf;
    }

    @Bean(name = "syncLocalListenerFactory")
    public ConcurrentKafkaListenerContainerFactory<String, SyncEvent> kafkaSyncLocalListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, SyncEvent> factory = new ConcurrentKafkaListenerContainerFactory();
        factory.setConsumerFactory(syncLocalConsumerFactory());
        factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE);
        factory.getContainerProperties().setAckOnError(false);
        factory.setErrorHandler(new SeekToCurrentErrorHandler(0));
        return factory;
    }

【问题讨论】:

  • Kafka 有一个名为 MirrorMaker 的内置工具来执行此操作
  • @cricket_007 我知道,但我不想运行其他应用程序(两种方式我都需要两个)。我想做同样的事情,但要与我更大的应用程序集成。我正在尝试进行逆向工程,但这并不容易,因为 MirrorMaker 是在 scala 中编写的,没有任何类似弹簧的功能,只是普通的 kafka :(

标签: spring apache-kafka kafka-consumer-api kafka-producer-api spring-kafka


【解决方案1】:

This website 描述了如何设置错误处理程序(使用SeekToCurrentErrorHandler),它可能会对您有所帮助。来自Spring documentation

SeekToCurrentErrorHandler:一个错误处理程序,用于查找每个主题的当前偏移量 剩余的记录。用于在消息后回退分区 失败,以便可以重播。

【讨论】:

  • 它没有改变任何东西。似乎 kafkaTemplate.send() 返回没有错误,但只是将消息排队等待在内存中传递,然后在一段时间后将其过期。如何改变这种行为?
  • 您的 sn-p 提到了使用名为“syncLocalListenerFactory”的KafkaListenerContainerFactory bean。你能发布那个bean的代码sn-p吗?它可能有助于我们查看您是如何定义重试处理程序和错误处理程序的。
  • 发现 kafkaTemplate 返回 Future,所以我添加了 call .get() ,它应该等待结果并在失败时抛出 RuntimeException 并在成功时抛出 Acknowlege::ack() 但仍然相同......跨度>
【解决方案2】:

发生这种情况是因为 kafka 事务不能跨集群。您的 @Transactional 注释没有任何意义,因此无论发布到集群 B 是否成功,偏移量都会提交给集群 A。

您目前可以为跨集群流实现的最佳保证是“至少一次”处理,您可以通过确保仅在集群 B 的目标代理确认消息后将偏移量提交给集群 A 来实现.

有关更多信息,请参阅我的博客文章 - https://medium.com/@harelopler/kafka-cross-cluster-stream-reaching-at-least-once-semantics-c74ed0eb1a54

【讨论】:

    猜你喜欢
    • 2017-12-12
    • 2023-03-20
    • 2018-07-09
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-07-16
    • 2018-11-21
    • 2023-01-25
    相关资源
    最近更新 更多