@@ -27,25 +27,16 @@ import org.apache.spark.sql.execution.datasources.{DataSourceStrategy, HadoopFsR
2727import org .apache .spark .sql .execution .datasources .v2 .DataSourceV2ScanRelation
2828import org .apache .spark .sql .execution .datasources .v2 .parquet .ParquetScan
2929import org .apache .spark .sql .internal .SQLConf
30- import org .apache .spark .sql .internal .SQLConf .LegacyBehaviorPolicy .{CORRECTED , LEGACY }
31- import org .apache .spark .sql .internal .SQLConf .ParquetOutputTimestampType .INT96
3230import org .apache .spark .sql .types ._
3331import org .apache .spark .tags .ExtendedSQLTest
3432import org .apache .spark .util .Utils
3533
3634import org .apache .hadoop .fs .Path
37- import org .apache .parquet .filter2 .predicate .{FilterApi , FilterPredicate , Operators }
35+ import org .apache .parquet .filter2 .predicate .{FilterApi , FilterPredicate }
3836import org .apache .parquet .filter2 .predicate .FilterApi ._
39- import org .apache .parquet .filter2 .predicate .Operators .{Column => _ , Eq , Gt , GtEq , Lt , LtEq , NotEq }
4037import org .apache .parquet .hadoop .{ParquetFileReader , ParquetInputFormat , ParquetOutputFormat }
4138import org .apache .parquet .hadoop .util .HadoopInputFile
4239
43- import java .sql .{Date , Timestamp }
44- import java .time .LocalDate
45-
46- import scala .reflect .ClassTag
47- import scala .reflect .runtime .universe .TypeTag
48-
4940abstract class GlutenParquetFilterSuite extends ParquetFilterSuite with GlutenSQLTestsBaseTrait {
5041 protected def checkFilterPredicate (
5142 predicate : Predicate ,
@@ -66,44 +57,6 @@ abstract class GlutenParquetFilterSuite extends ParquetFilterSuite with GlutenSQ
6657 getWorkspaceFilePath(" sql" , " core" , " src" , " test" , " resources" ).toString + " /" + name)
6758 }
6859
69- testGluten(" filter pushdown - timestamp" ) {
70- Seq (true , false ).foreach {
71- java8Api =>
72- Seq (CORRECTED , LEGACY ).foreach {
73- rebaseMode =>
74- val millisData = Seq (
75- " 1000-06-14 08:28:53.123" ,
76- " 1582-06-15 08:28:53.001" ,
77- " 1900-06-16 08:28:53.0" ,
78- " 2018-06-17 08:28:53.999" )
79- // INT96 doesn't support pushdown
80- withSQLConf(
81- SQLConf .DATETIME_JAVA8API_ENABLED .key -> java8Api.toString,
82- SQLConf .PARQUET_INT96_REBASE_MODE_IN_WRITE .key -> rebaseMode.toString,
83- SQLConf .PARQUET_OUTPUT_TIMESTAMP_TYPE .key -> INT96 .toString
84- ) {
85- import testImplicits ._
86- withTempPath {
87- file =>
88- millisData
89- .map(i => Tuple1 (Timestamp .valueOf(i)))
90- .toDF
91- .write
92- .format(dataSourceName)
93- .save(file.getCanonicalPath)
94- readParquetFile(file.getCanonicalPath) {
95- df =>
96- val schema = new SparkToParquetSchemaConverter (conf).convert(df.schema)
97- assertResult(None ) {
98- createParquetFilters(schema).createFilter(sources.IsNull (" _1" ))
99- }
100- }
101- }
102- }
103- }
104- }
105- }
106-
10760 testGluten(" SPARK-12218: 'Not' is included in Parquet filter pushdown" ) {
10861 import testImplicits ._
10962
@@ -428,153 +381,4 @@ class GlutenParquetV2FilterSuite extends GlutenParquetFilterSuite with GlutenSQL
428381 }
429382 }
430383 }
431-
432- /**
433- * Takes a sequence of products `data` to generate multi-level nested dataframes as new test data.
434- * It tests both non-nested and nested dataframes which are written and read back with Parquet
435- * datasource.
436- *
437- * This is different from [[ParquetTest.withParquetDataFrame ]] which does not test nested cases.
438- */
439- private def withNestedParquetDataFrame [T <: Product : ClassTag : TypeTag ](data : Seq [T ])(
440- runTest : (DataFrame , String , Any => Any ) => Unit ): Unit =
441- withNestedParquetDataFrame(spark.createDataFrame(data))(runTest)
442-
443- private def withNestedParquetDataFrame (inputDF : DataFrame )(
444- runTest : (DataFrame , String , Any => Any ) => Unit ): Unit = {
445- withNestedDataFrame(inputDF).foreach {
446- case (newDF, colName, resultFun) =>
447- withTempPath {
448- file =>
449- newDF.write.format(dataSourceName).save(file.getCanonicalPath)
450- readParquetFile(file.getCanonicalPath)(df => runTest(df, colName, resultFun))
451- }
452- }
453- }
454-
455- testGluten(" filter pushdown - date" ) {
456- implicit class StringToDate (s : String ) {
457- def date : Date = Date .valueOf(s)
458- }
459-
460- val data = Seq (" 1000-01-01" , " 2018-03-19" , " 2018-03-20" , " 2018-03-21" )
461- import testImplicits ._
462-
463- // Velox backend does not support rebaseMode being LEGACY.
464- Seq (false , true ).foreach {
465- java8Api =>
466- Seq (CORRECTED ).foreach {
467- rebaseMode =>
468- withSQLConf(
469- SQLConf .DATETIME_JAVA8API_ENABLED .key -> java8Api.toString,
470- SQLConf .PARQUET_REBASE_MODE_IN_WRITE .key -> rebaseMode.toString) {
471- val dates = data.map(i => Tuple1 (Date .valueOf(i))).toDF()
472- withNestedParquetDataFrame(dates) {
473- case (inputDF, colName, fun) =>
474- implicit val df : DataFrame = inputDF
475-
476- def resultFun (dateStr : String ): Any = {
477- val parsed = if (java8Api) LocalDate .parse(dateStr) else Date .valueOf(dateStr)
478- fun(parsed)
479- }
480-
481- val dateAttr : Expression = df(colName).expr
482- assert(df(colName).expr.dataType === DateType )
483-
484- checkFilterPredicate(dateAttr.isNull, classOf [Eq [_]], Seq .empty[Row ])
485- checkFilterPredicate(
486- dateAttr.isNotNull,
487- classOf [NotEq [_]],
488- data.map(i => Row .apply(resultFun(i))))
489-
490- checkFilterPredicate(
491- dateAttr === " 1000-01-01" .date,
492- classOf [Eq [_]],
493- resultFun(" 1000-01-01" ))
494- logWarning(s " java8Api: $java8Api, rebaseMode, $rebaseMode" )
495- checkFilterPredicate(
496- dateAttr <=> " 1000-01-01" .date,
497- classOf [Eq [_]],
498- resultFun(" 1000-01-01" ))
499- checkFilterPredicate(
500- dateAttr =!= " 1000-01-01" .date,
501- classOf [NotEq [_]],
502- Seq (" 2018-03-19" , " 2018-03-20" , " 2018-03-21" ).map(i => Row .apply(resultFun(i))))
503-
504- checkFilterPredicate(
505- dateAttr < " 2018-03-19" .date,
506- classOf [Lt [_]],
507- resultFun(" 1000-01-01" ))
508- checkFilterPredicate(
509- dateAttr > " 2018-03-20" .date,
510- classOf [Gt [_]],
511- resultFun(" 2018-03-21" ))
512- checkFilterPredicate(
513- dateAttr <= " 1000-01-01" .date,
514- classOf [LtEq [_]],
515- resultFun(" 1000-01-01" ))
516- checkFilterPredicate(
517- dateAttr >= " 2018-03-21" .date,
518- classOf [GtEq [_]],
519- resultFun(" 2018-03-21" ))
520-
521- checkFilterPredicate(
522- Literal (" 1000-01-01" .date) === dateAttr,
523- classOf [Eq [_]],
524- resultFun(" 1000-01-01" ))
525- checkFilterPredicate(
526- Literal (" 1000-01-01" .date) <=> dateAttr,
527- classOf [Eq [_]],
528- resultFun(" 1000-01-01" ))
529- checkFilterPredicate(
530- Literal (" 2018-03-19" .date) > dateAttr,
531- classOf [Lt [_]],
532- resultFun(" 1000-01-01" ))
533- checkFilterPredicate(
534- Literal (" 2018-03-20" .date) < dateAttr,
535- classOf [Gt [_]],
536- resultFun(" 2018-03-21" ))
537- checkFilterPredicate(
538- Literal (" 1000-01-01" .date) >= dateAttr,
539- classOf [LtEq [_]],
540- resultFun(" 1000-01-01" ))
541- checkFilterPredicate(
542- Literal (" 2018-03-21" .date) <= dateAttr,
543- classOf [GtEq [_]],
544- resultFun(" 2018-03-21" ))
545-
546- checkFilterPredicate(
547- ! (dateAttr < " 2018-03-21" .date),
548- classOf [GtEq [_]],
549- resultFun(" 2018-03-21" ))
550- checkFilterPredicate(
551- dateAttr < " 2018-03-19" .date || dateAttr > " 2018-03-20" .date,
552- classOf [Operators .Or ],
553- Seq (Row (resultFun(" 1000-01-01" )), Row (resultFun(" 2018-03-21" ))))
554-
555- Seq (3 , 20 ).foreach {
556- threshold =>
557- withSQLConf(
558- SQLConf .PARQUET_FILTER_PUSHDOWN_INFILTERTHRESHOLD .key -> s " $threshold" ) {
559- checkFilterPredicate(
560- In (
561- dateAttr,
562- Array (
563- " 2018-03-19" .date,
564- " 2018-03-20" .date,
565- " 2018-03-21" .date,
566- " 2018-03-22" .date).map(Literal .apply)),
567- if (threshold == 3 ) classOf [Operators .In [_]] else classOf [Operators .Or ],
568- Seq (
569- Row (resultFun(" 2018-03-19" )),
570- Row (resultFun(" 2018-03-20" )),
571- Row (resultFun(" 2018-03-21" )))
572- )
573- }
574- }
575- }
576- }
577- }
578- }
579- }
580384}
0 commit comments