【问题标题】:CodeGen grows beyond 64 KB error when normalizing large PySpark dataframe规范化大型 PySpark 数据帧时,CodeGen 出现超过 64 KB 的错误
【发布时间】:2017-04-27 04:57:48
【问题描述】:

我有一个包含 1300 万行和 800 列的 PySpark 数据框。我需要规范化这些数据,所以一直在使用这个代码,它适用于较小的开发数据集。

def z_score_w(col, w):
    avg_ = avg(col).over(w)
    stddev_ = stddev_pop(col).over(w)
    return (col - avg_) / stddev_

w = Window().partitionBy().rowsBetween(-sys.maxsize, sys.maxsize)    
norm_exprs = [z_score_w(signalsDF[x], w).alias(x) for x in signalsDF.columns]

normDF = signalsDF.select(norm_exprs)

但是,在使用完整数据集时,我遇到了 codegen 异常:

        at org.apache.spark.sql.catalyst.expressions.codegen.CodeGenerator$.org$apache$spark$sql$catalyst$expressions$codegen$CodeGenerator$$doCompile(CodeGenerator.scala:893
)
        at org.apache.spark.sql.catalyst.expressions.codegen.CodeGenerator$$anon$1.load(CodeGenerator.scala:950)
        at org.apache.spark.sql.catalyst.expressions.codegen.CodeGenerator$$anon$1.load(CodeGenerator.scala:947)
        at org.spark_project.guava.cache.LocalCache$LoadingValueReference.loadFuture(LocalCache.java:3599)
        at org.spark_project.guava.cache.LocalCache$Segment.loadSync(LocalCache.java:2379)
        ... 44 more
Caused by: org.codehaus.janino.JaninoRuntimeException: Code of method "(Lorg/apache/spark/sql/catalyst/expressions/GeneratedClass;[Ljava/lang/Object;)V" of class "org.apache.
spark.sql.catalyst.expressions.GeneratedClass$SpecificMutableProjection" grows beyond 64 KB
        at org.codehaus.janino.CodeContext.makeSpace(CodeContext.java:941)
        at org.codehaus.janino.CodeContext.write(CodeContext.java:836)
        at org.codehaus.janino.UnitCompiler.writeOpcode(UnitCompiler.java:10251)
        at org.codehaus.janino.UnitCompiler.pushConstant(UnitCompiler.java:8933)
        at org.codehaus.janino.UnitCompiler.compileGet2(UnitCompiler.java:4346)
        at org.codehaus.janino.UnitCompiler.access$7100(UnitCompiler.java:185)
        at org.codehaus.janino.UnitCompiler$10.visitBooleanLiteral(UnitCompiler.java:3267)

周围有几个Spark JIRA issues 看起来相似,但这些都标记为已解决。还有 this SO question 是相关的,但答案是另一种技术。

我有自己的解决方法,可以标准化数据框的成批列。这可行,但我最终得到了多个数据帧,然后我必须加入,这很慢。

所以,我的问题是 - 是否有另一种技术来规范我缺少的大型数据帧?

我正在使用 spark-2.0.1。

