【问题标题】:Spring Kafka - Encountering "Magic v0 does not support record headers" errorSpring Kafka - 遇到“Magic v0 不支持记录头”错误
【发布时间】:2018-05-01 22:36:22
【问题描述】:

我正在运行一个 Spring Boot 应用程序,并在 compile('org.springframework.kafka:spring-kafka:2.1.5.RELEASE') 中进行了处理

我正在尝试使用此版本反对 Cloudera 安装:

Cloudera Distribution of Apache Kafka Version 3.0.0-1.3.0.0.p0.40 Version 0.11.0+kafka3.0.0+50

我的 KafkaProducerConfig 类非常简单:

@Configuration
public class KafkaProducerConfig {
private static final Logger LOGGER = LoggerFactory.getLogger(KafkaProducerConfig.class);

@Value("${spring.kafka.bootstrap-servers}")
private String bootstrapServers;

@Value("${spring.kafka.template.default-topic}")
private String defaultTopicName;

@Value("${spring.kafka.producer.compression-type}")
private String producerCompressionType;

@Value("${spring.kafka.producer.client-id}")
private String producerClientId;

@Bean
public Map<String, Object> producerConfigs() {
    Map<String, Object> props = new HashMap<>();

    props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, this.bootstrapServers);
    props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
    props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);
    props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, this.producerCompressionType);
    props.put(ProducerConfig.CLIENT_ID_CONFIG, this.producerClientId);
    props.put(JsonSerializer.ADD_TYPE_INFO_HEADERS, false);

    return props;
}

@Bean
public ProducerFactory<String, Pdid> producerFactory() {
    return new DefaultKafkaProducerFactory<>(producerConfigs());
}

@Bean
public KafkaTemplate<String, Pdid> kafkaTemplate() {
    KafkaTemplate<String, Pdid> kafkaTemplate = new KafkaTemplate<>(producerFactory());

    kafkaTemplate.setDefaultTopic(this.defaultTopicName);

    return kafkaTemplate;
}

@PostConstruct
public void postConstruct() {
    LOGGER.info("Kafka producer configuration: " + this.producerConfigs().toString());
    LOGGER.info("Kafka topic name: " + this.defaultTopicName);
}

}

当我启动应用程序时,我收到:

2018-05-01 17:15:41.355 INFO 54674 --- [nio-9000-exec-2] o.a.kafka.common.utils.AppInfoParser : Kafka version : 1.0.1 2018-05-01 17:15:41.356 INFO 54674 --- [nio-9000-exec-2] o.a.kafka.common.utils.AppInfoParser : Kafka commitId : c0518aa65f25317e

然后,我发送一个有效载荷。它针对该主题显示在 Kafka 工具中。但是,在尝试摄取数据时,在 Kafka 端的日志中,我收到:

