You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Kafka Connect向Elasticsearch写入数据失败:Unknown magic byte错误

问题:Kafka Connect Elasticsearch Sink写入消息失败,报错Unknown magic byte!

尝试将Kafka的pageviews2主题消息发送至Elasticsearch索引,主题为空时连接器能正常创建索引,但写入消息时失败。已添加schema.ignore、key.ignore等参数,也尝试使用内置的value.converter,均无法解决问题。


使用的配置

连接器配置1

curl -X POST -H "Content-Type: application/json" --data '{"name": "pageviews3","config": {"connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector","tasks.max": "1","topics": "pageviews2","connection.url": "http://xxxxxxxxxx:9200","type.name": "kafka-connect","connection.username": "elastic","connection.password": "xxxxxxxxxxxx"}}' http://xxxxx:8083/connectors

忽略Schema Registry的连接器配置2

curl -X POST -H "Content-Type: application/json" --data '{"name": "page-views34","config": {"connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector","connection.url": "http://xxxx:9200","tasks.max": "1","topics": "pageviews2","type.name": "_doc","value.converter": "org.apache.kafka.connect.json.JsonConverter","value.converter.schemas.enable": "false","schema.ignore": "true","key.ignore": "true","connection.username": "elastic","connection.password": "xxxxxx"}}' http://xxxxx:8083/connectors

数据生成器配置

curl -i -X PUT http://localhost:8083/connectors/datagen_local_02/config \
     -H "Content-Type: application/json" \
     -d '{"connector.class": "io.confluent.kafka.connect.datagen.DatagenConnector","key.converter": "org.apache.kafka.connect.storage.StringConverter","kafka.topic": "pageviews2","quickstart": "pageviews","max.interval": 1000,"iterations": 10000000,"tasks.max": "1"}'

报错信息

org.apache.kafka.connect.errors.ConnectException: Tolerance exceeded in error handler
    at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:223)
    at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execute(RetryWithToleranceOperator.java:149)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.convertAndTransformRecord(WorkerSinkTask.java:513)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.convertMessages(WorkerSinkTask.java:493)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:332)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:234)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:203)
    at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:189)
    at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:244)
    at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:539)
    at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)
    at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136)
    at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635)
    at java.base/java.lang.Thread.run(Thread.java:833)
Caused by: org.apache.kafka.connect.errors.DataException: Failed to deserialize data for topic pageviews2 to Avro: 
    at io.confluent.connect.avro.AvroConverter.toConnectData(AvroConverter.java:124)
    at org.apache.kafka.connect.storage.Converter.toConnectData(Converter.java:88)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.lambda$convertAndTransformRecord$3(WorkerSinkTask.java:513)
    at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndRetry(RetryWithToleranceOperator.java:173)
    at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:207)
    ... 13 more
Caused by: org.apache.kafka.common.errors.SerializationException: Unknown magic byte!
    at io.confluent.kafka.serializers.AbstractKafkaSchemaSerDe.getByteBuffer(AbstractKafkaSchemaSerDe.java:244)
    at io.confluent.kafka.serializers.AbstractKafkaAvroDeserializer$DeserializationContext.<init>(AbstractKafkaAvroDeserializer.java:334)
    at io.confluent.kafka.serializers.AbstractKafkaAvroDeserializer.deserializeWithSchemaAndVersion(AbstractKafkaAvroDeserializer.java:202)
    at io.confluent.connect.avro.AvroConverter$Deserializer.deserialize(AvroConverter.java:172)
    at io.confluent.connect.avro.AvroConverter.toConnectData(AvroConverter.java:107)
    ... 17 more

解决方案

报错核心是Unknown magic byte!,说明Elasticsearch Sink连接器在尝试用Avro反序列化消息,但主题中的消息实际不是Avro格式。

从数据生成器配置可知,Datagen连接器生成的是无Schema的JSON格式消息(仅指定key.converter为StringConverter,未配置Avro转换器),但你的第一个连接器配置未指定value.converter,会默认使用集群全局的Avro转换器,导致反序列化失败;第二个配置虽指定了JsonConverter,但可能因旧连接器实例未删除、全局配置优先级等问题未生效。

解决步骤

  1. 删除旧连接器实例:避免新旧配置冲突
    curl -X DELETE http://xxxxx:8083/connectors/pageviews3
    curl -X DELETE http://xxxxx:8083/connectors/page-views34
    
  2. 创建匹配格式的连接器配置:确保转换器与消息格式完全匹配
    curl -X POST -H "Content-Type: application/json" --data '{
      "name": "pageviews-es-sink",
      "config": {
        "connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector",
        "tasks.max": "1",
        "topics": "pageviews2",
        "connection.url": "http://xxxx:9200",
        "type.name": "_doc",
        "key.converter": "org.apache.kafka.connect.storage.StringConverter",
        "value.converter": "org.apache.kafka.connect.json.JsonConverter",
        "value.converter.schemas.enable": "false",
        "connection.username": "elastic",
        "connection.password": "xxxxxx"
      }
    }' http://xxxxx:8083/connectors
    
  3. 验证消息格式:确认主题内消息为JSON格式
    kafka-console-consumer.sh --bootstrap-server <kafka-broker>:9092 --topic pageviews2 --from-beginning
    
  4. 检查全局配置:若Connect集群全局设置了Avro转换器,需确保连接器局部配置已覆盖全局(局部配置优先级更高)

内容的提问来源于stack exchange,提问作者Onur Ulusoy

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.25 02:38:15