将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)
问题原因
- 手动拼接的文件不符合标准Avro数据格式:
DataFileReader要求读取的是标准Avro数据文件,该文件包含固定魔法字节(Obj\x01)、文件元数据(含Schema)、同步标记及序列化记录数据。你手动写入0x0魔法字节+Schema字符串+占位符内容,完全不符合格式,导致魔法字节校验失败。 - 代码存在占位符错误:
out.write("${output file}")直接写入字符串占位符而非Kafka实际消息数据;FileUtils.writeByteArrayToFile的目标文件路径为空,生成的文件本身无效。 - 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
相关产品推荐
相关产品推荐

