【发布时间】: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