【问题标题】:New to Pyspark - importing a CSV and creating a parquet file with array columnsPyspark 新手 - 导入 CSV 并创建包含数组列的 parquet 文件
【发布时间】:2019-08-12 23:28:57
【问题描述】:

我是 Pyspark 的新手,我一直在努力完成我认为相当简单的事情。我正在尝试执行 ETL 过程,将 csv 文件转换为 parquet 文件。 CSV 文件有几个简单的列,但其中一列是分隔的整数数组,我想将其展开/解压缩到镶木地板文件中。这个 parquet 文件实际上由 .net 核心微服务使用,该微服务使用 Parquet Reader 进行下游计算。为了让这个问题简单,该列的结构是:

"geomap" 5:3:7|4:2:1|8:2:78 -> 这表示一个包含 3 个项目的数组,它在“|”处分割然后由值 (5,3,7), (4,2,1), (8,2,78) 构建一个元组

我尝试了各种流程和架构,但无法正确解决。通过 UDF,我正在创建列表列表或元组列表,但我无法正确获取架构或将数据解压缩到镶木地板写入操作中。我要么得到空值、错误或其他问题。我需要以不同的方式处理这个问题吗?相关代码如下。为简单起见,我只是显示问题列,因为我还有其他工作。这是我第一次尝试 Pyspark,因此很抱歉遗漏了一些明显的东西:

def convert_geo(geo):
   return [tuple(x.split(':')) for x in geo.split('|')]

compression_type = 'snappy'

schema = ArrayType(StructType([
    StructField("c1", IntegerType(), False),
    StructField("c2", IntegerType(), False),
    StructField("c3", IntegerType(), False)
]))

spark_convert_geo = udf(lambda z: convert_geo(z),schema)

source_path = '...path to csv'
destination_path = 'path for generated parquet file'

df = spark.read.option('delimiter',',').option('header','true').csv(source_path).withColumn("geomap",spark_convert_geo(col('geomap')).alias("geomap"))
df.write.mode("overwrite").format('parquet').option('compression', compression_type).save(destination_path)

编辑:根据添加 printSchema() 输出的请求,我也不确定这里有什么问题。我似乎仍然无法正确显示或渲染字符串拆分值。这包含所有列。我确实看到了 c1 和 c2 以及 c3 结构名称...

root |-- lrsegid: integer (nullable = true) |-- loadsourceid: integer (nullable = true) |-- agencyid: integer (nullable = true) |-- acres: float (nullable = true) |-- sourcemap: array (nullable = true) | |-- element: integer (containsNull = true) |-- geomap: array (nullable = true) | |-- element: struct (containsNull = true) | | |-- c1: integer (nullable = false) | | |-- c2: integer (nullable = false) | | |-- c3: integer (nullable = false) 

【问题讨论】:

  • 你能把df.printSchema的输出贴出来
  • 当然,我已经使用 printSchema() 的输出编辑了帖子。它包含了我为简单起见而省略的所有其他列。

标签: python apache-spark dataframe pyspark parquet


【解决方案1】:

问题在于convert_geo 函数返回一个包含字符元素的元组列表,而不是架构中指定的整数。如果您按以下方式进行修改,它将起作用:

def convert_geo(geo):
    return [tuple([int(y) for y in x.split(':')]) for x in geo.split('|')]

【讨论】:

  • 我可以发誓我尝试将架构全部设为 String type(),但它仍然无法正常工作。让我再检查一下。另外,如果 Tuple 中的第三项实际上需要是 double 怎么办?如何编辑 UDF 以使第三项成为不同的值类型?
  • 上述调整对我有用。如果你想要结构元素的不同数据类型,你可以用 for 循环和一些条件逻辑替换列表理解
  • 您说得对。标记为已回答。我在玩一堆不同的模式和结构,我一定从未尝试过将值类型与正确的模式定义匹配。我在 Python 中做的不多,我不确定列表理解是否有办法在创建时混合值类型。我假设一个 for 循环可能会慢一点?我想这取决于内部列表理解的实现。但是,是的,我的 parquet 文件现在有 3 个结构、相同的长度、相同的重复级别和正确的数据。
  • 感谢您的回答。考虑一下,您可能可以通过使用与元组长度相同的类型转换函数列表来避免 for 循环(然后用冒号拆分列表压缩并像以前一样使用列表推导)
  • 嗯,也许吧。你碰巧有这样的例子吗?如果没有,不用担心,我可以稍微考虑一下这个想法。再次感谢你的帮助。我是一名 .net 开发人员,在我的 python 中有些生疏。我相信我能弄明白。
猜你喜欢
  • 2019-12-06
  • 2018-05-28
  • 1970-01-01
  • 2016-10-22
  • 2023-02-21
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多