【问题标题】:Creating an Apache Spark RDD of a Class in PySpark在 PySpark 中创建一个类的 Apache Spark RDD
【发布时间】: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 'ma​​in' on

【问题讨论】:

    标签: apache-spark pyspark rdd case-class python-dataclasses


    【解决方案1】:

    Dataset API 不适用于 python。

    “Dataset 是数据的分布式集合。Dataset 是 Spark 1.6 中添加的一个新接口,它提供了 RDD 的优点(强类型化,能够使用强大的 lambda 函数)和 Spark SQL 优化执行引擎的优点。数据集可以从 JVM 对象构建,然后使用功能转换(map、flatMap、过滤器等)进行操作。数据集 API 在 Scala 和 Java 中可用。Python 不支持数据集 API。但是由于鉴于 Python 的动态特性,Dataset API 的许多好处已经可用(即您可以通过名称自然地访问行的字段 row.columnName)。R 的情况类似。”

    https://spark.apache.org/docs/latest/sql-programming-guide.html

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2018-02-14
      • 2015-08-31
      • 1970-01-01
      • 2016-08-08
      • 2020-04-02
      • 2014-05-13
      • 1970-01-01
      • 2015-06-15
      相关资源
      最近更新 更多