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

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. 1个字节的magic byte(固定为0x0)
  2. 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中执行以下操作验证:

  1. 执行print 'deployment';,此时Value格式应显示为AVRO
  2. 执行select * from deployment;,应能正常查询到结构化的Avro数据

3. 可选:清理旧消息

如果topic中已有不符合格式的旧消息,可以删除topic并重新创建,避免KSQL处理旧消息时仍出现错误:

docker exec kafka kafka-topics --bootstrap-server kafka:29092 --delete --topic deployment

内容的提问来源于stack exchange,提问作者Ankit Kumar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 03:42:02