通过Kinesis流传输Avro数据时遇ByteBuffer.allocate异常求助
解决Kinesis传输Avro数据时的IllegalArgumentException问题
你这明显是踩了Avro序列化的常见坑:用生成Avro完整文件的DataFileWriter来序列化单条记录,但Kinesis每条记录只需要独立的Avro二进制数据,不是带文件头、元数据的完整Avro文件格式。
问题根源
DataFileWriter的设计目标是生成可持久化的Avro文件,会自动给序列化结果加上文件头部(包含schema、元数据、同步标记等)。但Kinesis传输的是单条业务记录,当消费者用DataFileReader去读取这种被截断的"文件片段"时,就会因为无法解析完整的文件结构抛出异常。
修复步骤
1. 修改生产者的序列化逻辑
把DataFileWriter换成直接用BinaryEncoder序列化单条记录:
private static byte[] toAvro(User user) throws IOException { SpecificDatumWriter<User> datumWriter = new SpecificDatumWriter<>(User.class); try (ByteArrayOutputStream baos = new ByteArrayOutputStream()) { BinaryEncoder encoder = EncoderFactory.get().binaryEncoder(baos, null); datumWriter.write(user, encoder); encoder.flush(); baos.flush(); return baos.toByteArray(); } }
这里直接指定User类型给SpecificDatumWriter,比泛型的GenericContainer更高效且类型安全。
2. 修改消费者的反序列化逻辑
对应地,用BinaryDecoder反序列化单条记录,抛弃DataFileReader:
private User convert(ByteBuffer data) { try { SpecificDatumReader<User> datumReader = new SpecificDatumReader<>(User.class); BinaryDecoder decoder = DecoderFactory.get().binaryDecoder(data.array(), null); return datumReader.read(null, decoder); } catch (IOException e) { throw new RuntimeException("反序列化Avro记录失败", e); } }
同时调整processRecords方法,直接处理ByteBuffer避免编码转换问题:
@Override public void processRecords(List<Record> records, IRecordProcessorCheckpointer checkpointer) { System.out.println("Processing " + records.size() + " records"); List<User> data = records.stream() .map(Record::getData) .map(this::convert) .collect(Collectors.toList()); // 这里添加你的业务处理逻辑 }
额外说明
- 如果你确实需要在Kinesis中传输批量Avro数据,那得保证每条Kinesis记录是完整的Avro文件,但这种场景很少见——Kinesis本身更适合单条记录的流式传输。
- 之前本地测试正常是因为你读写的是完整的Avro文件,而通过Kinesis传输时,每条记录是错误的"文件片段",自然就触发了异常。
内容的提问来源于stack exchange,提问作者idmitriev
相关产品推荐
相关产品推荐

