srielau commented on code in PR #58087:
URL: https://github.com/apache/spark/pull/58087#discussion_r3824028716
##########
sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/orc/OrcUtils.scala:
##########
@@ -469,9 +469,16 @@ object OrcUtils extends Logging {
val typeDesc = new TypeDescription(ops.orcCategory)
typeDesc.setAttribute(CATALYST_TYPE_ATTRIBUTE_NAME, dt.typeName)
Some(typeDesc)
- case _: StringType =>
+ // Write CHAR/VARCHAR as ORC STRING plus spark.sql.catalyst.type, not
native
+ // ORC CHAR/VARCHAR. Native ORC maxLength would truncate/pad
independently of
+ // Spark store assignment. Hive-written native CHAR still round-trips
on read
+ // via toCatalystSchema. Stamp subclasses through charVarcharTypeName
so a
+ // later StringType arm cannot drop the length.
+ case s: StringType =>
val typeDesc = new TypeDescription(TypeDescription.Category.STRING)
- typeDesc.setAttribute(CATALYST_TYPE_ATTRIBUTE_NAME,
StringType.typeName)
+ typeDesc.setAttribute(
+ CATALYST_TYPE_ATTRIBUTE_NAME,
+ CharVarcharUtils.charVarcharTypeName(s).getOrElse(s.typeName))
Review Comment:
Collapsing onto `s: StringType` changes unbounded STRING stamping from
`StringType.typeName` (`string`) to `s.typeName`. Non-default collated STRING
will now persist on file-only ORC inference; Avro `avroStringSchema` still
stamps only CHAR/VARCHAR.
If preserving collation here is intended, please add a file-only ORC test
for `STRING COLLATE UTF8_LCASE`. If not,
`charVarcharTypeName(s).getOrElse(StringType.typeName)` keeps the old unbounded
contract.
##########
sql/core/src/test/scala/org/apache/spark/sql/CharVarcharTestSuite.scala:
##########
@@ -1341,6 +1396,385 @@ class BasicCharVarcharTestSuite extends
SharedSparkSession {
}
}
+ test("SPARK-58794: language surfaces keep CHAR/VARCHAR under
standardSemantics") {
+ withSQLConf(SQLConf.CHAR_VARCHAR_STANDARD_SEMANTICS.key -> "true") {
+ // CTAS / CREATE VIEW inherit projected CHAR/VARCHAR (D13).
+ withTable("std_src", "std_ctas") {
+ sql("CREATE TABLE std_src (c CHAR(5), v VARCHAR(5)) USING parquet")
+ sql("INSERT INTO std_src VALUES ('ab', 'ab')")
+ sql("CREATE TABLE std_ctas USING parquet AS SELECT c, v FROM std_src")
+ assert(spark.table("std_ctas").schema.map(_.dataType) ===
+ Seq(CharType(5), VarcharType(5)))
+ checkAnswer(
+ sql("SELECT concat('<', c, '>'), concat('<', v, '>') FROM std_ctas"),
+ Row("<ab >", "<ab>"))
+ val cteDf = sql(
+ "WITH t AS (SELECT c, v FROM std_src) SELECT typeof(c), typeof(v)
FROM t")
+ checkAnswer(cteDf, Row("char(5)", "varchar(5)"))
+ }
+ withTable("std_view_src") {
+ withView("std_cv_view", "std_cv_view_v") {
+ sql("CREATE TABLE std_view_src (c CHAR(4), v VARCHAR(4)) USING
parquet")
+ sql("INSERT INTO std_view_src VALUES ('xy', 'ab')")
+ sql("CREATE VIEW std_cv_view AS SELECT c FROM std_view_src")
+ sql("CREATE VIEW std_cv_view_v AS SELECT v FROM std_view_src")
+ assert(spark.table("std_cv_view").schema.head.dataType ===
CharType(4))
+ assert(spark.table("std_cv_view_v").schema.head.dataType ===
VarcharType(4))
+ checkAnswer(
+ sql("SELECT concat('<', c, '>') FROM std_cv_view"), Row("<xy >"))
+ checkAnswer(sql("SELECT v FROM std_cv_view_v"), Row("ab"))
+ }
+ }
+
+ // ALTER COLUMN equal-length CHAR/VARCHAR remains supported with
first-class types.
+ // VARCHAR widen / CHAR->VARCHAR are allowed by CheckAnalysis on V2
tables
+ // (see DSV2CharVarcharDDLTestSuite); V1 file-source ALTER only evolves
collation
+ // (same StringConstraint), so length changes stay rejected there.
+ withTable("std_alter") {
+ sql("CREATE TABLE std_alter (c CHAR(4), v VARCHAR(4)) USING parquet")
+ sql("ALTER TABLE std_alter CHANGE COLUMN c TYPE CHAR(4)")
+ sql("ALTER TABLE std_alter CHANGE COLUMN v TYPE VARCHAR(4)")
+ assert(spark.table("std_alter").schema.map(_.dataType) ===
+ Seq(CharType(4), VarcharType(4)))
+ intercept[AnalysisException] {
+ sql("ALTER TABLE std_alter CHANGE COLUMN c TYPE CHAR(5)")
+ }
+ intercept[AnalysisException] {
+ sql("ALTER TABLE std_alter CHANGE COLUMN v TYPE VARCHAR(5)")
+ }
+ }
+
+ // Session variables: DECLARE / SET keep the type and apply CAST
assignment.
+ sql("DECLARE OR REPLACE VARIABLE std_char_var CHAR(4)")
+ sql("DECLARE OR REPLACE VARIABLE std_varchar_var VARCHAR(4)")
+ try {
+ sql("SET VARIABLE std_char_var = 'ab'")
+ val charVarDf = sql("SELECT std_char_var AS c")
+ assert(charVarDf.schema.head.dataType === CharType(4))
+ checkAnswer(sql("SELECT concat('<', std_char_var, '>')"), Row("<ab
>"))
+ // Oversize by trailing blanks only is trimmed to fit CHAR(n).
+ sql("SET VARIABLE std_char_var = 'abcd '")
+ checkAnswer(sql("SELECT concat('<', std_char_var, '>')"),
Row("<abcd>"))
+ intercept[SparkRuntimeException] {
+ sql("SET VARIABLE std_char_var = 'abcde'").collect()
+ }
+
+ sql("SET VARIABLE std_varchar_var = 'ab'")
+ val varcharVarDf = sql("SELECT std_varchar_var AS v")
+ assert(varcharVarDf.schema.head.dataType === VarcharType(4))
+ checkAnswer(sql("SELECT concat('<', std_varchar_var, '>')"),
Row("<ab>"))
+ // Oversize by trailing blanks only is trimmed to fit.
+ sql("SET VARIABLE std_varchar_var = 'abcd '")
+ checkAnswer(sql("SELECT std_varchar_var"), Row("abcd"))
+ intercept[SparkRuntimeException] {
+ sql("SET VARIABLE std_varchar_var = 'abcde'").collect()
+ }
+ } finally {
+ sql("DROP TEMPORARY VARIABLE IF EXISTS std_char_var")
+ sql("DROP TEMPORARY VARIABLE IF EXISTS std_varchar_var")
+ }
+
+ // SQL scripting local variables keep CHAR/VARCHAR inside a compound
statement.
+ val localVarScript =
+ """
+ |BEGIN
+ | DECLARE c CHAR(4);
+ | DECLARE v VARCHAR(4);
+ | SET c = 'ab';
+ | SET v = 'cd';
+ | SELECT typeof(c), concat('<', c, '>'), typeof(v), concat('<', v,
'>');
+ |END
+ |""".stripMargin
+ val localVarDf = sql(localVarScript)
+ assert(localVarDf.schema.map(_.dataType) ===
+ Seq(StringType, StringType, StringType, StringType))
+ checkAnswer(localVarDf, Row("char(4)", "<ab >", "varchar(4)", "<cd>"))
+ // Trailing-blank trim on local SET into CHAR/VARCHAR.
+ checkAnswer(
+ sql(
+ """
+ |BEGIN
+ | DECLARE c CHAR(4);
+ | DECLARE v VARCHAR(4);
+ | SET c = 'abcd ';
+ | SET v = 'abcd ';
+ | SELECT concat('<', c, '>'), v;
+ |END
+ |""".stripMargin),
+ Row("<abcd>", "abcd"))
+ intercept[SparkRuntimeException] {
+ sql(
+ """
+ |BEGIN
+ | DECLARE c CHAR(4);
+ | SET c = 'abcde';
+ |END
+ |""".stripMargin).collect()
+ }
+ intercept[SparkRuntimeException] {
+ sql(
+ """
+ |BEGIN
+ | DECLARE v VARCHAR(4);
+ | SET v = 'abcde';
+ |END
+ |""".stripMargin).collect()
+ }
+
+ // Cursor FETCH INTO CHAR/VARCHAR locals applies store assignment (pad /
length).
+ withSQLConf(SQLConf.SQL_SCRIPTING_CURSOR_ENABLED.key -> "true") {
+ val cursorScript =
+ """
+ |BEGIN
+ | DECLARE fetched_c CHAR(4);
+ | DECLARE fetched_v VARCHAR(4);
+ | DECLARE cur CURSOR FOR
+ | SELECT cast('ab' AS CHAR(4)) AS c, cast('cd' AS VARCHAR(4))
AS v;
+ | OPEN cur;
+ | FETCH cur INTO fetched_c, fetched_v;
+ | SELECT typeof(fetched_c), concat('<', fetched_c, '>'),
+ | typeof(fetched_v), concat('<', fetched_v, '>');
+ | CLOSE cur;
+ |END
+ |""".stripMargin
+ checkAnswer(
+ sql(cursorScript),
+ Row("char(4)", "<ab >", "varchar(4)", "<cd>"))
+
+ // FETCH plain STRING into CHAR pads via assignment cast.
+ checkAnswer(
+ sql(
+ """
+ |BEGIN
+ | DECLARE fetched CHAR(4);
+ | DECLARE cur CURSOR FOR SELECT 'ab' AS c;
+ | OPEN cur;
+ | FETCH cur INTO fetched;
+ | SELECT typeof(fetched), concat('<', fetched, '>');
+ | CLOSE cur;
+ |END
+ |""".stripMargin),
+ Row("char(4)", "<ab >"))
+
+ // FETCH into a wider CHAR pads; trailing blanks trim into a shorter
target.
+ checkAnswer(
+ sql(
+ """
+ |BEGIN
+ | DECLARE fetched CHAR(5);
+ | DECLARE cur CURSOR FOR SELECT cast('xy' AS CHAR(2)) AS c;
+ | OPEN cur;
+ | FETCH cur INTO fetched;
+ | SELECT concat('<', fetched, '>');
+ | CLOSE cur;
+ |END
+ |""".stripMargin),
+ Row("<xy >"))
+ checkAnswer(
+ sql(
+ """
+ |BEGIN
+ | DECLARE fetched_c CHAR(4);
+ | DECLARE fetched_v VARCHAR(4);
+ | DECLARE cur CURSOR FOR
+ | SELECT cast('abcd ' AS CHAR(5)) AS c, cast('abcd ' AS
VARCHAR(5)) AS v;
+ | OPEN cur;
+ | FETCH cur INTO fetched_c, fetched_v;
+ | SELECT concat('<', fetched_c, '>'), fetched_v;
+ | CLOSE cur;
+ |END
+ |""".stripMargin),
+ Row("<abcd>", "abcd"))
+ intercept[SparkRuntimeException] {
+ sql(
+ """
+ |BEGIN
+ | DECLARE fetched VARCHAR(2);
+ | DECLARE cur CURSOR FOR SELECT cast('abcd' AS VARCHAR(4)) AS v;
+ | OPEN cur;
+ | FETCH cur INTO fetched;
+ | CLOSE cur;
+ |END
+ |""".stripMargin).collect()
+ }
+ }
+
+ // SQL FUNCTION params/RETURNS apply store assignment (pad, blank-trim,
overflow).
+ sql("CREATE OR REPLACE TEMPORARY FUNCTION std_char_fn() RETURNS CHAR(3)
RETURN 'a'")
+ sql(
+ """CREATE OR REPLACE TEMPORARY FUNCTION std_varchar_ret()
+ |RETURNS VARCHAR(3) RETURN 'ab'""".stripMargin)
+ sql(
+ """CREATE OR REPLACE TEMPORARY FUNCTION std_char_param(x CHAR(3))
+ |RETURNS CHAR(3) RETURN x""".stripMargin)
+ sql(
+ """CREATE OR REPLACE TEMPORARY FUNCTION std_varchar_param(x VARCHAR(3))
+ |RETURNS VARCHAR(3) RETURN x""".stripMargin)
+ try {
+ val fnDf = sql("SELECT std_char_fn() AS c")
+ assert(fnDf.schema.head.dataType === CharType(3))
+ checkAnswer(sql("SELECT concat('<', std_char_fn(), '>')"), Row("<a
>"))
+ val retVarcharDf = sql("SELECT std_varchar_ret() AS v")
+ assert(retVarcharDf.schema.head.dataType === VarcharType(3))
+ checkAnswer(retVarcharDf, Row("ab"))
+
+ // STRING -> CHAR(n) param: pad; trailing blanks trim; non-blank
overflow errors.
+ val charParamDf = sql("SELECT std_char_param('a') AS c")
+ assert(charParamDf.schema.head.dataType === CharType(3))
+ checkAnswer(sql("SELECT concat('<', std_char_param('a'), '>')"),
Row("<a >"))
+ checkAnswer(
+ sql("SELECT concat('<', std_char_param('abc '), '>')"),
+ Row("<abc>"))
+ intercept[SparkRuntimeException] {
+ sql("SELECT std_char_param('abcd')").collect()
+ }
+
+ val paramDf = sql("SELECT std_varchar_param('ab') AS v")
+ assert(paramDf.schema.head.dataType === VarcharType(3))
+ checkAnswer(paramDf, Row("ab"))
+ checkAnswer(sql("SELECT std_varchar_param('abc ')"), Row("abc"))
+ intercept[SparkRuntimeException] {
+ sql("SELECT std_varchar_param('abcd')").collect()
+ }
+ } finally {
+ sql("DROP TEMPORARY FUNCTION IF EXISTS std_char_fn")
+ sql("DROP TEMPORARY FUNCTION IF EXISTS std_varchar_ret")
+ sql("DROP TEMPORARY FUNCTION IF EXISTS std_char_param")
+ sql("DROP TEMPORARY FUNCTION IF EXISTS std_varchar_param")
+ }
+
+ // ORC catalog tables stamp the catalyst type so typeof survives
write/read.
+ withTable("std_orc") {
+ sql("CREATE TABLE std_orc (c CHAR(5), v VARCHAR(5)) USING orc")
+ sql("INSERT INTO std_orc VALUES ('ab', 'cd')")
+ assert(spark.table("std_orc").schema.map(_.dataType) ===
+ Seq(CharType(5), VarcharType(5)))
+ checkAnswer(
+ sql("SELECT concat('<', c, '>'), concat('<', v, '>') FROM std_orc"),
+ Row("<ab >", "<cd>"))
+ }
+
+ // File-only ORC inference recovers the catalyst type stamped on write.
+ withTempPath { dir =>
+ val path = dir.getCanonicalPath
+ spark.range(1).selectExpr("cast('ab' AS CHAR(4)) AS c")
+ .write.mode("overwrite").orc(path)
+ val orcDf = spark.read.orc(path)
+ assert(orcDf.schema.head.dataType === CharType(4))
+ checkAnswer(orcDf.selectExpr("concat('<', c, '>')"), Row("<ab >"))
+ // Reading with first-class types off replaces CHAR with STRING even
if the
+ // file was stamped under standardSemantics.
+ withSQLConf(SQLConf.CHAR_VARCHAR_STANDARD_SEMANTICS.key -> "false") {
+ val readOff = spark.read.orc(path)
+ assert(readOff.schema.head.dataType === StringType)
+ }
+ }
+ // First-class types off: CAST CHAR is STRING before the writer, so ORC
does not stamp CHAR.
+ withSQLConf(SQLConf.CHAR_VARCHAR_STANDARD_SEMANTICS.key -> "false") {
+ withTempPath { dir =>
+ val path = dir.getCanonicalPath
+ spark.range(1).selectExpr("cast('ab' AS CHAR(4)) AS c")
+ .write.mode("overwrite").orc(path)
+ assert(spark.read.orc(path).schema.head.dataType === StringType)
+ }
+ }
+ // preserveCharVarcharTypeInfo also keeps first-class types, so write
still stamps CHAR.
+ withSQLConf(
+ SQLConf.CHAR_VARCHAR_STANDARD_SEMANTICS.key -> "false",
+ SQLConf.PRESERVE_CHAR_VARCHAR_TYPE_INFO.key -> "true") {
+ withTempPath { dir =>
+ val path = dir.getCanonicalPath
+ spark.range(1).selectExpr("cast('ab' AS CHAR(4)) AS c")
+ .write.mode("overwrite").orc(path)
+ assert(spark.read.orc(path).schema.head.dataType === CharType(4))
+ }
+ }
+
+ // Avro conversion stamps CHAR/VARCHAR (including nested fields and CHAR
map keys).
+ // The avro data source lives in connector/avro; sql/core still owns
toAvroType,
+ // serializer, and deserializer, so round-trip them here with a
DataFileWriter.
+ val converters = org.apache.spark.sql.avro.SchemaConverters
+ Seq(CharType(5), VarcharType(7)).foreach { dt =>
+ val avro = converters.toAvroType(dt, nullable = false)
+ val back = converters.toSqlType(avro).dataType
+ assert(back === dt, s"Avro round-trip lost $dt, got $back")
+ }
+ val nestedSchema = new StructType()
+ .add("c", CharType(4))
+ .add("s", new StructType().add("f", VarcharType(3)))
+ .add("m", MapType(CharType(2), VarcharType(3)))
+ val nestedAvro = converters.toAvroType(nestedSchema, nullable = false)
+ assert(converters.toSqlType(nestedAvro).dataType === nestedSchema)
+ val badStamp = org.apache.avro.SchemaBuilder.builder().stringType()
+ badStamp.addProp("spark.sql.catalyst.type", "int")
+ val badStampErr = intercept[Exception] {
+ converters.toSqlType(badStamp)
+ }
+ assert(badStampErr.getMessage.contains("STRING subtype") ||
+ Option(badStampErr.getCause).exists(_.getMessage.contains("STRING
subtype")))
+ // Flag-off replace happens before Avro sees the type, so CHAR is not
stamped.
+ withSQLConf(SQLConf.CHAR_VARCHAR_STANDARD_SEMANTICS.key -> "false") {
+ val replaced =
CharVarcharUtils.replaceCharVarcharWithString(CharType(4))
+ assert(replaced === StringType)
+ val replacedAvro = converters.toAvroType(replaced, nullable = false)
+ assert(replacedAvro.getProp("spark.sql.catalyst.type") === null)
+ assert(converters.toSqlType(replacedAvro).dataType === StringType)
+ }
+
+ withTempPath { file =>
+ val ser = new org.apache.spark.sql.avro.AvroSerializer(
+ nestedSchema, nestedAvro, nullable = false)
+ val row = org.apache.spark.sql.catalyst.InternalRow(
+ org.apache.spark.unsafe.types.UTF8String.fromString("ab "),
+ org.apache.spark.sql.catalyst.InternalRow(
+ org.apache.spark.unsafe.types.UTF8String.fromString("xy")),
+ new org.apache.spark.sql.catalyst.util.ArrayBasedMapData(
+ new org.apache.spark.sql.catalyst.util.GenericArrayData(Array(
+ org.apache.spark.unsafe.types.UTF8String.fromString("k "))),
+ new org.apache.spark.sql.catalyst.util.GenericArrayData(Array(
+ org.apache.spark.unsafe.types.UTF8String.fromString("v")))))
+ val record = ser.serialize(row)
+ .asInstanceOf[org.apache.avro.generic.GenericRecord]
+ val writer = new org.apache.avro.file.DataFileWriter(
+ new
org.apache.avro.generic.GenericDatumWriter[org.apache.avro.generic.GenericRecord](
+ nestedAvro))
+ writer.create(nestedAvro, file)
+ writer.append(record)
+ writer.close()
+ val reader = new org.apache.avro.file.DataFileReader(
+ file,
+ new
org.apache.avro.generic.GenericDatumReader[org.apache.avro.generic.GenericRecord]())
+ try {
+ assert(converters.toSqlType(reader.getSchema).dataType ===
nestedSchema)
+ val deser = new org.apache.spark.sql.avro.AvroDeserializer(
+ nestedAvro, nestedSchema, "CORRECTED", false, "", -1)
+ val back = deser.deserialize(reader.next()).get
+ .asInstanceOf[org.apache.spark.sql.catalyst.InternalRow]
+ assert(back.getUTF8String(0).toString === "ab ")
Review Comment:
This only checks the top-level CHAR. The interesting part of this fixture is
the nested VARCHAR and `MapType(CharType(2), VarcharType(3))`. Please assert
those deserialized values too (and that the map key round-trips), otherwise the
schema converter can pass while ser/de of map keys regresses.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]