【发布时间】:2020-09-22 08:40:22
【问题描述】:
我必须将 Scala 代码转换为 python。
scala 代码将字符串的 RDD 转换为 case-class 的 RDD。代码如下:
case class Stock(
stockName: String,
dt: String,
openPrice: Double,
highPrice: Double,
lowPrice: Double,
closePrice: Double,
adjClosePrice: Double,
volume: Double
)
def parseStock(inputRecord: String, stockName: String): Stock = {
val column = inputRecord.split(",")
Stock(
stockName,
column(0),
column(1).toDouble,
column(2).toDouble,
column(3).toDouble,
column(4).toDouble,
column(5).toDouble,
column(6).toDouble)
}
def parseRDD(rdd: RDD[String], stockName: String): RDD[Stock] = {
val header = rdd.first
rdd.filter((data) => {
data(0) != header(0) && !data.contains("null")
})
.map(data => parseStock(data, stockName))
}
是否可以在 PySpark 中实现这一点?我尝试使用以下代码,它给出了错误
from dataclasses import dataclass
@dataclass(eq=True,frozen=True)
class Stock:
stockName : str
dt: str
openPrice: float
highPrice: float
lowPrice: float
closePrice: float
adjClosePrice: float
volume: float
def parseStock(inputRecord, stockName):
column = inputRecord.split(",")
return Stock(stockName,
column[0],
column[1],
column[2],
column[3],
column[4],
column[5],
column[6])
def parseRDD(rdd, stockName):
header = rdd.first()
res = rdd.filter(lambda data : data != header).map(lambda data : parseStock(data, stockName))
return res
错误 Py4JJavaError:调用 z:org.apache.spark.api.python.PythonRDD.collectAndServe 时出错。 :org.apache.spark.SparkException:作业因阶段失败而中止:阶段 21.0 中的任务 0 失败 1 次,最近一次失败:阶段 21.0 中丢失任务 0.0(TID 31,本地主机,执行程序驱动程序):org.apache.spark .api.python.PythonException: Traceback(最近一次调用最后一次):
文件“/content/spark-2.4.5-bin-hadoop2.7/python/lib/pyspark.zip/pyspark/worker.py”,第 364 行,在 main
func,分析器,反序列化器,序列化器 = read_command(pickleSer,infile)
文件“/content/spark-2.4.5-bin-hadoop2.7/python/lib/pyspark.zip/pyspark/worker.py”,第 69 行,在 read_command
command = serializer._read_with_length(file)
_read_with_length 中的文件“/content/spark-2.4.5-bin-hadoop2.7/python/lib/pyspark.zip/pyspark/serializers.py”,第 173 行
返回 self.loads(obj)
加载中的文件“/content/spark-2.4.5-bin-hadoop2.7/python/lib/pyspark.zip/pyspark/serializers.py”,第 587 行
返回pickle.loads(obj,编码=编码)
AttributeError: Can't get attribute 'main' on
【问题讨论】:
标签: apache-spark pyspark rdd case-class python-dataclasses