KSQL查询Stream时Kafka主题Avro反序列化错误求助
问题:KSQL查询Avro Stream时出现反序列化错误
错误信息
ERROR {"type":0,"deserializationError":{"target":"value","errorMessage":"Error deserializing message from topic: deployment","recordB64":"HnloOTgyaGl1cXN4aG85OVZnaXRodWIuY29tL21pbmR0aWNrbGUvY29udGVudC1zeW5jLXdvcmtmbG93KmNvbnRlbnQtc3luYy13b3JrZmxvdwZjbnQ=","cause":["Failed to deserialize data for topic deployment to Avro: ","Unknown magic byte!"],"topic":"deployment"},"recordProcessingError":null,"productionError":null,"serializationError":null,"kafkaStreamsThreadError":null} (processing.transient_DEPLOYMENT_603110738765559581.KsqlTopic.Source.deserializer) ksql-server | [2023-07-22 21:31:59,605] WARN stream-thread [_confluent-ksql-ksql-examples-service-idtransient_transient_DEPLOYMENT_603110738765559581_1690061516393-80259bf8-0a07-46fa-b86b-864a9a359fc7-StreamThread-1] task [0_1] Skipping record due to deserialization error. topic=[deployment] partition=[1] offset=[29] (org.apache.kafka.streams.processor.internals.RecordDeserializer)
环境与操作步骤
创建Stream语句
Create Stream deployment with (KAFKA_TOPIC='deployment', VALUE_FORMAT='AVRO', KEY_FORMAT='AVRO', PARTITIONS=2, VALUE_SCHEMA_ID=1, KEY_SCHEMA_ID=1);
Docker Compose配置
--- version: '2' services: zookeeper: image: 'confluentinc/cp-zookeeper:7.3.2' hostname: zookeeper container_name: zookeeper networks: - ksql-poc ports: - '32181:32181' environment: ZOOKEEPER_CLIENT_PORT: 32181 ZOOKEEPER_TICK_TIME: 2000 kafka: image: 'confluentinc/cp-kafka:7.3.2' hostname: kafka container_name: kafka networks: - ksql-poc ports: - '9092:9092' - '29092:29092' depends_on: - zookeeper environment: KAFKA_BROKER_ID: 1 KAFKA_ZOOKEEPER_CONNECT: zookeeper:32181 KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:29092,PLAINTEXT_HOST://localhost:9092 KAFKA_AUTO_CREATE_TOPICS_ENABLE: 'true' KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1 KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1 schema-registry: image: 'confluentinc/cp-schema-registry:7.3.2' hostname: schema-registry container_name: schema-registry networks: - ksql-poc depends_on: - zookeeper - kafka ports: - '8081:8081' environment: SCHEMA_REGISTRY_HOST_NAME: schema-registry SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS: 'kafka:29092' SCHEMA_REGISTRY_KAFKASTORE_CONNECTION_URL: zookeeper:32181 # Runs the Kafka KSQL Server ksql-server: image: 'confluentinc/cp-ksqldb-server:7.3.2' hostname: ksql-server container_name: ksql-server networks: - ksql-poc ports: - '8088:8088' depends_on: - kafka - schema-registry environment: KSQL_CONFIG_DIR: '/etc/ksql' KSQL_BOOTSTRAP_SERVERS: 'kafka:29092' KSQL_LISTENERS: 'http://0.0.0.0:8088' KSQL_KSQL_SERVICE_ID: 'ksql-examples-service-id' KSQL_KSQL_SCHEMA_REGISTRY_URL: 'http://schema-registry:8081' KSQL_KSQL_LOGGING_PROCESSING_TOPIC_AUTO_CREATE: 'true' KSQL_KSQL_LOGGING_PROCESSING_TOPIC_NAME: 'ksql_processing_log' KSQL_KSQL_LOGGING_PROCESSING_STREAM_AUTO_CREATE: 'true' KSQL_KSQL_LOGGING_PROCESSING_ROWS_INCLUDE: 'true' # KSQL_LOG4J_OPTS: '-Dlog4j.configuration=file:/etc/ksql/log4j-rolling.properties' # Runs the KSQL CLI ksql-cli: image: confluentinc/cp-ksqldb-cli:7.3.2 container_name: ksql-cli hostname: ksql-cli networks: - ksql-poc depends_on: - kafka - ksql-server entrypoint: /bin/sh tty: true networks: ksql-poc: driver: bridge ipam: driver: default
注册Schema的请求
curl -X POST -H "Content-Type: application/vnd.schemaregistry.v1+json" \ --data '{"schema": "{\"type\":\"record\",\"name\":\"Deployment\",\"namespace\":\"com.onrush.domain.event.schema\",\"fields\":[{\"name\":\"id\",\"type\":\"string\"},{\"name\":\"repo\",\"type\":\"string\"},{\"name\":\"name\",\"type\":\"string\"},{\"name\":\"teamId\",\"type\":\"string\"}]}"}' \ http://localhost:8081/subjects/deployment-value/versions
执行print 'deployment';的结果
Key format: ¯\_(ツ)_/¯ - no data processed Value format: KAFKA_STRING rowtime: 2023/07/23 10:07:13.993 Z, key: <null>, value: yh982hiuqsxho99Vgithub.com/xyz/content-sync-workflow*content-sync-workflowcnt, partition: 0 rowtime: 2023/07/23 10:07:14.155 Z, key: <null>, value: yh982hiuqsxho99Vgithub.com/xyz/content-sync-workflow*content-sync-workflowcnt, partition: 1 rowtime: 2023/07/23 10:07:14.164 Z, key: <null>, value: yh982hiuqsxho99Vgithub.com/xyz/content-sync-workflow*content-sync-workflowcnt, partition: 1
Go生产代码
avroData := map[string]interface{}{ "id": "yh982hiuqsxho99", "repo": "github.com/xyz/content-sync-workflow", "name": "content-sync-workflow", "teamId": "cnt", } binaryAvroData, err := codec.BinaryFromNative(nil, avroData) if err != nil { log.Fatal("Error encoding Avro data:", err) } message := &sarama.ProducerMessage{ Topic: topic, Value: sarama.ByteEncoder(binaryAvroData), }
核心问题
使用Go代码可以正常生产/消费该Avro格式消息,但在KSQL CLI执行select * from deployment;时出现反序列化错误,提示Unknown magic byte!,无法正常查询数据。
解决方案
问题根源
KSQL(以及Confluent生态中的Avro序列化/反序列化组件)要求Avro消息必须遵循Confluent Avro格式,该格式在原始Avro二进制数据前添加了:
- 1个字节的
magic byte(固定为0x0) - 4个字节的Schema ID(大端序)
而你的Go代码仅使用普通Avro编码器生成了原始Avro二进制数据,没有添加Confluent格式的头部信息,导致KSQL无法识别,触发Unknown magic byte!错误。从print 'deployment';的结果也能看出,KSQL将消息识别为KAFKA_STRING而非Avro,进一步验证了这一点。
修复步骤
1. 改用Confluent Avro编码器生产消息
在Go代码中,需要使用支持Confluent Avro格式的库(比如github.com/confluentinc/confluent-kafka-go/v2/kafka的AvroProducer)。示例代码如下:
package main import ( "log" "github.com/confluentinc/confluent-kafka-go/v2/kafka" "github.com/confluentinc/confluent-kafka-go/v2/schemaregistry" "github.com/confluentinc/confluent-kafka-go/v2/schemaregistry/avro" ) func main() { // 连接Schema Registry srClient, err := schemaregistry.NewClient(schemaregistry.NewConfig("http://localhost:8081")) if err != nil { log.Fatal(err) } // 创建Avro编码器 avroCodec, err := avro.NewEncoder(srClient, avro.NewEncoderConfig()) if err != nil { log.Fatal(err) } // 定义消息数据 avroData := map[string]interface{}{ "id": "yh982hiuqsxho99", "repo": "github.com/xyz/content-sync-workflow", "name": "content-sync-workflow", "teamId": "cnt", } // 编码为Confluent Avro格式 value, err := avroCodec.Encode("deployment-value", avroData) if err != nil { log.Fatal(err) } // 创建Kafka生产者 producer, err := kafka.NewProducer(&kafka.ConfigMap{"bootstrap.servers": "localhost:9092"}) if err != nil { log.Fatal(err) } defer producer.Close() // 生产消息 topic := "deployment" err = producer.Produce(&kafka.Message{ TopicPartition: kafka.TopicPartition{Topic: &topic, Partition: kafka.PartitionAny}, Value: value, }, nil) if err != nil { log.Fatal(err) } producer.Flush(15 * 1000) }
2. 验证修复效果
重新生产消息后,在KSQL CLI中执行以下操作验证:
- 执行
print 'deployment';,此时Value格式应显示为AVRO - 执行
select * from deployment;,应能正常查询到结构化的Avro数据
3. 可选:清理旧消息
如果topic中已有不符合格式的旧消息,可以删除topic并重新创建,避免KSQL处理旧消息时仍出现错误:
docker exec kafka kafka-topics --bootstrap-server kafka:29092 --delete --topic deployment
内容的提问来源于stack exchange,提问作者Ankit Kumar
相关产品推荐
相关产品推荐

