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

JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git


The following commit(s) were added to refs/heads/master by this push:
     new 5df8b9846e [common][spark] Fix ClassCastException when reading a 
VECTOR column as ARRAY (#8830)
5df8b9846e is described below

commit 5df8b9846e465c7514aa00f1baa9c93b9c3f144d
Author: XiaoHongbo <[email protected]>
AuthorDate: Thu Jul 23 21:00:35 2026 +0800

    [common][spark] Fix ClassCastException when reading a VECTOR column as 
ARRAY (#8830)
---
 .../apache/paimon/data/columnar/ColumnarArray.java |  4 ++++
 .../data/columnar/VectorizedColumnBatch.java       |  7 +++++-
 .../data/columnar/ColumnarRowWithVectorTest.java   | 27 ++++++++++++++++++++++
 .../apache/paimon/spark/data/SparkArrayData.scala  | 15 ++++++++----
 4 files changed, 47 insertions(+), 6 deletions(-)

diff --git 
a/paimon-common/src/main/java/org/apache/paimon/data/columnar/ColumnarArray.java
 
b/paimon-common/src/main/java/org/apache/paimon/data/columnar/ColumnarArray.java
index 31ea8bedae..aec42896e1 100644
--- 
a/paimon-common/src/main/java/org/apache/paimon/data/columnar/ColumnarArray.java
+++ 
b/paimon-common/src/main/java/org/apache/paimon/data/columnar/ColumnarArray.java
@@ -146,6 +146,10 @@ public final class ColumnarArray implements InternalArray, 
DataSetters, Serializ
 
     @Override
     public InternalArray getArray(int pos) {
+        if (data instanceof VecColumnVector) {
+            // A nested VECTOR is exposed as ARRAY; a vector is an array.
+            return ((VecColumnVector) data).getVector(offset + pos);
+        }
         InternalArray array = ((ArrayColumnVector) data).getArray(offset + 
pos);
         if (array instanceof ColumnarArray) {
             ((ColumnarArray) array).setFileIO(fileIO);
diff --git 
a/paimon-common/src/main/java/org/apache/paimon/data/columnar/VectorizedColumnBatch.java
 
b/paimon-common/src/main/java/org/apache/paimon/data/columnar/VectorizedColumnBatch.java
index 01c6037ca6..a0e51b2bbf 100644
--- 
a/paimon-common/src/main/java/org/apache/paimon/data/columnar/VectorizedColumnBatch.java
+++ 
b/paimon-common/src/main/java/org/apache/paimon/data/columnar/VectorizedColumnBatch.java
@@ -122,7 +122,12 @@ public class VectorizedColumnBatch implements Serializable 
{
     }
 
     public InternalArray getArray(int rowId, int colId) {
-        return ((ArrayColumnVector) columns[colId]).getArray(rowId);
+        ColumnVector column = columns[colId];
+        if (column instanceof VecColumnVector) {
+            // A VECTOR is exposed as ARRAY; a vector is an array.
+            return ((VecColumnVector) column).getVector(rowId);
+        }
+        return ((ArrayColumnVector) column).getArray(rowId);
     }
 
     public InternalVector getVector(int rowId, int colId) {
diff --git 
a/paimon-common/src/test/java/org/apache/paimon/data/columnar/ColumnarRowWithVectorTest.java
 
b/paimon-common/src/test/java/org/apache/paimon/data/columnar/ColumnarRowWithVectorTest.java
index d33ba4c02d..99401ad034 100644
--- 
a/paimon-common/src/test/java/org/apache/paimon/data/columnar/ColumnarRowWithVectorTest.java
+++ 
b/paimon-common/src/test/java/org/apache/paimon/data/columnar/ColumnarRowWithVectorTest.java
@@ -47,6 +47,33 @@ public class ColumnarRowWithVectorTest {
         assertThat(row.getVector(0).toFloatArray()).isEqualTo(new float[] 
{4.0f, 5.0f, 6.0f});
     }
 
+    @Test
+    public void testVectorReadAsArray() {
+        // VECTOR is exposed as ARRAY (e.g. Flink maps VECTOR -> ARRAY), so it 
is read via
+        // getArray; a vector is an array, so this must not throw a 
ClassCastException.
+        float[] values = new float[] {1.0f, 2.0f, 3.0f, 4.0f, 5.0f, 6.0f};
+        VectorizedColumnBatch batch = makeColumnBatch(values, 2, null);
+
+        ColumnarRow row = new ColumnarRow(batch);
+        row.setRowId(0);
+        assertThat(row.getArray(0).toFloatArray()).isEqualTo(new float[] 
{1.0f, 2.0f, 3.0f});
+
+        row.setRowId(1);
+        assertThat(row.getArray(0).toFloatArray()).isEqualTo(new float[] 
{4.0f, 5.0f, 6.0f});
+    }
+
+    @Test
+    public void testNestedVectorReadAsArray() {
+        // Nested VECTOR (ARRAY<VECTOR>/MAP<..,VECTOR>) reads via 
ColumnarArray.getArray.
+        float[] values = new float[] {1.0f, 2.0f, 3.0f, 4.0f, 5.0f, 6.0f};
+        VectorizedColumnBatch batch = makeColumnBatch(values, 2, null);
+        VecColumnVector vectorColumn = (VecColumnVector) batch.columns[0];
+
+        ColumnarArray outer = new ColumnarArray(vectorColumn, 0, 2);
+        assertThat(outer.getArray(0).toFloatArray()).isEqualTo(new float[] 
{1.0f, 2.0f, 3.0f});
+        assertThat(outer.getArray(1).toFloatArray()).isEqualTo(new float[] 
{4.0f, 5.0f, 6.0f});
+    }
+
     @Test
     public void testVectorNullable() {
         float[] values = new float[] {1.0f, 2.0f, 3.0f, 4.0f, 5.0f, 6.0f};
diff --git 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/data/SparkArrayData.scala
 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/data/SparkArrayData.scala
index 0fdf0dda19..5635f035f3 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/data/SparkArrayData.scala
+++ 
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/data/SparkArrayData.scala
@@ -20,7 +20,7 @@ package org.apache.paimon.spark.data
 
 import org.apache.paimon.data.{Blob, BlobView, InternalArray}
 import org.apache.paimon.spark.DataConverter
-import org.apache.paimon.types.{ArrayType => PaimonArrayType, BigIntType, 
BlobType, DataType => PaimonDataType, DataTypeChecks, RowType}
+import org.apache.paimon.types.{ArrayType => PaimonArrayType, BigIntType, 
BlobType, DataType => PaimonDataType, DataTypeChecks, RowType, VectorType}
 import org.apache.paimon.utils.InternalRowUtils
 
 import org.apache.spark.sql.catalyst.InternalRow
@@ -141,10 +141,15 @@ abstract class AbstractSparkArrayData extends 
SparkArrayData {
   override def getStruct(ordinal: Int, numFields: Int): InternalRow = 
DataConverter
     .fromPaimon(paimonArray.getRow(ordinal, numFields), 
elementType.asInstanceOf[RowType])
 
-  override def getArray(ordinal: Int): ArrayData = DataConverter.fromPaimon(
-    paimonArray.getArray(ordinal),
-    elementType.asInstanceOf[PaimonArrayType],
-    blobAsDescriptor)
+  override def getArray(ordinal: Int): ArrayData = elementType match {
+    // A nested VECTOR is exposed as a Spark array; read it as a vector, not 
an array.
+    case vectorType: VectorType =>
+      DataConverter.fromPaimon(paimonArray.getVector(ordinal), vectorType)
+    case arrayType: PaimonArrayType =>
+      DataConverter.fromPaimon(paimonArray.getArray(ordinal), arrayType, 
blobAsDescriptor)
+    case other =>
+      throw new UnsupportedOperationException("Not an array type: " + other)
+  }
 
   override def getMap(ordinal: Int): MapData =
     DataConverter.fromPaimon(paimonArray.getMap(ordinal), elementType)

Reply via email to