【问题讨论】:

    标签: apache-spark pyspark apache-spark-sql pyspark-sql window-functions


    【解决方案1】:

    请查看此链接,我们通过在代码中添加检查点解决了此错误。

    检查点只是将数据或数据帧写回磁盘并读回。

    https://stackoverflow.com/a/55208567/7241837

    检查点详情

    https://github.com/JerryLead/SparkInternals/blob/master/markdown/english/6-CacheAndCheckpoint.md

    问:什么样的RDD需要checkpoint?

    the computation takes a long time
    the computing chain is too long
    depends too many RDDs
    

    其实把ShuffleMapTask的输出保存到本地磁盘也是checkpoint,不过只是partition的数据输出而已。

    问:什么时候检查点?

    如上所述,每次需要缓存计算分区时,都会将其缓存到内存中。但是,检查点不遵循相同的原则。相反,它会等到一个作业结束,然后启动另一个作业来完成检查点。需要检查点的 RDD 将被计算两次;因此建议在 rdd.checkpoint() 之前执行 rdd.cache()。在这种情况下,第二个作业不会重新计算 RDD。相反,它只会读取缓存。实际上,Spark 提供了 rdd.persist(StorageLevel.DISK_ONLY) 方法,就像在磁盘上缓存一样。因此,它在第一次计算时将 RDD 缓存在磁盘上,但是这种持久化和检查点是不同的,我们稍后会讨论区别。

    问:如何实现检查点?

    程序如下:

    RDD 将是:[ 已初始化 --> 标记为检查点 --> 正在进行检查点-> 检查点]。最后,它将是 检查点。

    数据帧类似:将数据帧写入磁盘或 s3 并在新数据帧中读回数据。

    初始化

    在驱动端,调用 rdd.checkpoint() 后,RDD 将由 RDDCheckpointData 管理。用户应设置检查点的存储路径(在 hdfs 上)。

    标记为检查点

    初始化后,RDDCheckpointData会标记RDD MarkedForCheckpoint。

    检查点正在进行中

    当作业完成时,将调用 finalRdd.doCheckpoint()。 finalRDD 向后扫描计算链。当遇到需要检查点的 RDD 时,会将 RDD 标记为 CheckpointingInProgress,然后将配置文件(用于写入 hdfs)如 core-site.xml 广播到其他工作节点的 blockManager。之后,将启动一个作业来完成检查点:

      rdd.context.runJob(rdd, CheckpointRDD.writeToFile(path.toString,  broadcastedConf))
    

    检查点

    job 完成 checkpoint 后,会清理 RDD 的所有依赖,并将 RDD 设置为 checkpointed。然后,添加一个补充依赖,并将父 RDD 设置为 CheckpointRDD。 checkpointRDD将来会用于从文件系统中读取checkpoint文件,然后生成RDD分区

    有趣的是:

    两个 RDD 在驱动程序中被检查点,但只有结果(见下面的代码)被成功检查点。不确定这是一个错误还是只是下游 RDD 会被有意设置检查点。

    val data1 = Array[(Int, Char)]((1, 'a'), (2, 'b'), (3, 'c'),
        (4, 'd'), (5, 'e'), (3, 'f'), (2, 'g'), (1, 'h'))
       val pairs1 = sc.parallelize(data1, 3)
    
       val data2 = Array[(Int, Char)]((1, 'A'), (2, 'B'), (3, 'C'), (4, 'D'))
       val pairs2 = sc.parallelize(data2, 2)
    
       pairs2.checkpoint
    
       val result = pairs1.join(pairs2)
       result.checkpoint
    

    【讨论】:

      【解决方案2】:

      一个明显的问题是您使用窗口函数的方式。以下框架:

      Window().partitionBy().rowsBetween(-sys.maxsize, sys.maxsize)    
      

      在实践中有点没用。如果没有分区列,它首先将所有数据重新洗牌到单个分区。这种缩放方法仅适用于按组执行缩放。

      Spark 提供了两个可用于扩展特征的类:

      • pyspark.ml.feature.StandardScaler
      • pyspark.mllib.feature.StandardScaler

      不幸的是,两者都需要Vector 数据作为输入。使用机器学习

      from pyspark.ml.feature import StandardScaler as MLScaler, VectorAssembler
      from pyspark.ml import Pipeline
      
      scaled = Pipeline(stages=[
          VectorAssembler(inputCols=df.columns, outputCol="features"), 
          MLScaler(withMean=True, inputCol="features", outputCol="scaled")
      ]).fit(df).transform(df).select("scaled")
      

      如果您需要原始形状,这需要进一步扩展 scaled 列。

      使用 MLlib:

      from pyspark.mllib.feature import StandardScaler as MLLibScaler
      from pyspark.mllib.linalg import DenseVector
      
      rdd = df.rdd.map(DenseVector)
      scaler = MLLibScaler(withMean=True, withStd=True)
      
      scaler.fit(rdd).transform(rdd).map(lambda v: v.array.tolist()).toDF(df.columns)
      

      如果存在与列数相关的代码生成问题,后一种方法可能更有用。

      另一种方法可以解决这个问题来计算全局统计数据

      from pyspark.sql.functions import avg, col, stddev_pop, struct
      
      stats = df.agg(*[struct(avg(c), stddev_pop(c)) for c in df.columns]).first()
      

      然后选择:

      df.select(*[
          ((col(c) - mean) / std).alias(c)
          for (c, (mean, std)) in zip(df.columns, stats)
      ])
      

      按照您的 cmets,您认为可以使用 NumPy 和一些基本转换来表达的最简单解决方案:

      rdd = df.rdd.map(np.array)  # Convert to RDD of NumPy vectors
      stats = rdd.stats()  # Compute mean and std
      scaled = rdd.map(lambda v: (v - stats.mean()) / stats.stdev())  # Normalize
      

      并转换回DataFrame:

      scaled.map(lambda x: x.tolist()).toDF(df.columns)
      

      【讨论】:

        猜你喜欢
        • 2018-11-26
        • 2020-11-15
        • 1970-01-01
        • 2017-08-30
        • 2018-02-12
        • 2021-11-13
        • 2018-09-15
        • 2018-08-10
        • 1970-01-01
        相关资源
        最近更新 更多