Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 4 additions & 2 deletions c/src/main/java/org/apache/arrow/c/StructVectorLoader.java
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@
import org.apache.arrow.memory.BufferAllocator;
import org.apache.arrow.util.Collections2;
import org.apache.arrow.vector.BaseVariableWidthViewVector;
import org.apache.arrow.vector.ExtensionTypeVector;
import org.apache.arrow.vector.FieldVector;
import org.apache.arrow.vector.TypeLayout;
import org.apache.arrow.vector.complex.StructVector;
Expand Down Expand Up @@ -124,9 +125,10 @@ private void loadBuffers(
Iterator<Long> variadicBufferCounts) {
checkArgument(nodes.hasNext(), "no more field nodes for field %s and vector %s", field, vector);
ArrowFieldNode fieldNode = nodes.next();
// variadicBufferLayoutCount will be 0 for vectors of a type except BaseVariableWidthViewVector
FieldVector storageVector = ExtensionTypeVector.getStorageVector(vector);
// Only view storage has variadic buffers.
long variadicBufferLayoutCount = 0;
if (vector instanceof BaseVariableWidthViewVector) {
if (storageVector instanceof BaseVariableWidthViewVector) {
if (variadicBufferCounts.hasNext()) {
variadicBufferLayoutCount = variadicBufferCounts.next();
} else {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@
import java.util.List;
import org.apache.arrow.memory.ArrowBuf;
import org.apache.arrow.vector.BaseVariableWidthViewVector;
import org.apache.arrow.vector.ExtensionTypeVector;
import org.apache.arrow.vector.FieldVector;
import org.apache.arrow.vector.TypeLayout;
import org.apache.arrow.vector.complex.StructVector;
Expand Down Expand Up @@ -103,12 +104,13 @@ private void appendNodes(
List<Long> variadicBufferCounts) {
nodes.add(
new ArrowFieldNode(vector.getValueCount(), includeNullCount ? vector.getNullCount() : -1));
FieldVector storageVector = ExtensionTypeVector.getStorageVector(vector);
List<ArrowBuf> fieldBuffers = vector.getFieldBuffers();
long variadicBufferCount = getVariadicBufferCount(vector);
long variadicBufferCount = getVariadicBufferCount(storageVector);
int expectedBufferCount =
(int) (TypeLayout.getTypeBufferCount(vector.getField().getType()) + variadicBufferCount);
// only update variadicBufferCounts for vectors that have variadic buffers
if (vector instanceof BaseVariableWidthViewVector) {
if (storageVector instanceof BaseVariableWidthViewVector) {
variadicBufferCounts.add(variadicBufferCount);
}
if (fieldBuffers.size() != expectedBufferCount) {
Expand Down
70 changes: 70 additions & 0 deletions c/src/test/java/org/apache/arrow/c/RoundtripTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@
import org.apache.arrow.vector.DateMilliVector;
import org.apache.arrow.vector.DecimalVector;
import org.apache.arrow.vector.DurationVector;
import org.apache.arrow.vector.ExtensionTypeVector;
import org.apache.arrow.vector.FieldVector;
import org.apache.arrow.vector.FixedSizeBinaryVector;
import org.apache.arrow.vector.Float2Vector;
Expand Down Expand Up @@ -76,6 +77,7 @@
import org.apache.arrow.vector.ValueVector;
import org.apache.arrow.vector.VarBinaryVector;
import org.apache.arrow.vector.VarCharVector;
import org.apache.arrow.vector.VariableWidthFieldVector;
import org.apache.arrow.vector.VectorSchemaRoot;
import org.apache.arrow.vector.ViewVarBinaryVector;
import org.apache.arrow.vector.ViewVarCharVector;
Expand All @@ -91,6 +93,8 @@
import org.apache.arrow.vector.complex.StructVector;
import org.apache.arrow.vector.complex.UnionVector;
import org.apache.arrow.vector.complex.impl.UnionMapWriter;
import org.apache.arrow.vector.extension.JsonType;
import org.apache.arrow.vector.extension.OpaqueType;
import org.apache.arrow.vector.extension.UuidType;
import org.apache.arrow.vector.holders.IntervalDayHolder;
import org.apache.arrow.vector.holders.NullableLargeVarBinaryHolder;
Expand All @@ -100,13 +104,17 @@
import org.apache.arrow.vector.types.Types.MinorType;
import org.apache.arrow.vector.types.pojo.ArrowType;
import org.apache.arrow.vector.types.pojo.ArrowType.ExtensionType;
import org.apache.arrow.vector.types.pojo.ExtensionTypeRegistry;
import org.apache.arrow.vector.types.pojo.Field;
import org.apache.arrow.vector.types.pojo.FieldType;
import org.apache.arrow.vector.types.pojo.Schema;
import org.apache.arrow.vector.util.TransferPair;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.Arguments;
import org.junit.jupiter.params.provider.MethodSource;

public class RoundtripTest {
private static final String EMPTY_SCHEMA_PATH = "";
Expand Down Expand Up @@ -848,6 +856,68 @@ public void testExtensionTypeVector() {
}
}

static Stream<Arguments> viewExtensionCases() {
byte[][] inlined = {
"{}".getBytes(StandardCharsets.UTF_8), null, "null".getBytes(StandardCharsets.UTF_8)
};
byte[][] external = {
"{\"message\":\"a JSON value longer than twelve bytes\"}".getBytes(StandardCharsets.UTF_8),
null
};
return Stream.of(
new JsonType(ArrowType.Utf8View.INSTANCE),
new OpaqueType(ArrowType.Utf8View.INSTANCE, "json", "test"),
new OpaqueType(ArrowType.BinaryView.INSTANCE, "binary", "test"))
.flatMap(
type ->
Stream.of(
Arguments.of(type, new byte[0][], 3),
Arguments.of(type, inlined, 3),
Arguments.of(type, external, 4)));
}

@ParameterizedTest
@MethodSource("viewExtensionCases")
public void testViewExtensionTypeVector(ExtensionType type, byte[][] values, int bufferCount) {
ExtensionType previous = ExtensionTypeRegistry.lookup(type.extensionName());
ExtensionTypeRegistry.register(type);
Schema schema = new Schema(Collections.singletonList(Field.nullable("a", type)));
try (VectorSchemaRoot root = VectorSchemaRoot.create(schema, allocator)) {
ExtensionTypeVector<?> vector = (ExtensionTypeVector<?>) root.getVector("a");
setVector((VariableWidthFieldVector) vector.getUnderlyingVector(), values);
root.setRowCount(values.length);

try (ArrowSchema arrowSchema = ArrowSchema.allocateNew(allocator);
ArrowArray arrowArray = ArrowArray.allocateNew(allocator)) {
Data.exportVector(allocator, vector, null, arrowArray, arrowSchema);
assertEquals(bufferCount, arrowArray.snapshot().n_buffers);
try (FieldVector imported =
Data.importVector(childAllocator, arrowArray, arrowSchema, null)) {
assertViewExtensionVector(vector, imported);
}
}

try (VectorSchemaRoot imported = vectorSchemaRootRoundtrip(root)) {
assertEquals(root.getSchema(), imported.getSchema());
assertEquals(root.getRowCount(), imported.getRowCount());
assertViewExtensionVector(vector, imported.getVector("a"));
}
} finally {
ExtensionTypeRegistry.unregister(type);
if (previous != null) {
ExtensionTypeRegistry.register(previous);
}
}
}

private void assertViewExtensionVector(ExtensionTypeVector<?> expected, FieldVector actual) {
assertEquals(expected.getField(), actual.getField());
ExtensionTypeVector<?> imported = assertInstanceOf(ExtensionTypeVector.class, actual);
assertTrue(
VectorEqualsVisitor.vectorEquals(
expected.getUnderlyingVector(), imported.getUnderlyingVector()));
}

@Test
public void testVectorSchemaRoot() {
VectorSchemaRoot imported;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -76,6 +76,24 @@ public T getUnderlyingVector() {
return underlyingVector;
}

/** Get the storage vector, unwrapping extension vectors. */
public static FieldVector getStorageVector(FieldVector vector) {
while (vector instanceof ExtensionTypeVector) {
vector = ((ExtensionTypeVector<?>) vector).getUnderlyingVector();
}
return vector;
}

@Override
public int getExportedCDataBufferCount() {
return this.underlyingVector.getExportedCDataBufferCount();
}

@Override
public void exportCDataBuffers(List<ArrowBuf> buffers, ArrowBuf buffersPtr, long nullValue) {
this.underlyingVector.exportCDataBuffers(buffers, buffersPtr, nullValue);
}

@Override
public void allocateNew() throws OutOfMemoryException {
this.underlyingVector.allocateNew();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -106,11 +106,12 @@ private void loadBuffers(
Iterator<ArrowFieldNode> nodes,
CompressionCodec codec,
Iterator<Long> variadicBufferCounts) {
FieldVector storageVector = ExtensionTypeVector.getStorageVector(vector);
checkArgument(nodes.hasNext(), "no more field nodes for field %s and vector %s", field, vector);
ArrowFieldNode fieldNode = nodes.next();
// variadicBufferLayoutCount will be 0 for vectors of a type except BaseVariableWidthViewVector
// Only view storage has variadic buffers.
long variadicBufferLayoutCount = 0;
if (vector instanceof BaseVariableWidthViewVector) {
if (storageVector instanceof BaseVariableWidthViewVector) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This now expects a variadicBufferCounts entry for extension vectors whose storage is a view vector, but c/StructVectorUnloader (which feeds this loader in Data.importIntoVectorSchemaRoot) still checks vector instanceof BaseVariableWidthViewVector on the wrapper, so it never emits one.

Importing a C Data batch with an arrow.json or OpaqueType column on Utf8View/BinaryView storage now fails with IlllegalStateException: No variadicBufferCounts available for BaseVariableWidthViewVector when all values are inlined (<= 12 bytes) or the batch is empty. The same OpaqueVector(Utf8View) case loads on main, so this is a regression for existing extension types, not only for the new one. It also affects ArrowArrayStreamReader.loadNextBatch and the dataset NativeScanner.

Could you apply the same unwrapping in StructVectorUnloader, and add a C Data round-trip test for an extension type on view storage? RoundtripTest.testExtensionTypeVector only covers UuidType.

if (variadicBufferCounts.hasNext()) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Same problem in the export direction: this now adds a variadic count for extension vectors on view storage, but c/StructVectorLoader.loadBuffers still checks the wrapper and never consumes it.

Data.exportVectorSchemaRoot (and ArrayStreamExporter) on a root with a JsonVector(Utf8View) column then throws IllegalArgumentException: not all nodes, buffers and variadicBufferCounts were consumed. It fails even with zero variadic buffers, a case that succeeds on main.

StructVectorLoader needs the same change as VectorLoader. Since all four loader/unloader classes have to agree on which vectors carry a variadic count, a single shared helper for the unwrap (for example a static storage-vector accessor on ExtensionTypeVector) would keep them in sync.

variadicBufferLayoutCount = variadicBufferCounts.next();
} else {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -104,14 +104,15 @@ private void appendNodes(
List<ArrowFieldNode> nodes,
List<ArrowBuf> buffers,
List<Long> variadicBufferCounts) {
FieldVector storageVector = ExtensionTypeVector.getStorageVector(vector);
nodes.add(
new ArrowFieldNode(vector.getValueCount(), includeNullCount ? vector.getNullCount() : -1));
List<ArrowBuf> fieldBuffers = vector.getFieldBuffers();
long variadicBufferCount = getVariadicBufferCount(vector);
long variadicBufferCount = getVariadicBufferCount(storageVector);
int expectedBufferCount =
(int) (TypeLayout.getTypeBufferCount(vector.getField().getType()) + variadicBufferCount);
// only update variadicBufferCounts for vectors that have variadic buffers
if (vector instanceof BaseVariableWidthViewVector) {
if (storageVector instanceof BaseVariableWidthViewVector) {
variadicBufferCounts.add(variadicBufferCount);
}
if (fieldBuffers.size() != expectedBufferCount) {
Expand Down
132 changes: 132 additions & 0 deletions vector/src/main/java/org/apache/arrow/vector/extension/JsonType.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,132 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.apache.arrow.vector.extension;

import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.DeserializationFeature;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import java.util.Collections;
import java.util.Objects;
import org.apache.arrow.memory.BufferAllocator;
import org.apache.arrow.vector.FieldVector;
import org.apache.arrow.vector.types.pojo.ArrowType;
import org.apache.arrow.vector.types.pojo.ExtensionTypeRegistry;
import org.apache.arrow.vector.types.pojo.Field;
import org.apache.arrow.vector.types.pojo.FieldType;

/**
* Canonical extension type for UTF-8 encoded RFC 8259 JSON values.
*
* <p>The storage type is {@link ArrowType.Utf8}, {@link ArrowType.LargeUtf8}, or {@link
* ArrowType.Utf8View}. Values use the corresponding string vector; this type does not parse or
* validate individual JSON values.
*
* <p>Register the type before reading schemas containing {@code arrow.json}:
*
* <pre>{@code
* JsonType.ensureRegistered();
* Field field = Field.nullable("json", new JsonType(ArrowType.Utf8.INSTANCE));
* try (JsonVector vector = (JsonVector) field.createVector(allocator)) {
* VarCharVector storage = (VarCharVector) vector.getUnderlyingVector();
* storage.setSafe(0, "{}".getBytes(java.nio.charset.StandardCharsets.UTF_8));
* vector.setValueCount(1);
* Text value = vector.getObject(0);
* }
* }</pre>
*/
public class JsonType extends ArrowType.ExtensionType {
public static final String EXTENSION_NAME = "arrow.json";
private static final ObjectMapper MAPPER =
new ObjectMapper().enable(DeserializationFeature.FAIL_ON_TRAILING_TOKENS);
private final ArrowType storageType;

/** Register a prototype that can deserialize all supported JSON storage types. */
public static void ensureRegistered() {
ExtensionTypeRegistry.register(new JsonType(ArrowType.Utf8.INSTANCE));
}

/**
* Create a JSON type backed by the specified string type.
*
* @param storageType Utf8, LargeUtf8, or Utf8View
* @throws IllegalArgumentException if the storage type is not a supported string type
*/
public JsonType(ArrowType storageType) {
Objects.requireNonNull(storageType, "storageType");
if (!(storageType instanceof ArrowType.Utf8)
&& !(storageType instanceof ArrowType.LargeUtf8)
&& !(storageType instanceof ArrowType.Utf8View)) {
throw new IllegalArgumentException(
"arrow.json requires Utf8, LargeUtf8, or Utf8View storage, got " + storageType);
}
this.storageType = storageType;
}

@Override
public ArrowType storageType() {
return storageType;
}

@Override
public String extensionName() {
return EXTENSION_NAME;
}

@Override
public boolean extensionEquals(ExtensionType other) {
return other instanceof JsonType && storageType.equals(other.storageType());
}

@Override
public String serialize() {
return "";
}

@Override
public ArrowType deserialize(ArrowType storageType, String serializedData) {
JsonType type = new JsonType(storageType);
if (serializedData == null) {
throw new InvalidExtensionMetadataException("arrow.json metadata must not be null");
}
if (!serializedData.isEmpty()) {
try {
JsonNode metadata = MAPPER.readTree(serializedData);
if (metadata == null || !metadata.isObject()) {
throw new InvalidExtensionMetadataException("arrow.json metadata must be a JSON object");
}
} catch (JsonProcessingException e) {
throw new InvalidExtensionMetadataException("arrow.json metadata is invalid", e);
}
}
return type;
}

@Override
public boolean isComplex() {
return false;
}

@Override
public FieldVector getNewVector(String name, FieldType fieldType, BufferAllocator allocator) {
Field field = new Field(name, fieldType, Collections.emptyList());
FieldType storageFieldType =
new FieldType(fieldType.isNullable(), storageType, fieldType.getDictionary(), null);
FieldVector storage = storageFieldType.createNewSingleVector(name, allocator, null);
return new JsonVector(field, allocator, storage);
}
}
Loading
Loading