Skip to content

Commit 153293e

Browse files
committed
fix unit test
1 parent 8426522 commit 153293e

1 file changed

Lines changed: 24 additions & 20 deletions

File tree

sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/FileSourceStrategy.scala

Lines changed: 24 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -22,9 +22,11 @@ import scala.collection.mutable.ArrayBuffer
2222
import org.apache.hadoop.fs.{BlockLocation, FileStatus, LocatedFileStatus, Path}
2323

2424
import org.apache.spark.internal.Logging
25+
import org.apache.spark.rdd.RDD
2526
import org.apache.spark.sql._
26-
import org.apache.spark.sql.catalyst.expressions
27+
import org.apache.spark.sql.catalyst.{expressions, InternalRow}
2728
import org.apache.spark.sql.catalyst.expressions._
29+
import org.apache.spark.sql.catalyst.expressions.codegen.GenerateUnsafeProjection
2830
import org.apache.spark.sql.catalyst.planning.PhysicalOperation
2931
import org.apache.spark.sql.catalyst.plans.logical.LogicalPlan
3032
import org.apache.spark.sql.execution.DataSourceScanExec
@@ -111,25 +113,27 @@ private[sql] object FileSourceStrategy extends Strategy with Logging {
111113

112114
val optimizerMetadataOnly =
113115
readDataColumns.isEmpty && files.sparkSession.sessionState.conf.optimizerMetadataOnly
114-
val scanRdd = if (optimizerMetadataOnly) {
115-
val partitionValues = selectedPartitions.map(_.values)
116-
files.sqlContext.sparkContext.parallelize(partitionValues, 1)
117-
} else {
118-
val readFile = files.fileFormat.buildReaderWithPartitionValues(
119-
sparkSession = files.sparkSession,
120-
dataSchema = files.dataSchema,
121-
partitionSchema = files.partitionSchema,
122-
requiredSchema = prunedDataSchema,
123-
filters = pushedDownFilters,
124-
options = files.options,
125-
hadoopConf = files.sparkSession.sessionState.newHadoopConfWithOptions(files.options))
126-
127-
val plannedPartitions = getFilePartitions(files, selectedPartitions)
128-
new FileScanRDD(
129-
files.sparkSession,
130-
readFile,
131-
plannedPartitions)
132-
}
116+
val scanRdd: RDD[InternalRow] = if (optimizerMetadataOnly) {
117+
val partitionSchema = files.partitionSchema.toAttributes
118+
lazy val converter = GenerateUnsafeProjection.generate(partitionSchema, partitionSchema)
119+
val partitionValues = selectedPartitions.map(_.values)
120+
files.sqlContext.sparkContext.parallelize(partitionValues, 1).map(converter(_))
121+
} else {
122+
val readFile = files.fileFormat.buildReaderWithPartitionValues(
123+
sparkSession = files.sparkSession,
124+
dataSchema = files.dataSchema,
125+
partitionSchema = files.partitionSchema,
126+
requiredSchema = prunedDataSchema,
127+
filters = pushedDownFilters,
128+
options = files.options,
129+
hadoopConf = files.sparkSession.sessionState.newHadoopConfWithOptions(files.options))
130+
131+
val plannedPartitions = getFilePartitions(files, selectedPartitions)
132+
new FileScanRDD(
133+
files.sparkSession,
134+
readFile,
135+
plannedPartitions)
136+
}
133137

134138
val meta = Map(
135139
"Format" -> files.fileFormat.toString,

0 commit comments

Comments
 (0)