This is an automated email from the ASF dual-hosted git repository.
biyan 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 36dcebc20 [spark] Support varchar/char type (#3361)
36dcebc20 is described below
commit 36dcebc206cf85a7abadad3a25c9724e6a3ff737
Author: Zouxxyy <[email protected]>
AuthorDate: Tue May 21 18:34:27 2024 +0800
[spark] Support varchar/char type (#3361)
---
docs/content/spark/quick-start.md | 12 +++-
.../org/apache/paimon/spark/SparkInternalRow.java | 6 +-
.../org/apache/paimon/spark/SparkTypeUtils.java | 8 ++-
.../paimon/spark/PaimonPartitionManagement.scala | 5 +-
.../apache/paimon/spark/SparkInternalRowTest.java | 10 +++-
.../org/apache/paimon/spark/SparkReadITCase.java | 9 ++-
.../org/apache/paimon/spark/SparkTypeTest.java | 4 ++
.../org/apache/paimon/spark/sql/DDLTestBase.scala | 67 ++++++++++++++++++++++
.../spark/sql/PaimonPartitionManagementTest.scala | 11 +++-
.../paimon/spark/sql/SparkVersionSupport.scala | 2 +
10 files changed, 121 insertions(+), 13 deletions(-)
diff --git a/docs/content/spark/quick-start.md
b/docs/content/spark/quick-start.md
index d8922baf3..253ed0b2f 100644
--- a/docs/content/spark/quick-start.md
+++ b/docs/content/spark/quick-start.md
@@ -290,7 +290,17 @@ All Spark's data types are available in package
`org.apache.spark.sql.types`.
</tr>
<tr>
<td><code>StringType</code></td>
- <td><code>VarCharType</code>, <code>CharType</code></td>
+ <td><code>VarCharType(Integer.MAX_VALUE)</code></td>
+ <td>true</td>
+ </tr>
+ <tr>
+ <td><code>VarCharType(length)</code></td>
+ <td><code>VarCharType(length)</code></td>
+ <td>true</td>
+ </tr>
+ <tr>
+ <td><code>CharType(length)</code></td>
+ <td><code>CharType(length)</code></td>
<td>true</td>
</tr>
<tr>
diff --git
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkInternalRow.java
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkInternalRow.java
index a54f3e339..787a1329a 100644
---
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkInternalRow.java
+++
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkInternalRow.java
@@ -40,6 +40,7 @@ import org.apache.spark.sql.types.BinaryType;
import org.apache.spark.sql.types.BooleanType;
import org.apache.spark.sql.types.ByteType;
import org.apache.spark.sql.types.CalendarIntervalType;
+import org.apache.spark.sql.types.CharType;
import org.apache.spark.sql.types.DateType;
import org.apache.spark.sql.types.Decimal;
import org.apache.spark.sql.types.DecimalType;
@@ -53,6 +54,7 @@ import org.apache.spark.sql.types.StringType;
import org.apache.spark.sql.types.StructType;
import org.apache.spark.sql.types.TimestampType;
import org.apache.spark.sql.types.UserDefinedType;
+import org.apache.spark.sql.types.VarcharType;
import org.apache.spark.unsafe.types.CalendarInterval;
import org.apache.spark.unsafe.types.UTF8String;
@@ -205,7 +207,9 @@ public class SparkInternalRow extends
org.apache.spark.sql.catalyst.InternalRow
if (dataType instanceof DoubleType) {
return getDouble(ordinal);
}
- if (dataType instanceof StringType) {
+ if (dataType instanceof StringType
+ || dataType instanceof CharType
+ || dataType instanceof VarcharType) {
return getUTF8String(ordinal);
}
if (dataType instanceof DecimalType) {
diff --git
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkTypeUtils.java
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkTypeUtils.java
index 6d88b47c7..850183cdb 100644
---
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkTypeUtils.java
+++
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/SparkTypeUtils.java
@@ -75,12 +75,16 @@ public class SparkTypeUtils {
@Override
public DataType visit(CharType charType) {
- return DataTypes.StringType;
+ return new
org.apache.spark.sql.types.CharType(charType.getLength());
}
@Override
public DataType visit(VarCharType varCharType) {
- return DataTypes.StringType;
+ if (varCharType.getLength() == VarCharType.MAX_LENGTH) {
+ return DataTypes.StringType;
+ } else {
+ return new
org.apache.spark.sql.types.VarcharType(varCharType.getLength());
+ }
}
@Override
diff --git
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PaimonPartitionManagement.scala
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PaimonPartitionManagement.scala
index 8c5192e01..bebb1ac4a 100644
---
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PaimonPartitionManagement.scala
+++
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/PaimonPartitionManagement.scala
@@ -23,10 +23,11 @@ import org.apache.paimon.operation.FileStoreCommit
import org.apache.paimon.table.FileStoreTable
import org.apache.paimon.table.sink.BatchWriteBuilder
import org.apache.paimon.types.RowType
-import org.apache.paimon.utils.{FileStorePathFactory, RowDataPartitionComputer}
+import org.apache.paimon.utils.RowDataPartitionComputer
import org.apache.spark.sql.Row
import org.apache.spark.sql.catalyst.{CatalystTypeConverters, InternalRow}
+import org.apache.spark.sql.catalyst.util.CharVarcharUtils
import org.apache.spark.sql.connector.catalog.SupportsPartitionManagement
import org.apache.spark.sql.types.StructType
@@ -51,7 +52,7 @@ trait PaimonPartitionManagement extends
SupportsPartitionManagement {
override def dropPartition(internalRow: InternalRow): Boolean = {
// convert internalRow to row
val row: Row = CatalystTypeConverters
- .createToScalaConverter(partitionSchema())
+
.createToScalaConverter(CharVarcharUtils.replaceCharVarcharWithString(partitionSchema()))
.apply(internalRow)
.asInstanceOf[Row]
val rowDataPartitionComputer = new RowDataPartitionComputer(
diff --git
a/paimon-spark/paimon-spark-common/src/test/java/org/apache/paimon/spark/SparkInternalRowTest.java
b/paimon-spark/paimon-spark-common/src/test/java/org/apache/paimon/spark/SparkInternalRowTest.java
index a152a12e9..14080cd8a 100644
---
a/paimon-spark/paimon-spark-common/src/test/java/org/apache/paimon/spark/SparkInternalRowTest.java
+++
b/paimon-spark/paimon-spark-common/src/test/java/org/apache/paimon/spark/SparkInternalRowTest.java
@@ -28,6 +28,7 @@ import org.apache.paimon.data.Timestamp;
import org.apache.paimon.utils.DateTimeUtils;
import org.apache.spark.sql.catalyst.CatalystTypeConverters;
+import org.apache.spark.sql.catalyst.util.CharVarcharUtils;
import org.junit.jupiter.api.Test;
import java.math.BigDecimal;
@@ -54,6 +55,8 @@ public class SparkInternalRowTest {
GenericRow.of(
1,
fromString("jingsong"),
+ fromString("apache"),
+ fromString("paimon"),
22.2,
new GenericMap(
Stream.of(
@@ -79,9 +82,12 @@ public class SparkInternalRowTest {
Decimal.fromBigDecimal(BigDecimal.valueOf(65782123123.01), 38, 2),
Decimal.fromBigDecimal(BigDecimal.valueOf(62123123.5),
10, 1));
+ // CatalystTypeConverters does not support char and varchar, we need
to replace char and
+ // varchar with string
Function1<Object, Object> sparkConverter =
CatalystTypeConverters.createToScalaConverter(
- SparkTypeUtils.fromPaimonType(ALL_TYPES));
+ CharVarcharUtils.replaceCharVarcharWithString(
+ SparkTypeUtils.fromPaimonType(ALL_TYPES)));
org.apache.spark.sql.Row sparkRow =
(org.apache.spark.sql.Row)
sparkConverter.apply(new
SparkInternalRow(ALL_TYPES).replace(rowData));
@@ -90,6 +96,8 @@ public class SparkInternalRowTest {
"{"
+ "\"id\":1,"
+ "\"name\":\"jingsong\","
+ + "\"char\":\"apache\","
+ + "\"varchar\":\"paimon\","
+ "\"salary\":22.2,"
+
"\"locations\":{\"key1\":{\"posX\":1.2,\"posY\":2.3},\"key2\":{\"posX\":2.4,\"posY\":3.5}},"
+ "\"strArray\":[\"v1\",\"v5\"],"
diff --git
a/paimon-spark/paimon-spark-common/src/test/java/org/apache/paimon/spark/SparkReadITCase.java
b/paimon-spark/paimon-spark-common/src/test/java/org/apache/paimon/spark/SparkReadITCase.java
index cc6ddb5e4..2dfea46c1 100644
---
a/paimon-spark/paimon-spark-common/src/test/java/org/apache/paimon/spark/SparkReadITCase.java
+++
b/paimon-spark/paimon-spark-common/src/test/java/org/apache/paimon/spark/SparkReadITCase.java
@@ -177,7 +177,8 @@ public class SparkReadITCase extends SparkReadTestBase {
spark.sql("CREATE TABLE testCreateTableAs AS SELECT * FROM
testCreateTable");
List<Row> result = spark.sql("SELECT * FROM
testCreateTableAs").collectAsList();
-
assertThat(result.stream().map(Row::toString)).containsExactlyInAnyOrder("[1,a,b]");
+ assertThat(result.stream().map(Row::toString))
+ .containsExactlyInAnyOrder("[1,a,b ]");
// partitioned table
spark.sql(
@@ -224,11 +225,13 @@ public class SparkReadITCase extends SparkReadTestBase {
+ " 'file.format' = 'parquet',\n"
+ " 'path' = '%s')\n"
+ "]]",
- showCreateString("testTableAs", "a BIGINT", "b
STRING", "c STRING"),
+ showCreateString(
+ "testTableAs", "a BIGINT", "b
VARCHAR(10)", "c CHAR(10)"),
new Path(warehousePath,
"default.db/testTableAs")));
List<Row> resultProp = spark.sql("SELECT * FROM
testTableAs").collectAsList();
-
assertThat(resultProp.stream().map(Row::toString)).containsExactlyInAnyOrder("[1,a,b]");
+ assertThat(resultProp.stream().map(Row::toString))
+ .containsExactlyInAnyOrder("[1,a,b ]");
// primary key
spark.sql(
diff --git
a/paimon-spark/paimon-spark-common/src/test/java/org/apache/paimon/spark/SparkTypeTest.java
b/paimon-spark/paimon-spark-common/src/test/java/org/apache/paimon/spark/SparkTypeTest.java
index b9740b1a3..c9f82a178 100644
---
a/paimon-spark/paimon-spark-common/src/test/java/org/apache/paimon/spark/SparkTypeTest.java
+++
b/paimon-spark/paimon-spark-common/src/test/java/org/apache/paimon/spark/SparkTypeTest.java
@@ -39,6 +39,8 @@ public class SparkTypeTest {
1)) // posX and posY have field id 0 and
1, here we start from 2
.field("id", DataTypes.INT().notNull())
.field("name", DataTypes.STRING()) /* optional by default
*/
+ .field("char", DataTypes.CHAR(10))
+ .field("varchar", DataTypes.VARCHAR(10))
.field("salary", DataTypes.DOUBLE().notNull())
.field(
"locations",
@@ -79,6 +81,8 @@ public class SparkTypeTest {
"StructType("
+ "StructField(id,IntegerType,true),"
+ "StructField(name,StringType,true),"
+ + "StructField(char,CharType(10),true),"
+ + "StructField(varchar,VarcharType(10),true),"
+ "StructField(salary,DoubleType,true),"
+ nestedRowMapType
+ ","
diff --git
a/paimon-spark/paimon-spark-common/src/test/scala/org/apache/paimon/spark/sql/DDLTestBase.scala
b/paimon-spark/paimon-spark-common/src/test/scala/org/apache/paimon/spark/sql/DDLTestBase.scala
index 6fda6b26c..f62b87a60 100644
---
a/paimon-spark/paimon-spark-common/src/test/scala/org/apache/paimon/spark/sql/DDLTestBase.scala
+++
b/paimon-spark/paimon-spark-common/src/test/scala/org/apache/paimon/spark/sql/DDLTestBase.scala
@@ -20,6 +20,7 @@ package org.apache.paimon.spark.sql
import org.apache.paimon.spark.PaimonSparkTestBase
+import org.apache.spark.sql.Row
import org.junit.jupiter.api.Assertions
abstract class DDLTestBase extends PaimonSparkTestBase {
@@ -84,4 +85,70 @@ abstract class DDLTestBase extends PaimonSparkTestBase {
"SparkCatalog can only create paimon table, but current provider is
parquet"))
}
}
+
+ test("Paimon DDL: create table with char/varchar/string") {
+ Seq("orc", "avro").foreach(
+ format => {
+ withTable("paimon_tbl") {
+ spark.sql(
+ s"""
+ |CREATE TABLE paimon_tbl (id int, col_s1 char(9), col_s2
varchar(10), col_s3 string)
+ |USING PAIMON
+ |TBLPROPERTIES ('file.format' = '$format')
+ |""".stripMargin)
+
+ spark.sql(s"""
+ |insert into paimon_tbl values
+ |(1, 'Wednesday', 'Wednesday', 'Wednesday'),
+ |(2, 'Friday', 'Friday', 'Friday')
+ |""".stripMargin)
+
+ // check description
+ checkAnswer(
+ spark
+ .sql(s"DESC paimon_tbl")
+ .select("col_name", "data_type")
+ .where("col_name LIKE 'col_%'")
+ .orderBy("col_name"),
+ Row("col_s1", "char(9)") :: Row("col_s2", "varchar(10)") :: Row(
+ "col_s3",
+ "string") :: Nil
+ )
+
+ // check select
+ if (format == "orc" && !gteqSpark3_4) {
+ // Orc reader will right trim the char type, e.g. "Friday " =>
"Friday" (see orc's `CharTreeReader`)
+ // and Spark has a conf `spark.sql.readSideCharPadding` to auto
padding char only since 3.4 (default true)
+ // So when using orc with Spark3.4-, here will return "Friday"
+ checkAnswer(
+ spark.sql(s"select col_s1 from paimon_tbl where id = 2"),
+ Row("Friday") :: Nil
+ )
+ // Spark will auto create the filter like
Filter(isnotnull(col_s1#124) AND (col_s1#124 = Friday ))
+ // for char type, so here will not return any rows
+ checkAnswer(
+ spark.sql(s"select col_s1 from paimon_tbl where col_s1 =
'Friday'"),
+ Nil
+ )
+ } else {
+ checkAnswer(
+ spark.sql(s"select col_s1 from paimon_tbl where id = 2"),
+ Row("Friday ") :: Nil
+ )
+ checkAnswer(
+ spark.sql(s"select col_s1 from paimon_tbl where col_s1 =
'Friday'"),
+ Row("Friday ") :: Nil
+ )
+ }
+ checkAnswer(
+ spark.sql(s"select col_s2 from paimon_tbl where col_s2 =
'Friday'"),
+ Row("Friday") :: Nil
+ )
+ checkAnswer(
+ spark.sql(s"select col_s3 from paimon_tbl where col_s3 =
'Friday'"),
+ Row("Friday") :: Nil
+ )
+ }
+ })
+ }
}
diff --git
a/paimon-spark/paimon-spark-common/src/test/scala/org/apache/paimon/spark/sql/PaimonPartitionManagementTest.scala
b/paimon-spark/paimon-spark-common/src/test/scala/org/apache/paimon/spark/sql/PaimonPartitionManagementTest.scala
index 34f65131f..5f4c90cfe 100644
---
a/paimon-spark/paimon-spark-common/src/test/scala/org/apache/paimon/spark/sql/PaimonPartitionManagementTest.scala
+++
b/paimon-spark/paimon-spark-common/src/test/scala/org/apache/paimon/spark/sql/PaimonPartitionManagementTest.scala
@@ -161,12 +161,17 @@ class PaimonPartitionManagementTest extends
PaimonSparkTestBase {
checkAnswer(
spark.sql("select * from T"),
- Row("a", "b", 1L, 20230816L, "1132") :: Row("a", "b", 1L,
20230816L, "1133") :: Row(
+ Row("a", "b ", 1L, 20230816L, "1132") :: Row(
"a",
- "b",
+ "b ",
+ 1L,
+ 20230816L,
+ "1133") :: Row("a", "b ", 2L, 20230817L, "1132") ::
Row(
+ "a",
+ "b ",
2L,
20230817L,
- "1132") :: Row("a", "b", 2L, 20230817L, "1134") :: Nil
+ "1134") :: Nil
)
}
}
diff --git
a/paimon-spark/paimon-spark-common/src/test/scala/org/apache/paimon/spark/sql/SparkVersionSupport.scala
b/paimon-spark/paimon-spark-common/src/test/scala/org/apache/paimon/spark/sql/SparkVersionSupport.scala
index c8f5711a6..5ac408934 100644
---
a/paimon-spark/paimon-spark-common/src/test/scala/org/apache/paimon/spark/sql/SparkVersionSupport.scala
+++
b/paimon-spark/paimon-spark-common/src/test/scala/org/apache/paimon/spark/sql/SparkVersionSupport.scala
@@ -24,4 +24,6 @@ trait SparkVersionSupport {
lazy val sparkVersion: String = SPARK_VERSION
lazy val gteqSpark3_3: Boolean = sparkVersion >= "3.3"
+
+ lazy val gteqSpark3_4: Boolean = sparkVersion >= "3.4"
}