【发布时间】:2020-05-13 17:24:06
【问题描述】:
我需要从 S3 读取一个 csv 文件,它具有字符串、双精度数据,但我将读取为字符串,它将提供仅包含字符串的动态框架。我想为每一行做下面
- 连接几列并创建新列
- 添加新列
- 将第 3 列中的值从字符串转换为日期
- 将第 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。我的问题如下 -
- 这是正确的方法吗?
- 因为我是我不断变化的架构,主要是有没有其他方法
【问题讨论】:
-
似乎是一个奇怪的赏金问题。