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)