Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@ import java.util.Properties

import org.apache.spark.sql.{Column, Row}
import org.apache.spark.sql.catalyst.expressions.Literal
import org.apache.spark.sql.types.{ArrayType, DecimalType, FloatType, ShortType}
import org.apache.spark.sql.types.{ArrayType, DecimalType, FloatType, NullType, ShortType}
import org.apache.spark.tags.DockerTest

/**
Expand Down Expand Up @@ -161,6 +161,10 @@ class PostgresIntegrationSuite extends DockerJDBCIntegrationSuite {
conn.prepareStatement("INSERT INTO custom_type (type_array, type) VALUES" +
"('{1,fds,fdsa}','fdasfasdf')").executeUpdate()

conn.prepareStatement(
"CREATE FUNCTION test_null() RETURNS VOID AS $$ BEGIN RETURN; END; $$ LANGUAGE plpgsql")
.executeUpdate()

}

test("Type mapping for various types") {
Expand Down Expand Up @@ -457,4 +461,14 @@ class PostgresIntegrationSuite extends DockerJDBCIntegrationSuite {
assert(infinitySeq.head.getTime == maxTimestamp)
assert(negativeInfinitySeq.head.getTime == minTimeStamp)
}


test("SPARK-47407: Support java.sql.Types.NULL for NullType") {
val df = spark.read.format("jdbc")
.option("url", jdbcUrl)
.option("query", "SELECT test_null()")
.load()
assert(df.schema.head.dataType === NullType)
checkAnswer(df, Seq(Row(null)))
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -219,9 +219,10 @@ object JdbcUtils extends Logging with SQLConfHelper {
case java.sql.Types.VARBINARY => BinaryType
case java.sql.Types.VARCHAR if conf.charVarcharAsString => StringType
case java.sql.Types.VARCHAR => VarcharType(precision)
case java.sql.Types.NULL => NullType
case _ =>
// For unmatched types:
// including java.sql.Types.ARRAY,DATALINK,DISTINCT,JAVA_OBJECT,NULL,OTHER,REF_CURSOR,
// including java.sql.Types.ARRAY,DATALINK,DISTINCT,JAVA_OBJECT,OTHER,REF_CURSOR,
// TIME_WITH_TIMEZONE,TIMESTAMP_WITH_TIMEZONE, and among others.
val jdbcType = classOf[JDBCType].getEnumConstants()
.find(_.getVendorTypeNumber == sqlType)
Expand Down Expand Up @@ -542,6 +543,9 @@ object JdbcUtils extends Logging with SQLConfHelper {
array => new GenericArrayData(elementConversion(array.getArray)))
row.update(pos, array)

case NullType =>
(_: ResultSet, row: InternalRow, pos: Int) => row.update(pos, null)

case _ => throw QueryExecutionErrors.unsupportedJdbcTypeError(dt.catalogString)
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,7 @@ private object PostgresDialect extends JdbcDialect with SQLConfHelper {
// timetz represents time with time zone, currently it maps to Types.TIME.
// We need to change to Types.TIME_WITH_TIMEZONE if the upstream changes.
Some(TimestampType)
case Types.OTHER if "void".equalsIgnoreCase(typeName) => Some(NullType)
case Types.OTHER => Some(StringType)
case _ if "text".equalsIgnoreCase(typeName) => Some(StringType) // sqlType is Types.VARCHAR
case Types.ARRAY =>
Expand Down