Apache Arrow写入Map字段后读取报错:not all nodes and buffers were consumed
Apache Arrow Java:写入Map字段后读取报"not all nodes and buffers were consumed"错误
问题描述
使用Apache Arrow的Java代码将Map字段写入文件,读取阶段加载Record Batch时抛出not all nodes and buffers were consumed错误。写入阶段调试确认VectorSchemaRoot数据正确,但读取报错。
相关代码
File file = new File("/temp/test.arrow"); Schema schema; try (BufferAllocator allocator = new RootAllocator()) { Field keyField = new Field("id", FieldType.notNullable(new ArrowType.Int(64, true)), null); Field valueField = new Field("value", FieldType.nullable(new ArrowType.Int(64, true)), null); Field structField = new Field("entry", FieldType.notNullable(ArrowType.Struct.INSTANCE), List.of(keyField, valueField)); Field mapIntToIntField = new Field("mapFieldIntToInt", FieldType.notNullable(new ArrowType.Map(false)), List.of(structField)); schema = new Schema(Arrays.asList(mapIntToIntField)); try ( VectorSchemaRoot vectorSchemaRoot = VectorSchemaRoot.create(schemaPerson, allocator); MapVector mapVector = (MapVector) vectorSchemaRoot.getVector("mapFieldIntToInt")) { UnionMapWriter mapWriter = mapVector.getWriter(); mapWriter.setPosition(0); mapWriter.startMap(); for (int i = 0; i < 3; i++) { mapWriter.startEntry(); mapWriter.key().bigInt().writeBigInt(i); mapWriter.value().bigInt().writeBigInt(i * 7); mapWriter.endEntry(); } mapWriter.endMap(); mapWriter.setValueCount(1); vectorSchemaRoot.setRowCount(1); try ( FileOutputStream fileOutputStream = new FileOutputStream(file); ArrowFileWriter writer = new ArrowFileWriter(vectorSchemaRoot, null, fileOutputStream.getChannel())) { writer.start(); writer.writeBatch(); writer.end(); } catch (IOException e) { e.printStackTrace(); } } } // Deserialize Arrow data from a file try ( BufferAllocator rootAllocator = new RootAllocator(); FileInputStream fileInputStream = new FileInputStream(file); ArrowFileReader reader = new ArrowFileReader(fileInputStream.getChannel(), rootAllocator)) { for (ArrowBlock arrowBlock : reader.getRecordBlocks()) { reader.loadRecordBatch(arrowBlock); // error thrown here VectorSchemaRoot vectorSchemaRootRecover = reader.getVectorSchemaRoot(); System.out.print(vectorSchemaRootRecover.contentToTSVString()); int totalCount = vectorSchemaRootRecover.getRowCount(); Schema schema = vectorSchemaRootRecover.getSchema(); for (int i = 0; i < totalCount; i++) { for (Field field : schema.getFields()) { String fieldName = field.getName(); Object arrowValue = vectorSchemaRootRecover.getVector(field).getObject(i); System.out.println("fieldName: " + fieldName + ", arrowValue: " + arrowValue); } } } } catch (IOException e) { e.printStackTrace(); }
解决提示
- 修正变量名错误:创建
VectorSchemaRoot时误用了未定义的schemaPerson,需改为之前创建的schema。这个错误会导致写入的Schema与实际Vector不匹配,引发读取解析失败。 - 符合Arrow Map类型Schema规范:Arrow的Map类型要求子字段必须命名为
entries(而非代码中的entry),这是强制规范。修改structField的名称:Field structField = new Field("entries", FieldType.notNullable(ArrowType.Struct.INSTANCE), List.of(keyField, valueField)); - 读取时重置VectorSchemaRoot:每次处理完一批数据后,调用
vectorSchemaRootRecover.clear(),避免残留数据干扰下一次解析。 - 使用更可靠的读取方式:替代手动遍历
RecordBlocks,改用reader.loadNextBatch()简化代码,避免手动加载的潜在错误:VectorSchemaRoot vectorSchemaRootRecover = reader.getVectorSchemaRoot(); while (reader.loadNextBatch()) { System.out.print(vectorSchemaRootRecover.contentToTSVString()); // 其他数据处理逻辑 vectorSchemaRootRecover.clear(); }
内容的提问来源于stack exchange,提问作者user2268908
相关产品推荐
相关产品推荐

