This is an automated email from the ASF dual-hosted git repository.

caicancai pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/calcite.git


The following commit(s) were added to refs/heads/main by this push:
     new 9a03982534 [CALCITE-6512] Support Arrow List type
9a03982534 is described below

commit 9a0398253453b1e8c9bd398ef334859508866ef5
Author: Cancai Cai <[email protected]>
AuthorDate: Thu Jun 4 10:33:35 2026 +0800

    [CALCITE-6512] Support Arrow List type
---
 .../adapter/arrow/ArrowDirectEnumerator.java       | 65 ++++++++++++++++++++++
 .../calcite/adapter/arrow/ArrowEnumerable.java     |  3 +-
 .../adapter/arrow/ArrowFieldTypeFactory.java       | 18 ++++--
 .../apache/calcite/adapter/arrow/ArrowTable.java   | 46 ++++++++++-----
 .../adapter/arrow/ArrowAdapterDataTypesTest.java   | 18 ++++++
 .../calcite/adapter/arrow/ArrowDataTest.java       | 57 +++++++++++++++++++
 6 files changed, 187 insertions(+), 20 deletions(-)

diff --git 
a/arrow/src/main/java/org/apache/calcite/adapter/arrow/ArrowDirectEnumerator.java
 
b/arrow/src/main/java/org/apache/calcite/adapter/arrow/ArrowDirectEnumerator.java
new file mode 100644
index 0000000000..2ab896f09c
--- /dev/null
+++ 
b/arrow/src/main/java/org/apache/calcite/adapter/arrow/ArrowDirectEnumerator.java
@@ -0,0 +1,65 @@
+/*
+ * 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.calcite.adapter.arrow;
+
+import org.apache.calcite.util.ImmutableIntList;
+import org.apache.calcite.util.Util;
+
+import org.apache.arrow.vector.ipc.ArrowFileReader;
+import org.apache.arrow.vector.ipc.message.ArrowRecordBatch;
+
+import java.io.IOException;
+
+/**
+ * Enumerator that reads projected Arrow value-vectors directly.
+ */
+class ArrowDirectEnumerator extends AbstractArrowEnumerator {
+  private final Runnable onClose;
+
+  ArrowDirectEnumerator(ArrowFileReader arrowFileReader, ImmutableIntList 
fields,
+      Runnable onClose) {
+    super(arrowFileReader, fields);
+    this.onClose = onClose;
+  }
+
+  @Override protected void evaluateOperator(ArrowRecordBatch arrowRecordBatch) 
{
+  }
+
+  @Override public boolean moveNext() {
+    if (currRowIndex >= rowCount - 1) {
+      final boolean hasNextBatch;
+      try {
+        hasNextBatch = arrowFileReader.loadNextBatch();
+      } catch (IOException e) {
+        throw Util.toUnchecked(e);
+      }
+      if (hasNextBatch) {
+        currRowIndex = 0;
+        this.valueVectors.clear();
+        loadNextArrowBatch();
+      }
+      return hasNextBatch;
+    } else {
+      currRowIndex++;
+      return true;
+    }
+  }
+
+  @Override public void close() {
+    onClose.run();
+  }
+}
diff --git 
a/arrow/src/main/java/org/apache/calcite/adapter/arrow/ArrowEnumerable.java 
b/arrow/src/main/java/org/apache/calcite/adapter/arrow/ArrowEnumerable.java
index 516822567e..735c75c8ed 100644
--- a/arrow/src/main/java/org/apache/calcite/adapter/arrow/ArrowEnumerable.java
+++ b/arrow/src/main/java/org/apache/calcite/adapter/arrow/ArrowEnumerable.java
@@ -56,8 +56,7 @@ class ArrowEnumerable extends AbstractEnumerable<Object> {
         return new ArrowFilterEnumerator(arrowFileReader, fields, filter,
             onClose);
       }
-      throw new IllegalArgumentException(
-          "The arrow enumerator must have either a filter or a projection");
+      return new ArrowDirectEnumerator(arrowFileReader, fields, onClose);
     } catch (Exception e) {
       throw Util.toUnchecked(e);
     }
