Java Apache Arrow中Struct列表类型Field填充失败,求正确实现
问题
我有一个以二维表为主的数据集,其中名为Attributes的列每个单元格包含一个Struct类型的List。每个Struct包含三个字段:Attribute Tag、Attribute Type和Attribute Value。
Attributes字段的定义如下:
/** * Attribute Tag - Two character tag. */ public static final Field ATTRIBUTE_TAG_FIELD = new Field("AttributeTag", FieldType.notNullable(new ArrowType.FixedSizeBinary(2)), null); /** * Attribute Type - One character type. */ public static final Field ATTRIBUTE_TYPE_FIELD = new Field( "AttributeType", new FieldType(false, new ArrowType.FixedSizeBinary(1), null), null ); /** * String representation of the Attribute value. */ public static final Field ATTRIBUTE_VALUE_FIELD = new Field("AttributeValue", FieldType.notNullable(new ArrowType.Utf8()), null); /** * The field is a nullable List of Structs each with an attribute tag, type and value. */ public static final Field ATTRIBUTES_FIELD = new Field("Attributes", FieldType.nullable(new ArrowType.List()), List.of( new Field("Attribute", FieldType.nullable(new ArrowType.Struct()), List.of( ATTRIBUTE_TAG_FIELD, ATTRIBUTE_TYPE_FIELD, ATTRIBUTE_VALUE_FIELD))));
我编写了一段代码尝试从源数据填充Attributes,代码运行无报错,但attributes ListVector中没有任何值。代码如下:
final ListVector attributes = (ListVector) ATTRIBUTES_FIELD.createVector(allocator); // 这是要填充到attributes向量中的属性源数据 final List<SAMRecord.SAMTagAndValue> recordAttributes = samRecord.getAttributes(); if (recordAttributes != null && recordAttributes.size() > 0 ) { final UnionListWriter listWriter = attributes.getWriter(); listWriter.allocate(); IntStream.range(0, recordAttributes.size()).forEachOrdered(attributeIndex -> { listWriter.setPosition(attributeIndex); listWriter.startList(); // 将属性值放入Arrow结构体 final SAMRecord.SAMTagAndValue samTagAndValue = recordAttributes.get(attributeIndex); // 我认为问题出在这里,调试时似乎创建了与向量无关的新写入器? final BaseWriter.StructWriter structWriter = listWriter.struct("Attribute"); structWriter.start(); final byte[] tagBytes = samTagAndValue.tag.getBytes(StandardCharsets.UTF_8); // todo 从值中获取类型 final byte[] typeBytes = "S".getBytes(StandardCharsets.UTF_8); final byte[] valueBytes = samTagAndValue.value.toString().getBytes(StandardCharsets.UTF_8); ArrowBuf tempBuf = allocator.buffer(tagBytes.length); tempBuf.setBytes(0, tagBytes); structWriter.varChar("AttributeTag").writeVarChar(0, tagBytes.length, tempBuf); tempBuf.close(); tempBuf = allocator.buffer(typeBytes.length); structWriter.varChar("AttributeType").writeVarChar(0, typeBytes.length, tempBuf); tempBuf.close(); tempBuf = allocator.buffer(valueBytes.length); structWriter.varChar("AttributeValue").writeVarChar(0, valueBytes.length, tempBuf); tempBuf.close(); structWriter.end(); }); listWriter.setValueCount(recordAttributes.size()); listWriter.end(); }
请问为何attributes ListVector中没有值?正确的实现方式是什么?
问题分析与解决方案
问题原因
- ListWriter生命周期管理错误:每个列表元素的
startList()必须对应endList(),原代码只调用了startList()但未在循环中闭合,导致数据无法持久化到向量。 - 字段类型不匹配:定义中
AttributeTag是FixedSizeBinary(2)、AttributeType是FixedSizeBinary(1),但代码使用varChar()写入,类型不匹配会导致数据被静默丢弃。 - StructWriter使用有误:未正确关联当前列表位置的元素,写入操作未映射到目标向量的内部结构。
正确实现代码
final ListVector attributes = (ListVector) ATTRIBUTES_FIELD.createVector(allocator); final List<SAMRecord.SAMTagAndValue> recordAttributes = samRecord.getAttributes(); if (recordAttributes != null && !recordAttributes.isEmpty()) { try (UnionListWriter listWriter = attributes.getWriter()) { // 预分配对应数量的空间 listWriter.allocate(recordAttributes.size()); listWriter.setValueCount(recordAttributes.size()); for (int attributeIndex = 0; attributeIndex < recordAttributes.size(); attributeIndex++) { listWriter.setPosition(attributeIndex); // 明确每个列表包含1个Struct元素 listWriter.startList(1); SAMRecord.SAMTagAndValue samTagAndValue = recordAttributes.get(attributeIndex); BaseWriter.StructWriter structWriter = listWriter.struct("Attribute"); structWriter.start(); // 处理FixedSizeBinary类型的AttributeTag byte[] tagBytes = samTagAndValue.tag.getBytes(StandardCharsets.UTF_8); // 保证字节长度符合FixedSizeBinary(2)要求 byte[] fixedTag = new byte[2]; System.arraycopy(tagBytes, 0, fixedTag, 0, Math.min(tagBytes.length, 2)); try (ArrowBuf tempBuf = allocator.buffer(2)) { tempBuf.setBytes(0, fixedTag); structWriter.fixedSizeBinary("AttributeTag").writeFixedSizeBinary(tempBuf); } // 处理FixedSizeBinary类型的AttributeType byte[] typeBytes = "S".getBytes(StandardCharsets.UTF_8); byte[] fixedType = new byte[1]; fixedType[0] = typeBytes.length > 0 ? typeBytes[0] : 'S'; try (ArrowBuf tempBuf = allocator.buffer(1)) { tempBuf.setBytes(0, fixedType); structWriter.fixedSizeBinary("AttributeType").writeFixedSizeBinary(tempBuf); } // 处理Utf8类型的AttributeValue String valueStr = samTagAndValue.value.toString(); byte[] valueBytes = valueStr.getBytes(StandardCharsets.UTF_8); try (ArrowBuf tempBuf = allocator.buffer(valueBytes.length)) { tempBuf.setBytes(0, valueBytes); structWriter.varChar("AttributeValue").writeVarChar(0, valueBytes.length, tempBuf); } structWriter.end(); // 闭合当前列表元素 listWriter.endList(); } // 刷新数据到向量 listWriter.flush(); } }
关键改进点
- 列表生命周期管理:每个
startList(1)对应endList(),明确每个列表包含的元素数量。 - 类型匹配:针对
FixedSizeBinary类型使用专属写入器,避免类型不匹配导致的数据丢失。 - 资源安全:用try-with-resources自动管理
ArrowBuf和UnionListWriter,避免内存泄漏。 - 数据合法性校验:确保写入的字节长度符合
FixedSizeBinary的定义要求。
内容的提问来源于stack exchange,提问作者Mark
相关产品推荐
相关产品推荐

