【问题标题】:OutOfMemoryError while generating 60 million JSON file using PySpark [duplicate]使用 PySpark 生成 6000 万个 JSON 文件时出现 OutOfMemoryError [重复]
【发布时间】:2019-09-07 10:43:08
【问题描述】:

通过 jdbc 连接,我能够使用下面的 PySpark 代码从 Oracle db 成功生成 6000 万条记录 CSV 文件。

然后现在我想要以 JSON 格式输出,所以我添加了这行代码:df1.toPandas().to_json("/home/user1/empdata.json", orient='records'),但是在生成 json 时出现了 OutOfMemoryError。

如果需要任何代码更改,请任何人推荐我。

from pyspark.sql import SparkSession

spark = SparkSession \
    .builder \
    .appName("Emp data Extract") \
    .config("spark.some.config.option", " ") \
    .getOrCreate()

def generateData():
    try:
        jdbcUrl = "jdbc:oracle:thin:USER/pwd@//hostname:1521/dbname"
        jdbcDriver = "oracle.jdbc.driver.OracleDriver"
        df1 = spark.read.format('jdbc').options(url=jdbcUrl, dbtable="(SELECT * FROM EMP) alias1", driver=jdbcDriver, fetchSize="2000").load()
        #df1.coalesce(1).write.format("csv").option("header", "true").save("/home/user1/empdata" , index=False)
        df1.toPandas().to_json("/home/user1/empdata.json", orient='records')
    except Exception as err:
        print(err)
        raise
    # finally:
    # conn.close()

if __name__ == '__main__':
    generateData()

错误日志:

2019-04-15 05:17:06 WARN  Utils:66 - Truncated the string representation of a plan since it was too large. This behavior can be adjusted by setting 'spark.debug.maxToStringFields' in SparkEnv.conf.
[Stage 0:>                                                          (0 + 1) / 1]2019-04-15 05:20:22 ERROR Executor:91 - Exception in task 0.0 in stage 0.0 (TID 0)
java.lang.OutOfMemoryError: Java heap space
        at java.util.Arrays.copyOf(Arrays.java:3236)
        at java.io.ByteArrayOutputStream.grow(ByteArrayOutputStream.java:118)
        at java.io.ByteArrayOutputStream.ensureCapacity(ByteArrayOutputStream.java:93)
        at java.io.ByteArrayOutputStream.write(ByteArrayOutputStream.java:153)
        at net.jpountz.lz4.LZ4BlockOutputStream.flushBufferedData(LZ4BlockOutputStream.java:220)
        at net.jpountz.lz4.LZ4BlockOutputStream.write(LZ4BlockOutputStream.java:173)
        at java.io.DataOutputStream.write(DataOutputStream.java:107)
        at org.apache.spark.sql.catalyst.expressions.UnsafeRow.writeToStream(UnsafeRow.java:552)
        at org.apache.spark.sql.execution.SparkPlan$$anonfun$2.apply(SparkPlan.scala:256)
        at org.apache.spark.sql.execution.SparkPlan$$anonfun$2.apply(SparkPlan.scala:247)
        at org.apache.spark.rdd.RDD$$anonfun$mapPartitionsInternal$1$$anonfun$apply$25.apply(RDD.scala:836)
        at org.apache.spark.rdd.RDD$$anonfun$mapPartitionsInternal$1$$anonfun$apply$25.apply(RDD.scala:836)
        at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:49)
        at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:324)
        at org.apache.spark.rdd.RDD.iterator(RDD.scala:288)
        at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:87)
        at org.apache.spark.scheduler.Task.run(Task.scala:109)
        at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:345)
        at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
        at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
        at java.lang.Thread.run(Thread.java:748)
2019-04-15 05:20:22 ERROR SparkUncaughtExceptionHandler:91 - Uncaught exception in thread Thread[Executor task launch worker for task 0,5,main]
java.lang.OutOfMemoryError: Java heap space

根据管理员的要求,我正在更新我的 cmets:这是一些不同的问题,也存在其他 outoutmemory 问题,但在不同的情况下会出现。错误可能相同,但问题不同。就我而言,我得到了由于大量数据。

【问题讨论】:

  • 文件太大,堆内存无法处理。处理此问题的一种方法是使用缓冲的 IO 流。
  • 好的,感谢您的回复。但同时我能够生成一个包含 60 m 记录的 csv。您可以看到 csv 的注释行。现在只有 json 的问题。
  • 对不起,我对 Spark 的了解有限,我提出了笼统的建议。我认为@Arnon Rotem-Gal-Oz 提到的应该对你有用。

标签: apache-spark pyspark python-3.6


【解决方案1】:

如果你想以 JSON 格式保存,你应该使用 Spark 的 write 命令 - 你目前所做的是将所有数据带到驱动程序并尝试将其加载到 pandas 数据帧中

df1.write.format('json').save('/path/file_name.json')

如果你需要一个文件,你可以试试

df1.coalesce(1).write.format('json').save('/path/file_name.json')

【讨论】:

  • 是的,现在我得到了输出,但是在 json 文件中缺少大括号 [ ] ,它的大小为 20gb。我无法手动更新,是否有任何选项可以包括这些。
  • 我的意思是,json 的开始和结束大括号。
  • 我不认为你可以 - 但如果你以后需要在 pandas 中阅读它,你可以使用 pd.read_json("filenanme.json",lines=True)
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-09-20
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-02-14
相关资源
最近更新 更多