diff --git 
a/arrow/src/main/java/org/apache/calcite/adapter/arrow/ArrowFieldTypeFactory.java
 
b/arrow/src/main/java/org/apache/calcite/adapter/arrow/ArrowFieldTypeFactory.java
index 1693637e8c..30c738bece 100644
--- 
a/arrow/src/main/java/org/apache/calcite/adapter/arrow/ArrowFieldTypeFactory.java
+++ 
b/arrow/src/main/java/org/apache/calcite/adapter/arrow/ArrowFieldTypeFactory.java
@@ -22,6 +22,7 @@
 
 import org.apache.arrow.vector.types.FloatingPointPrecision;
 import org.apache.arrow.vector.types.pojo.ArrowType;
+import org.apache.arrow.vector.types.pojo.Field;
 
 /**
  * Arrow field type.
@@ -32,19 +33,20 @@ private ArrowFieldTypeFactory() {
     throw new UnsupportedOperationException("Utility class");
   }
 
-  public static RelDataType toType(ArrowType arrowType, JavaTypeFactory 
typeFactory) {
-    RelDataType sqlType = of(arrowType, typeFactory);
+  public static RelDataType toType(Field field, JavaTypeFactory typeFactory) {
+    RelDataType sqlType = of(field, typeFactory);
     return typeFactory.createTypeWithNullability(sqlType, true);
   }
 
   /**
-   * Converts an Arrow type to a Calcite RelDataType.
+   * Converts an Arrow field to a Calcite RelDataType.
    *
-   * @param arrowType the Arrow type to convert
+   * @param field the Arrow field to convert
    * @param typeFactory the factory to create the Calcite type
    * @return the corresponding Calcite RelDataType
    */
