【问题标题】:Apache kafka consumer poll infinite loop on unexpected messagesApache kafka消费者轮询意外消息的无限循环
【发布时间】:2018-03-06 05:59:09
【问题描述】:

我正在对 Kafka 消费者实现进行集成测试。 我使用 wurstmeister/kafka docker 镜像和 Apache Kafka 消费者。 让我兴奋的场景是当我向某个主题发送“意外”消息时。 kafkaConsumer.poll(POLLING_TIMEOUT) 似乎在 RUN 模式下进入无限循环。但是,当我调试时,它会在我暂停并返回时起作用。

我在发送预期的消息时没有这个问题(不要在反序列化时抛出异常)。

这是我对 kafka 的 docker-compose 配置:

kafka:
  image: wurstmeister/kafka
  links:
    - zookeeper
  ports:
    - "9092:9092"
  environment:
    KAFKA_ADVERTISED_HOST_NAME: localhost
    KAFKA_ADVERTISED_PORT: 9092
    KAFKA_CREATE_TOPICS: "ProductLocation:1:1,ProductInformation:1:1,InventoryAvailableToSell:1:1"
    KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
  volumes:
    - /var/run/docker.sock:/var/run/docker.sock

Java 通用消费者:

@Override
public Collection<T> consume() {
    String eventToBePublishedName = ERROR_WHILE_RESETTING_OFFSET;
    boolean success = false;

    try {
        kafkaConsumer.resume(kafkaAssignments);
        if (isPollingTypeFull) {
            // dummy poll because its needed before resetting offset.
            // https://stackoverflow.com/questions/41008610/kafkaconsumer-0-10-java-api-error-message-no-current-assignment-for-partition
            kafkaConsumer.poll(POLLING_TIMEOUT);
            resetOffset();
        } else if (!offsetGotResetFirstTime) {
            resetOffset();
            offsetGotResetFirstTime = true;
        }

        eventToBePublishedName = ERROR_WHILE_POLLING;

        ConsumerRecords<Object, T> records;

        List<T> output = new ArrayList<>();

        do {
            records = kafkaConsumer.poll(POLLING_TIMEOUT);
            records.forEach(cr -> {
                T val = cr.value();
                if (val != null) {
                    output.add(cr.value());
                }
            });
        } while (records.count() > 0);

        eventToBePublishedName = CONSUMING;
        success = true;
        kafkaConsumer.pause(kafkaAssignments);
        return output;
    } finally {
        applicationEventPublisher.publishEvent(
                new OperationResultApplicationEvent(
                        this, OperationType.ConsumingOfMessages, eventToBePublishedName, success));
    }
}

反序列化:

public T deserialize(String topic, byte[] data) {
    try {
        JsonNode jsonNode = mapper.readTree(data);
        JavaType javaType = mapper.getTypeFactory().constructType(getValueClass());
        JsonNode value = jsonNode.get("value");
        return mapper.readValue(value.toString(), javaType);
    } catch (IllegalArgumentException | IOException | SerializationException e) {
        LOGGER.error("Can't deserialize data [" + Arrays.toString(data)
                + "] from topic [" + topic + "]", e);
        return null;
    }
}

在我的集成测试中,我通过发送到带有时间戳的主题名称为每个测试创建一个主题。这会创建新主题并使测试无状态。

这就是我配置 Kafka 消费者的方式:

Properties properties = new Properties();
    properties.put("bootstrap.servers", kafkaConfiguration.getServer());
    properties.put("group.id", kafkaConfiguration.getGroupId());
    properties.put("key.deserializer", kafkaConfiguration.getKeyDeserializer().getName());
    properties.put("value.deserializer", kafkaConfiguration.getValueDeserializer().getName());

【问题讨论】:

    标签: java docker apache-kafka docker-compose


    【解决方案1】:

    捕获异常并将您的承诺偏移量提前 +1 以跳过“毒丸”消息。

    【讨论】:

    • 你说的是哪个异常??我说投票是无限循环的。这是因为该主题的先前消费者尚未关闭。两者都使用非线程安全的 Kafka 消费者 ==> 只需关闭消费者即可修复它。我只是返回 null,然后对轮询的消费者记录的结果进行过滤,如代码示例中所示。偏移量是自动提交的。
    • 我说的是如何修复以前的消费者应用程序。我虽然你说之前的消费者死亡是因为它无法反序列化意外消息。在这种情况下,您应该修复该应用程序以捕获 SerDes 异常(或任何其他异常)并关闭并退出或将偏移量提前 1 并继续。
    • 消费者不会死(因此不会抛出异常),它会在轮询期间保持挂起。这与 Apache kafka 消费者线程安全有关。在创建一个新的(相同的主题和组 id)之前,我们需要关闭一个特定主题和组 id 的 kafka 消费者。这就是问题所在。
    【解决方案2】:

    如果您遇到这种情况,只需close消费者使用后,或pause使用后和resume开始使用前。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2018-09-29
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多