我已经创建了一个空的dataframe,并开始添加它,通过读取每个文件。但其中一个文件的列数比前一个文件多。如何仅为所有其他文件选择第一个文件中的列?
from pyspark.sql import SparkSession
from pyspark.sql import SQLContext
from pyspark.sql.types import StructType
import os, glob
spark = SparkSession.builder.\
config("spark.jars.packages","saurfang:spark-sas7bdat:2.0.0-s_2.11")\
.enableHiveSupport().getOrCreate()
fpath=''
schema = StructType([])
sc = spark.sparkContext
df_spark=spark.createDataFrame(sc.emptyRDD(), schema)
files=glob.glob(fpath +'*.sas7bdat')
for i,f in enumerate(files):
if i == 0:
df=spark.read.format('com.github.saurfang.sas.spark').load(f)
df_spark= df
else:
df=spark.read.format('com.github.saurfang.sas.spark').load(f)
df_spark=df_spark.union(df)发布于 2019-09-07 06:45:13
您可以在创建dataframe时提供自己的架构。例如,我有两个文件emp1.csv & emp2.csv具有不同的模式。
id,empname,empsalary
1,Vikrant,55550
id,empname,empsalary,age,country
2,Raghav,10000,32,India
schema = StructType([
StructField("id", IntegerType(), True),
StructField("name", StringType(), True),
StructField("salary", IntegerType(), True)])
file_path="file:///home/vikct001/user/vikrant/inputfiles/testfiles/emp*.csv"
df=spark.read.format("com.databricks.spark.csv").option("header", "true").schema(schema).load(file_path)指定模式不仅可以解决数据类型和格式问题,而且还需要提高性能。
如果需要删除格式错误的记录,还有其他选项,但这也会删除具有空值或不符合所提供模式的记录。它可以跳过那些记录,也可以具有多个分隔符和垃圾字符或一个空文件。
.option("mode", "DROPMALFORMED")FAILFAST模式会在发现格式错误的记录时抛出异常。
.option("mode", "FAILFAST")您还可以使用map函数来选择您选择的元素,并在构建dataframe时排除其他元素。
df=spark.read.format('com.databricks.spark.csv').option("header", "true").load(file_path).rdd.map(lambda x :(x[0],x[1],x[2])).toDF(["id","name","salary"])在这两种情况下,您都需要将头设置为“true”,否则它将包含您的csv头作为数据文件的第一条记录。
发布于 2019-09-07 01:35:39
您可以从第一个文件的架构中获取字段名,然后使用字段名数组从所有其他文件中选择列。
fields = df.schema.fieldNames可以使用字段数组从所有其他数据集中选择列。下面是scala代码。
df=spark.read.format('com.github.saurfang.sas.spark').load(f).select(fields(0),fields.drop(1):_*)https://stackoverflow.com/questions/57824016
复制相似问题