尝试在目录中使用 for loop over all files 并仅从文件中读取所需的列。
Example:
#files path list
file_lst=['<path1>','<path2>']
from pyspark.sql.functions import *
from pyspark.sql.types import *
#define schema for the required columns
schema = StructType([StructField("column1",StringType(),True),StructField("column2",StringType(),True)])
#create an empty dataframe
df=spark.createDataFrame([],schema)
#loop through files with reading header from the file then select only req cols
#union all dataframes
for i in file_lst:
tmp_df=spark.read.option("header","true").csv(i).select("column1","column2")
df=df.unionAll(tmp_df)
#display results
df.show()
如果您的目录中的文件在所有文件中以特定顺序的column1,column2,column3..etc(required columns),那么您可以尝试如下:
spark.read.option("header","true").csv("<directory>").select("column1","column2","column3").show()