-  private static RelDataType of(ArrowType arrowType, JavaTypeFactory 
typeFactory) {
+  private static RelDataType of(Field field, JavaTypeFactory typeFactory) {
+    ArrowType arrowType = field.getType();
     switch (arrowType.getTypeID()) {
     case Int:
       int bitWidth = ((ArrowType.Int) arrowType).getBitWidth();
@@ -82,6 +84,12 @@ private static RelDataType of(ArrowType arrowType, 
JavaTypeFactory typeFactory)
           ((ArrowType.Decimal) arrowType).getScale());
     case Time:
       return typeFactory.createSqlType(SqlTypeName.TIME);
+    case List:
+      if (field.getChildren().size() != 1) {
+        throw new IllegalArgumentException("Arrow List type must have one 
child field: " + field);
+      }
+      RelDataType elementType = toType(field.getChildren().get(0), 
typeFactory);
+      return typeFactory.createArrayType(elementType, -1);
     case Timestamp:
       ArrowType.Timestamp timestampType = (ArrowType.Timestamp) arrowType;
       int timestampPrecision;
diff --git 
a/arrow/src/main/java/org/apache/calcite/adapter/arrow/ArrowTable.java 
b/arrow/src/main/java/org/apache/calcite/adapter/arrow/ArrowTable.java
index ba1568bcd2..5afb74e51d 100644
--- a/arrow/src/main/java/org/apache/calcite/adapter/arrow/ArrowTable.java
+++ b/arrow/src/main/java/org/apache/calcite/adapter/arrow/ArrowTable.java
@@ -121,18 +121,7 @@ public Enumerable<Object> query(DataContext root, 
ImmutableIntList fields,
 
     if (conditions.isEmpty()) {
       filter = null;
-
-      final List<ExpressionTree> expressionTrees = new ArrayList<>();
-      for (int fieldOrdinal : fields) {
-        Field field = schema.getFields().get(fieldOrdinal);
-        TreeNode node = TreeBuilder.makeField(field);
-        expressionTrees.add(TreeBuilder.makeExpression(node, field));
-      }
-      try {
-        projector = Projector.make(schema, expressionTrees);
-      } catch (GandivaException e) {
-        throw Util.toUnchecked(e);
-      }
+      projector = makeProjector(fields);
     } else {
       projector = null;
 
@@ -208,11 +197,42 @@ private static RelDataType deduceRowType(Schema schema,
     final RelDataTypeFactory.Builder builder = typeFactory.builder();
     for (Field field : schema.getFields()) {
       builder.add(field.getName(),
-          ArrowFieldTypeFactory.toType(field.getType(), typeFactory));
+          ArrowFieldTypeFactory.toType(field, typeFactory));
     }
     return builder.build();
   }
 
+  private @Nullable Projector makeProjector(ImmutableIntList fields) {
+    if (containsListField(fields)) {
+      // Returning null selects ArrowEnumerable's direct vector-read path.
+      // Use that path for list fields because Gandiva does not support 
identity
+      // projection expressions over Arrow List vectors.
+      return null;
+    }
+
+    final List<ExpressionTree> expressionTrees = new ArrayList<>();
+    for (int fieldOrdinal : fields) {
+      Field field = schema.getFields().get(fieldOrdinal);
+      TreeNode node = TreeBuilder.makeField(field);
+      expressionTrees.add(TreeBuilder.makeExpression(node, field));
+    }
+    try {
+      return Projector.make(schema, expressionTrees);
+    } catch (GandivaException e) {
+      throw Util.toUnchecked(e);
+    }
+  }
+
+  private boolean containsListField(ImmutableIntList fields) {
+    for (int fieldOrdinal : fields) {
+      if (schema.getFields().get(fieldOrdinal).getType().getTypeID()
+          == ArrowType.ArrowTypeID.List) {
+        return true;
+      }
+    }
+    return false;
+  }
+
   /** Converts a single {@link ConditionToken} into a Gandiva {@link 
TreeNode}. */
   private TreeNode convertConditionToGandiva(ConditionToken token) {
     final List<TreeNode> treeNodes = new ArrayList<>(2);
diff --git 
a/arrow/src/test/java/org/apache/calcite/adapter/arrow/ArrowAdapterDataTypesTest.java
 
b/arrow/src/test/java/org/apache/calcite/adapter/arrow/ArrowAdapterDataTypesTest.java
index 566b4c531e..317d4dc26e 100644
--- 
a/arrow/src/test/java/org/apache/calcite/adapter/arrow/ArrowAdapterDataTypesTest.java
+++ 
b/arrow/src/test/java/org/apache/calcite/adapter/arrow/ArrowAdapterDataTypesTest.java
@@ -61,9 +61,27 @@ static void initializeArrowState(@TempDir Path sharedTempDir)
     ArrowDataTest arrowDataGenerator = new ArrowDataTest();
     arrowDataGenerator.writeArrowDataType(dataLocationFile);
 
+    File listDataLocationFile = 
arrowFilesDirectory.resolve("arrowlist.arrow").toFile();
+    ArrowDataTest arrowListDataGenerator = new ArrowDataTest();
+    arrowListDataGenerator.writeArrowListData(listDataLocationFile);
+
     arrow = ImmutableMap.of("model", 
modelFileTarget.toAbsolutePath().toString());
   }
 
+  @Test void testListProject() {
+    String sql = "select \"intListField\" from arrowlist";
+    String plan = "PLAN=ArrowToEnumerableConverter\n"
+        + "  ArrowTableScan(table=[[ARROW, ARROWLIST]], fields=[[0]])\n\n";
+    String result = "intListField=[0, 1]\n"
+        + "intListField=null\n"
+        + "intListField=[2, null]\n";
+    CalciteAssert.that()
+        .with(arrow)
+        .query(sql)
+        .returns(result)
+        .explainContains(plan);
+  }
+
   @Test void testTinyIntProject() {
     String sql = "select \"tinyIntField\" from arrowdatatype";
     String plan = "PLAN=ArrowToEnumerableConverter\n"
diff --git 
a/arrow/src/test/java/org/apache/calcite/adapter/arrow/ArrowDataTest.java 
b/arrow/src/test/java/org/apache/calcite/adapter/arrow/ArrowDataTest.java
index 8241a80dd1..7fd9f19c5e 100644
--- a/arrow/src/test/java/org/apache/calcite/adapter/arrow/ArrowDataTest.java
+++ b/arrow/src/test/java/org/apache/calcite/adapter/arrow/ArrowDataTest.java
@@ -41,6 +41,8 @@
 import org.apache.arrow.vector.TinyIntVector;
 import org.apache.arrow.vector.VarCharVector;
 import org.apache.arrow.vector.VectorSchemaRoot;
+import org.apache.arrow.vector.complex.ListVector;
+import org.apache.arrow.vector.complex.impl.UnionListWriter;
 import org.apache.arrow.vector.ipc.ArrowFileWriter;
 import org.apache.arrow.vector.types.DateUnit;
 import org.apache.arrow.vector.types.FloatingPointPrecision;
@@ -68,6 +70,8 @@
 import java.util.Calendar;
 import java.util.List;
 
+import static 
org.apache.arrow.vector.complex.BaseRepeatedValueVector.DATA_VECTOR_NAME;
+
 /**
  * Class that can be used to generate Arrow sample data into a data directory.
  */
@@ -150,6 +154,15 @@ private Schema makeArrowDateTypeSchema() {
     return new Schema(childrenBuilder.build(), null);
   }
 
+
+  private Schema makeArrowListSchema() {
+    FieldType listType = FieldType.nullable(new ArrowType.List());
+    FieldType elementType = FieldType.nullable(new ArrowType.Int(32, true));
+    Field elementField = new Field(DATA_VECTOR_NAME, elementType, null);
+    Field listField = new Field("intListField", listType, 
ImmutableList.of(elementField));
+    return new Schema(ImmutableList.of(listField), null);
+  }
+
   private Schema makeArrowSchema() {
     ImmutableList.Builder<Field> childrenBuilder = ImmutableList.builder();
     FieldType intType = FieldType.nullable(new ArrowType.Int(32, true));
@@ -329,6 +342,25 @@ public void writeArrowDataType(File file) throws 
IOException {
     fileOutputStream.close();
   }
 
+
+  public void writeArrowListData(File file) throws IOException {
+    Schema arrowSchema = makeArrowListSchema();
+    try (RootAllocator allocator = new RootAllocator(Integer.MAX_VALUE);
+         VectorSchemaRoot vectorSchemaRoot =
+             VectorSchemaRoot.create(arrowSchema, allocator);
+         FileOutputStream fileOutputStream = new FileOutputStream(file);
+         ArrowFileWriter arrowFileWriter =
+             new ArrowFileWriter(vectorSchemaRoot, null,
+                 fileOutputStream.getChannel())) {
+      arrowFileWriter.start();
+      int rowCount = 3;
+      vectorSchemaRoot.setRowCount(rowCount);
+      listField(vectorSchemaRoot.getVector("intListField"), rowCount);
+      arrowFileWriter.writeBatch();
+      arrowFileWriter.end();
+    }
+  }
+
   private void tinyIntField(FieldVector fieldVector, int rowCount) {
     TinyIntVector tinyIntVector = (TinyIntVector) fieldVector;
     tinyIntVector.setInitialCapacity(rowCount);
@@ -465,6 +497,31 @@ private void timeField(FieldVector fieldVector, int 
rowCount) {
     fieldVector.setValueCount(rowCount);
   }
 
+
+  private void listField(FieldVector fieldVector, int rowCount) {
+    ListVector listVector = (ListVector) fieldVector;
+    listVector.setInitialCapacity(rowCount);
+    listVector.allocateNew();
+    UnionListWriter writer = listVector.getWriter();
+    for (int i = 0; i < rowCount; i++) {
+      writer.setPosition(i);
+      if (i == 1) {
+        writer.writeNull();
+      } else {
+        writer.startList();
+        writer.writeInt(i);
+        if (i == 2) {
+          writer.writeNull();
+        } else {
+          writer.writeInt(i + 1);
+        }
+        writer.endList();
+      }
+    }
+    writer.setValueCount(rowCount);
+    fieldVector.setValueCount(rowCount);
+  }
+
   private void timestampSecField(FieldVector fieldVector, int rowCount) {
     TimeStampSecVector tsVector = (TimeStampSecVector) fieldVector;
     tsVector.setInitialCapacity(rowCount);

Reply via email to