【发布时间】:2020-02-01 06:11:52
【问题描述】:
我有一个尝试使用 Spark 解决的用例。用例是我必须调用一个需要batchSize 和token 的API,然后它会返回下一页的令牌。它给了我一个 JSON 对象的列表。现在我必须调用这个 API,直到所有结果都返回并以 parquet 格式将它们全部写入 s3。返回对象的大小范围为 0 到 1 亿。
我的方法是,我首先得到一批 100 万个对象,我将它们转换为数据集,然后使用写入镶木地板
dataSet.repartition(1).write.mode(SaveMode.Append)
.option("mapreduce.fileoutputcommitter.algorithm.version", "2")
.parquet(s"s3a://somepath/")
然后重复这个过程,直到我的 API 说没有更多数据,即 token 为空
因此,这些 API 调用必须在驱动程序上按顺序运行。一旦我得到一百万,我就会写信给 s3。
我在驱动程序上看到了这些内存问题。
Application application_1580165903122_19411 failed 1 times due to AM Container for appattempt_1580165903122_19411_000001 exited with exitCode: -104
Diagnostics: Container [pid=28727,containerID=container_1580165903122_19411_01_000001] is running beyond physical memory limits. Current usage: 6.6 GB of 6.6 GB physical memory used; 16.5 GB of 13.9 GB virtual memory used. Killing container.
Dump of the process-tree for container_1580165903122_19411_01_000001 :
从某种意义上说,我看到了一些奇怪的行为,有时 3000 万可以正常工作,有时会因此而失败。有时甚至 100 万个也会失败。
我想知道我是否犯了一些非常愚蠢的错误,还是有更好的方法来解决这个问题?
【问题讨论】:
标签: scala apache-spark