【问题标题】:Correctly reading the types from file in PySpark从 PySpark 中的文件中正确读取类型
【发布时间】:2016-03-11 00:23:07
【问题描述】:

我有一个制表符分隔的文件,其中包含以下行

id1 name1   ['a', 'b']  3.0 2.0 0.0 1.0

即一个id,一个名字,一个带有一些字符串的列表,以及一系列4个float属性。 我正在阅读这个文件作为

rdd = sc.textFile('myfile.tsv') \
    .map(lambda row: row.split('\t'))
df = sqlc.createDataFrame(rdd, schema)

我将架构指定为

schema = StructType([
    StructField('id', StringType(), True),
    StructField('name', StringType(), True),
    StructField('list', ArrayType(StringType()), True),
    StructField('att1', FloatType(), True),
    StructField('att2', FloatType(), True),
    StructField('att3', FloatType(), True),
    StructField('att4', FloatType(), True)
])

问题是,从 DataFrame 上的 collect 判断,列表和属性都没有被正确读取。事实上,我得到了所有的None

Row(id=u'id1', brand_name=u'name1', list=None, att1=None, att2=None, att3=None, att4=None)

我做错了什么?

【问题讨论】:

  • 绝对确定所有列都是制表符分隔的? (可能看起来像一个愚蠢的问题,但你永远不会知道)。如果有疑问,请执行 hexdump;一个空格的十六进制代码是 20,一个制表符是 09
  • @jDo 绝对确定并检查过

标签: python apache-spark dataframe pyspark


【解决方案1】:

它被正确阅读,它只是不像你期望的那样工作。 Schema 参数声明了哪些类型是,以避免昂贵的模式推断,而不是如何转换数据。提供与声明的架构匹配的输入是您的责任。

这也可以由数据源处理(查看spark-csvinferSchema 选项)。但它不会处理像数组这样的复杂类型。

由于您的架构大多是扁平的,并且您知道类型,您可以尝试这样的事情:

df = rdd.toDF([f.name for f in schema.fields])

exprs = [
    # You should excluding casting
    # on other complex types as well
    col(f.name).cast(f.dataType) if f.dataType.typeName() != "array" 
    else col(f.name)
    for f in schema.fields
]

df.select(*exprs)

并使用字符串处理函数或 UDF 分别处理复杂类型。或者,由于无论如何您都是在 Python 中读取数据,因此只需在创建 DF 之前强制执行所需的类型。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2016-06-20
    • 2021-08-12
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-09-07
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多