【问题标题】:Converting One Column (Fixed-Field-Width) Dataframe to Multicolumn Dataframe (Databricks, pyspark )将一列(固定字段宽度)数据帧转换为多列数据帧(Databricks,pyspark)
【发布时间】: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


【解决方案1】:

您不需要在之前添加列,您可以使用 substring 动态执行此操作并使用 select 循环遍历字典:

from pyspark.sql import functions as F
out = df.select("value",*[F.substring("value",*v).alias(k) 
                     for k,v in fixed_width_column_defs.items()])

【讨论】:

  • 谢谢!那成功了。我正在尝试理解“*v”?
  • @TimothyClotworthy 子字符串采用 3 个参数,1:col,2:position 3:length,我们有 col,但没有其他 2,因此我们将值放在 v 中,这是一个要填充的元组参数 2 和参数 3,尝试打印[*(2,3)],你会看到元组已经被解压到容器中,这里是一个列表。这可能会有所帮助:)
  • 对不起,从上面代码的上下文中,我如何打印出 [*(2,3)]?
  • @TimothyClotworthy 抱歉,我的意思是独立打印以理解:print([*(2,3)]) 甚至 print(*(2,3)) 以查看使用 * 的元组的解包。更多细节在这里:stackoverflow.com/questions/400739/…
  • hmm,所以我对我认为是另一个拆包感到困惑(这个:*[F.substring("value",*v).alias(k) for k,v in fixed_width_column_defs.项目()])?因为它在一个相当复杂的表达式前面有星号......谢谢!
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2017-09-18
  • 2021-11-16
  • 1970-01-01
  • 2021-12-21
  • 2023-01-21
  • 1970-01-01
相关资源
最近更新 更多