[KafkaApi-131] Error when handling request {replica_id=-1,max_wait_time=500,min_bytes=1,topics=[{topic=profiles-pdid,partitions=[{partition=0,fetch_offset=7,max_bytes=1048576}]}]}java.lang.IllegalArgumentException: Magic v0 does not support record headers
at org.apache.kafka.common.record.MemoryRecordsBuilder.appendWithOffset(MemoryRecordsBuilder.java:385)
at org.apache.kafka.common.record.MemoryRecordsBuilder.append(MemoryRecordsBuilder.java:568)
at org.apache.kafka.common.record.AbstractRecords.convertRecordBatch(AbstractRecords.java:117)
at org.apache.kafka.common.record.AbstractRecords.downConvert(AbstractRecords.java:98)
at org.apache.kafka.common.record.FileRecords.downConvert(FileRecords.java:245)
at kafka.server.KafkaApis$$anonfun$kafka$server$KafkaApis$$convertedPartitionData$1$1$$anonfun$apply$5.apply(KafkaApis.scala:523)
at kafka.server.KafkaApis$$anonfun$kafka$server$KafkaApis$$convertedPartitionData$1$1$$anonfun$apply$5.apply(KafkaApis.scala:521)
at scala.Option.map(Option.scala:146)
at kafka.server.KafkaApis$$anonfun$kafka$server$KafkaApis$$convertedPartitionData$1$1.apply(KafkaApis.scala:521)
at kafka.server.KafkaApis$$anonfun$kafka$server$KafkaApis$$convertedPartitionData$1$1.apply(KafkaApis.scala:511)
at scala.Option.flatMap(Option.scala:171)
at kafka.server.KafkaApis.kafka$server$KafkaApis$$convertedPartitionData$1(KafkaApis.scala:511)
at kafka.server.KafkaApis$$anonfun$kafka$server$KafkaApis$$createResponse$2$1.apply(KafkaApis.scala:559)
at kafka.server.KafkaApis$$anonfun$kafka$server$KafkaApis$$createResponse$2$1.apply(KafkaApis.scala:558)
at scala.collection.Iterator$class.foreach(Iterator.scala:891)
at scala.collection.AbstractIterator.foreach(Iterator.scala:1334)
at scala.collection.IterableLike$class.foreach(IterableLike.scala:72)
at scala.collection.AbstractIterable.foreach(Iterable.scala:54)
at kafka.server.KafkaApis.kafka$server$KafkaApis$$createResponse$2(KafkaApis.scala:558)
at kafka.server.KafkaApis$$anonfun$kafka$server$KafkaApis$$fetchResponseCallback$1$1.apply$mcVI$sp(KafkaApis.scala:579)
at kafka.server.ClientQuotaManager.recordAndThrottleOnQuotaViolation(ClientQuotaManager.scala:196)
at kafka.server.KafkaApis.sendResponseMaybeThrottle(KafkaApis.scala:2014)
at kafka.server.KafkaApis.kafka$server$KafkaApis$$fetchResponseCallback$1(KafkaApis.scala:578)
at kafka.server.KafkaApis$$anonfun$kafka$server$KafkaApis$$processResponseCallback$1$1.apply$mcVI$sp(KafkaApis.scala:598)
at kafka.server.ClientQuotaManager.recordAndThrottleOnQuotaViolation(ClientQuotaManager.scala:196)
at kafka.server.ClientQuotaManager.recordAndMaybeThrottle(ClientQuotaManager.scala:188)
at kafka.server.KafkaApis.kafka$server$KafkaApis$$processResponseCallback$1(KafkaApis.scala:597)
at kafka.server.KafkaApis$$anonfun$handleFetchRequest$1.apply(KafkaApis.scala:614)
at kafka.server.KafkaApis$$anonfun$handleFetchRequest$1.apply(KafkaApis.scala:614)
at kafka.server.ReplicaManager.fetchMessages(ReplicaManager.scala:639)
at kafka.server.KafkaApis.handleFetchRequest(KafkaApis.scala:606)
at kafka.server.KafkaApis.handle(KafkaApis.scala:98)
at kafka.server.KafkaRequestHandler.run(KafkaRequestHandler.scala:66)
at java.lang.Thread.run(Thread.java:748)

我从 Producer 应用程序端尝试了以下方法:

  1. 降级到 Spring Kafka 2.0.4。我希望降到 Kafka 版本 0.11.0 可以帮助解决这个问题,但它没有任何效果。
  2. 已验证节点都是相同的版本。根据我的管理员的说法,他们是。
  3. 经我的管理员验证,我们没有混合安装。再一次,有人告诉我我们不这样做。
  4. 基于类似的 Stack Overflow 问题,我回到 2.1.5 并尝试将 JsonSerializer.ADD_TYPE_INFO_HEADERS 设置为 false。我想也许它会删除日志所指的标题。再一次,没有去,同样的错误被记录下来。

我希望我遗漏了一些明显的东西。是否需要打开/关闭任何其他标题以帮助解决任何人都知道的 Magic v0 问题?

