【问题标题】:Kafka ConsumerInterceptor onCommit not being called when using transactions使用事务时未调用 Kafka ConsumerInterceptor onCommit
【发布时间】:2022-10-17 20:56:57
【问题描述】:

我在 Spring Boot 应用程序中使用 Spring Kafka。我正在尝试使用 Kafka ConsumerInterceptor 在提交偏移量时进行拦截。

这似乎工作生产者事务未启用,但事务已打开 Interceptor::onCommit 不再被调用。

以下最小示例一切都按预期工作:

@SpringBootApplication
@EnableKafka
class Application {
    @KafkaListener(topics = ["test"])
    fun onMessage(message: String) {
        log.warn("onMessage: $message")
    }

拦截器:

class Interceptor : ConsumerInterceptor<String, String> {
    override fun onCommit(offsets: MutableMap<TopicPartition, OffsetAndMetadata>) {
        log.warn("onCommit: $offsets")
    }

    override fun onConsume(records: ConsumerRecords<String, String>): ConsumerRecords<String, String> {
        log.warn("onConsume: $records")
        return records
    }
}

应用配置:

spring:
  kafka:
    consumer:
      enable-auto-commit: false
      auto-offset-reset: earliest
      properties:
        "interceptor.classes": com.example.Interceptor
      group-id: test-group
    listener:
      ack-mode: record

在使用@EmbeddedKafka 的测试中:

    @Test
    fun sendMessage() {
        kafkaTemplate.send("test", "id", "sent message").get() // block so we don't end before the consumer gets the message
    }

这输出了我所期望的:

onConsume: org.apache.kafka.clients.consumer.ConsumerRecords@6a646f3c
onMessage: sent message
onCommit: {test-0=OffsetAndMetadata{offset=1, leaderEpoch=null, metadata=''}}

但是,当我通过提供transaction-id-prefix 启用事务时,不再调用InterceptoronCommit

我更新的配置仅添加:

spring:
  kafka:
    producer:
      transaction-id-prefix: tx-id-

并且测试更新为将send 包装在事务中:

    @Test
    fun sendMessage() {
        kafkaTemplate.executeInTransaction {
            kafkaTemplate.send("test", "a", "sent message").get()
        }
    }

通过此更改,我的日志输出现在只有

onConsume: org.apache.kafka.clients.consumer.ConsumerRecords@738b5968
onMessage: sent message

调用InterceptoronConsume 方法,@KafkaListener 接收消息,但从未调用onCommit

有谁碰巧知道这里发生了什么?我对我应该在这里看到的内容的期望是否不正确?

【问题讨论】:

    标签: spring-boot kotlin apache-kafka spring-kafka


    【解决方案1】:

    使用事务时不会通过消费者提交偏移量(恰好一次语义)。相反,偏移量是通过生产者提交的。

    /**
     * Sends a list of specified offsets to the consumer group coordinator, and also marks
     * those offsets as part of the current transaction. These offsets will be considered
     * committed only if the transaction is committed successfully. The committed offset should
     * be the next message your application will consume, i.e. lastProcessedMessageOffset + 1.
     * <p>
     * This method should be used when you need to batch consumed and produced messages
     * together, typically in a consume-transform-produce pattern. Thus, the specified
     * {@code groupMetadata} should be extracted from the used {@link KafkaConsumer consumer} via
     * {@link KafkaConsumer#groupMetadata()} to leverage consumer group metadata. This will provide
     * stronger fencing than just supplying the {@code consumerGroupId} and passing in {@code new ConsumerGroupMetadata(consumerGroupId)},
     * however note that the full set of consumer group metadata returned by {@link KafkaConsumer#groupMetadata()}
     * requires the brokers to be on version 2.5 or newer to understand.
     *
     * <p>
     * Note, that the consumer should have {@code enable.auto.commit=false} and should
     * also not commit offsets manually (via {@link KafkaConsumer#commitSync(Map) sync} or
     * {@link KafkaConsumer#commitAsync(Map, OffsetCommitCallback) async} commits).
     * This method will raise {@link TimeoutException} if the producer cannot send offsets before expiration of {@code max.block.ms}.
     * Additionally, it will raise {@link InterruptException} if interrupted.
     *
     * @throws IllegalStateException if no transactional.id has been configured or no transaction has been started.
     * @throws ProducerFencedException fatal error indicating another producer with the same transactional.id is active
     * @throws org.apache.kafka.common.errors.UnsupportedVersionException fatal error indicating the broker
     *         does not support transactions (i.e. if its version is lower than 0.11.0.0) or
     *         the broker doesn't support latest version of transactional API with all consumer group metadata
     *         (i.e. if its version is lower than 2.5.0).
     * @throws org.apache.kafka.common.errors.UnsupportedForMessageFormatException fatal error indicating the message
     *         format used for the offsets topic on the broker does not support transactions
     * @throws org.apache.kafka.common.errors.AuthorizationException fatal error indicating that the configured
     *         transactional.id is not authorized, or the consumer group id is not authorized.
     * @throws org.apache.kafka.clients.consumer.CommitFailedException if the commit failed and cannot be retried
     *         (e.g. if the consumer has been kicked out of the group). Users should handle this by aborting the transaction.
     * @throws org.apache.kafka.common.errors.FencedInstanceIdException if this producer instance gets fenced by broker due to a
     *                                                                  mis-configured consumer instance id within group metadata.
     * @throws org.apache.kafka.common.errors.InvalidProducerEpochException if the producer has attempted to produce with an old epoch
     *         to the partition leader. See the exception for more details
     * @throws KafkaException if the producer has encountered a previous fatal or abortable error, or for any
     *         other unexpected error
     * @throws TimeoutException if the time taken for sending offsets has surpassed max.block.ms.
     * @throws InterruptException if the thread is interrupted while blocked
     */
    public void sendOffsetsToTransaction(Map<TopicPartition, OffsetAndMetadata> offsets,
                                         ConsumerGroupMetadata groupMetadata) throws ProducerFencedException {
    
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-11-30
      • 1970-01-01
      • 1970-01-01
      • 2019-08-12
      • 1970-01-01
      相关资源
      最近更新 更多