【问题标题】:Pyspark Schema update/alter DataframePyspark Schema 更新/更改数据框
【发布时间】:2020-05-13 17:24:06
【问题描述】:

我需要从 S3 读取一个 csv 文件,它具有字符串、双精度数据,但我将读取为字符串,它将提供仅包含字符串的动态框架。我想为每一行做下面

  1. 连接几列并创建新列
  2. 添加新列
  3. 将第 3 列中的值从字符串转换为日期
  4. 将第 4、5、6 列的值分别从字符串转换为十进制
Storename,code,created_date,performancedata,accumulateddata,maxmontlydata
GHJ 0,GHJ0000001,2020-03-31,0015.5126-,0024.0446-,0017.1811-
MULT,C000000001,2020-03-31,0015.6743-,0024.4533-,0018.0719-

下面是我目前写的代码

def ConvertToDec(myString):
    pattern = re.compile("[0-9]{0,4}[\\.]?[0-9]{0,4}[-]?")
    myString=myString.strip()
    doubleVal="";
    if  myString and  not pattern.match(myString):
       doubleVal=-9999.9999;
    else:
     doubleVal=-Decimal(myString);
    return doubleVal

def rowwise_function(row): 
    row_dict = row.asDict()
    data='d';
    if not row_dict['code']:
        data=row_dict['code']
    else:  
        data='CD'
    if not row_dict['performancedata']:
        data= data +row_dict['performancedata']
    else:  
        data=data + 'HJ'
    // new columns
    row_dict['LC_CODE']=data
    row_dict['CD_CD']=123
    row_dict['GBL']=123.345
    if rec["created_date"]:
        rec["created_date"]= convStr =datetime.datetime.strptime(rec["created_date"], '%Y-%m-%d')
    if rec["performancedata"]
        rec["performancedata"] = ConvertToDec(rec["performancedata"])

    newrow = Row(**row_dict)
    return newrow

store_df = spark.read.option("header","true").csv("C:\\STOREDATA.TXT", sep="|")
ratings_rdd = store_df.rdd
ratings_rdd_new = ratings_rdd.map(lambda row: rowwise_function(row))
updatedDF=spark.createDataFrame(ratings_rdd_new)

基本上,我正在创建几乎新的 DataFrame。我的问题如下 -

  1. 这是正确的方法吗?
  2. 因为我是我不断变化的架构,主要是有没有其他方法

【问题讨论】:

  • 似乎是一个奇怪的赏金问题。

标签: dataframe pyspark


【解决方案1】:

使用 Spark dataframes/sql,为什么要使用 rdd?您不需要执行任何低级别的数据操作,所有操作都是列级别的,因此数据帧更容易/更高效地使用。

创建新列 - .withColumn(<col_name>, <expression/value>) (refer) 所有的if都可以.filter (refer)

整个ConvertToDec 可以使用strip 和ast 模块或float 更好地编写。

【讨论】:

  • 如何更改现有列?像 ConverttoDec 我该怎么做?如果我不应该使用 row_wise 函数
  • 您可以使用:df['col_name'] 选择一列并进行必要的更改。我建议使用 withColumns 以便您同时拥有旧列和新列,最后只需添加一个选择来选择所需的列。
  • 如果您更熟悉 sql,那么您可以创建一个视图并使用它。Additional reference
  • 感谢您的回答,这很有帮助,您的意思是说我可以使用“withColumn”更改现有列和新列?你能举例说明 withColumn 进动的十进制解析吗
  • 按照您的方式进行行迭代,效率会很低。它是执行数据框列操作的长形式。使用withColumn 通过执行操作来添加新列,然后只需使用.select 最后选择所需的列...示例代码:decimal check example,检查并根据需要替换函数.. string to decimal
猜你喜欢
  • 1970-01-01
  • 2019-06-18
  • 2018-01-09
  • 1970-01-01
  • 1970-01-01
  • 2019-05-11
  • 1970-01-01
  • 2016-03-08
  • 1970-01-01
相关资源
最近更新 更多