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

将Kafka字节数据导入BigQuery时遇Avro文件无效异常求助

问题描述

尝试将Kafka主题中最多100条字节格式的记录导入BigQuery,已有avsc格式的Schema,执行以下步骤后读取Avro文件报错:

  • 使用Kafka控制台消费者消费100条消息并保存到文件
  • 编写代码生成包含魔法标记|Schema|记录的Avro文件
  • 编写测试工具读取Avro数据时抛出org.apache.avro.InvalidAvroMagicException: Not an Avro data file

生成Avro文件的代码

public class AvroSerializer {
  public static final byte MAGIC_BYTE = 0x0;

  public void serialize() throws Exception {
    Schema schema;
    ByteArrayOutputStream out = new ByteArrayOutputStream();
    try {
      schema =
          new Schema.Parser()
              .parse(
                  new File(
                      "${path to schema.avcs}"));
      byte[] kafakTopicData =
          FileUtils.readFileToByteArray(
              new File(
                  "${path to kafka topic dump using kafka console consumer}"));
      // MAGIC_BYTE | schemaId-bytes | avro_payload
      out.write(MAGIC_BYTE);
      out.write(schema.toString().getBytes());
      out.write("${output file}");
      FileUtils.writeByteArrayToFile(
          new File(
              ""),
          out.toByteArray());
    } catch (Exception ex) {
      throw new Exception(ex);
    }
  }
}

读取数据的代码

public void decryptAvro() {
    Schema schema = null;
    try {
      schema =
          new Schema.Parser()
              .parse(
                  new File(
                      "${path to schema.avsc}"));
      DatumReader<GenericRecord> datumReader = new GenericDatumReader<>(schema);
      DataFileReader<GenericRecord> dataFileReader =
          new DataFileReader<GenericRecord>(
              new File(
                  "${path to output file created in earlier step}"),
              datumReader);
      GenericRecord hcpClaims = null;

      while (dataFileReader.hasNext()) {
        hcpClaims = dataFileReader.next(hcpClaims);
        System.out.println(hcpClaims);
      }
    } catch (Exception e) {
      e.printStackTrace();
    }
  }

错误信息

org.apache.avro.InvalidAvroMagicException: Not an Avro data file.
    at org.apache.avro.file.DataFileStream.validateMagic(DataFileStream.java:115)
    at org.apache.avro.file.DataFileStream.initialize(DataFileStream.java:123)
    at org.apache.avro.file.DataFileReader.<init>(DataFileReader.java:143)
    at org.apache.avro.file.DataFileReader.<init>(DataFileReader.java:113)
    at com.optum.clm.avroutils.AvroReader.decryptAvro(AvroReader.java:22)
问题原因
  1. 手动拼接的文件不符合标准Avro数据格式:DataFileReader要求读取的是标准Avro数据文件,该文件包含固定魔法字节(Obj\x01)、文件元数据(含Schema)、同步标记及序列化记录数据。你手动写入0x0魔法字节+Schema字符串+占位符内容,完全不符合格式,导致魔法字节校验失败。
  2. 代码存在占位符错误:out.write("${output file}")直接写入字符串占位符而非Kafka实际消息数据;FileUtils.writeByteArrayToFile的目标文件路径为空,生成的文件本身无效。
  3. Kafka控制台消费的文件可能不是原始字节数据:默认Kafka控制台消费者会将消息转成字符串打印,而非保存原始字节流,导致读取的kafakTopicData不是正确的Avro序列化字节。
解决步骤

步骤1:正确获取Kafka原始字节消息

不要用默认控制台消费者,改用以下命令获取原始字节(避免字符串转换):

kafka-console-consumer.sh --bootstrap-server <kafka地址> --topic <主题名> --from-beginning --max-messages 100 --property print.value=false --property print.key=false --property value.deserializer=org.apache.kafka.common.serialization.ByteArrayDeserializer > raw_kafka_bytes.bin

注:若控制台输出包含额外字符,建议直接用Kafka API消费并保存原始字节,可靠性更高。

步骤2:用Avro官方API生成标准Avro文件

替换原AvroSerializer代码,使用DataFileWriter生成符合格式的Avro文件:

public class AvroSerializer {
    public void serialize() throws Exception {
        Schema schema = new Schema.Parser().parse(new File("${path to schema.avsc}"));
        File kafkaDataFile = new File("${path to raw kafka bytes file}");
        File outputAvroFile = new File("${path to output avro file}");

        List<byte[]> kafkaMessages = readRawKafkaMessages(kafkaDataFile);

        // 使用DataFileWriter写入标准Avro文件
        DatumWriter<GenericRecord> datumWriter = new GenericDatumWriter<>(schema);
        try (DataFileWriter<GenericRecord> dataFileWriter = new DataFileWriter<>(datumWriter)) {
            dataFileWriter.create(schema, outputAvroFile);
            GenericDatumReader<GenericRecord> datumReader = new GenericDatumReader<>(schema);
            Decoder decoder = DecoderFactory.get().binaryDecoder(new byte[0], null);
            for (byte[] messageBytes : kafkaMessages) {
                decoder.configure(messageBytes, 0, messageBytes.length);
                GenericRecord record = datumReader.read(null, decoder);
                dataFileWriter.append(record);
            }
        }
    }

    // 读取Kafka原始字节文件,假设换行是消息分隔符(需根据实际保存方式调整)
    private List<byte[]> readRawKafkaMessages(File file) throws IOException {
        List<byte[]> messages = new ArrayList<>();
        try (BufferedInputStream bis = new BufferedInputStream(new FileInputStream(file))) {
            ByteArrayOutputStream baos = new ByteArrayOutputStream();
            int b;
            while ((b = bis.read()) != -1) {
                if (b == '\n') {
                    messages.add(baos.toByteArray());
                    baos.reset();
                } else {
                    baos.write(b);
                }
            }
            if (baos.size() > 0) {
                messages.add(baos.toByteArray());
            }
        }
        return messages;
    }
}

步骤3:处理Confluent格式的Kafka消息(若适用)

如果Kafka消息是Confluent Avro序列化器生成的(格式为0x0魔法字节 + 4字节大端Schema ID + Avro payload),需先提取payload再写入文件,解析逻辑如下:

private byte[] extractConfluentAvroPayload(byte[] messageBytes) {
    if (messageBytes[0] != 0x0) {
        throw new IllegalArgumentException("Not a Confluent Avro message");
    }
    // 跳过魔法字节和4字节Schema ID,剩余部分为Avro payload
    return Arrays.copyOfRange(messageBytes, 5, messageBytes.length);
}
验证

运行修改后的序列化代码生成Avro文件,再用原读取代码测试,即可正常读取GenericRecord数据。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 07:07:20