我们有其他应用程序在同一环境中成功写入其他主题,但它们是手工制作必要的 Spring bean 的旧应用程序。此外,这些应用程序还使用更旧的客户端 (0.8.2.2),并且它们使用 StringSerializer 作为 Producer 值而不是 JSON。我需要我的数据是 JSON 格式,当我们在一个应该支持 0.11.x 的系统上时,我真的不想降级到 0.8.2.2。

【问题讨论】:

    标签: apache-kafka cloudera flume spring-kafka


    【解决方案1】:

    但它们是手工制作必要的 Spring bean 的旧应用程序。

    在 org.apache.kafka.common.record.FileRecords。 downConvert (FileRecords.java:245)

    我不熟悉 kafka 代理的内部结构,但它“听起来”像是主题是使用旧代理创建的,并且它们的格式不支持标头,而不是代理版本本身(提示:downConvert)。

    您是否与干净的经纪人一起尝试过这个?

    只要您不尝试使用代理不支持的功能,1.0.x 客户端就可以与旧代理(回到 0.10.2.x IIRC)通信。您的代理是 0.11(确实支持标头)这一事实进一步表明问题出在主题记录格式。

    我已经成功地测试了 up/down broker/client/topic 版本没有问题,只要你使用通用功能子集。

    JsonSerializer.ADD_TYPE_INFO_HEADERS 为 false。

    这应该会阻止框架设置任何标题;您需要显示您的生产者代码(以及所有配置)。

    您还可以将ProducerInterceptor 添加到生产者配置并检查onSend() 方法中的ProducerRecord headers 属性,以确定输出消息中设置的标头。

    如果您使用的是 spring-messaging 消息(template.setn(Message&lt;?&gt; m),默认情况下将映射标题)。使用原始 template.send() 方法不会设置任何标头(除非您发送带有标头的 ProducerRecord

    【讨论】:

    • 感谢您的建议/澄清。我已经更新了我的原始帖子,以获得完整的 Kafka 生产者配置,而不仅仅是我最初发布的部分。
    • 谢谢,但是您在模板上使用的是哪个send() 方法?
    • 我们使用的是kafkaTemplate.sendDefault,默认主题设置在application.yml文件中。
    【解决方案2】:

    问题的解决方案是两件事的结合:

    1. 按照 Gary Russell 的建议,我需要添加 JsonSerializer.ADD_TYPE_INFO_HEADERS to false
    2. 我需要刷新所有已放入主题的记录之前配置已放入我的应用程序。之前的记录有标题,这会破坏 Flume 消费者。

    【讨论】:

    • 啊哈 - 所以downConvert 是因为主题中已有的消息有标题,但消费者不支持它们。
    【解决方案3】:

    我已将我的应用程序升级到 Spring boot 2x,但我遇到了一些与 kafka 客户端依赖项的兼容性问题(请参阅 Spring-boot and Spring-Kafka compatibility matrix),所以我也必须升级它。另一方面,我在服务器上运行了一个较旧的代理(kafka 0.10),然后我无法向它发送消息。我还意识到,即使将 JsonSerializer.ADD_TYPE_INFO_HEADERS 设置为 false,kafka 生产者也在内部设置标头,并且由于魔法是根据 kafka 的版本(在 RecordBatch 中)固定的,所以在这种情况下没有办法不倒下进入MemoryRecordsBuilder.appendWithOffset 上的条件:if (magic &lt; RecordBatch.MAGIC_VALUE_V2 &amp;&amp; headers != null &amp;&amp; headers.length &gt; 0) throw new IllegalArgumentException("Magic v" + magic + " does not support record headers");。 最后,我解决这个问题的唯一方法是升级我的 kafka 服务器。

    【讨论】:

      猜你喜欢
      • 2018-06-05
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-02-25
      • 2020-06-28
      • 1970-01-01
      相关资源
      最近更新 更多