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

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中没有值?正确的实现方式是什么?


问题分析与解决方案

问题原因

  1. ListWriter生命周期管理错误:每个列表元素的startList()必须对应endList(),原代码只调用了startList()但未在循环中闭合,导致数据无法持久化到向量。
  2. 字段类型不匹配:定义中AttributeTag是FixedSizeBinary(2)、AttributeType是FixedSizeBinary(1),但代码使用varChar()写入,类型不匹配会导致数据被静默丢弃。
  3. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 11:02:04