【发布时间】:2021-08-16 14:30:12
【问题描述】:
我有一个从文本文件中提取的单列数据框。所以从这里:
oneColDF = (spark.read
.format("text")
.load(file_path))
display(oneColDF)
,导致这个df:
大约有一百行。这种格式没有可靠的分隔符(例如,空格不起作用,因为某些字段中有空格),但是,列是固定宽度的,所以我知道每个字段的列名和宽度(所有字段都是字符串)。我知道这一点是因为我得到了字典:
fixed_width_column_defs = {
"submitted_at": (1, 15),
"order_id": (16, 40),
"customer_id": (56, 40),
"sales_rep_id": (96, 40),
"sales_rep_ssn": (136, 15),
"sales_rep_first_name": (151, 15),
"sales_rep_last_name": (166, 15),
"sales_rep_address": (181, 40),
"sales_rep_city": (221, 20),
"sales_rep_state": (241, 2),
"sales_rep_zip": (243, 5),
"shipping_address_attention": (248, 30),
"shipping_address_address": (278, 40),
"shipping_address_city": (318, 20),
"shipping_address_state": (338, 2),
"shipping_address_zip": (340, 5),
"product_id": (345, 40),
"product_quantity": (385, 5),
"product_sold_price": (390, 20)
}
所以我可以像这样添加空列:
multiColDF= oneColDF.withColumn('submitted_at',
lit(None).cast(StringType())).withColumn('order_id', lit(None).cast(StringType())).withColumn('customer_id', lit(None).cast(StringType())).withColumn('sales_rep_id', lit(None).cast(StringType())).withColumn('sales_rep_ssn', lit(None).cast(StringType())).withColumn('sales_rep_first_name', lit(None).cast(StringType())).withColumn('sales_rep_last_name', lit(None).cast(StringType())).withColumn('sales_rep_address', lit(None).cast(StringType())).withColumn('sales_rep_city', lit(None).cast(StringType())).withColumn('sales_rep_state', lit(None).cast(StringType())).withColumn('sales_rep_zip', lit(None).cast(StringType())).withColumn('shipping_address_attention', lit(None).cast(StringType())).withColumn('shipping_address_address', lit(None).cast(StringType())).withColumn('shipping_address_city', lit(None).cast(StringType())).withColumn('shipping_address_state', lit(None).cast(StringType())).withColumn('shipping_address_zip', lit(None).cast(StringType())).withColumn('product_id', lit(None).cast(StringType())).withColumn('product_quantity', lit(None).cast(StringType())).withColumn('product_sold_price', lit(None).cast(StringType()))
,所以现在数据框包含所有列:
所以我想弄清楚如何遍历数据框以使用值列中的适当数据更新所有列。
for row in multiColDF.rdd.collect():
addTheDatafromValColToEachCol()
我被困在这一点上。我将不胜感激任何想法,无论是建立在我所做的事情上,还是建立在一个更简单的解决方案之上。谢谢。
【问题讨论】:
-
建议将其重命名为“使用 spark sql 访问固定字段数据文件”。老前辈会明白的。
标签: python apache-spark pyspark