【问题标题】:Spark How to write to parquet file from data using synchronous APISpark如何使用同步API从数据写入parquet文件
【发布时间】:2020-02-01 06:11:52
【问题描述】:

我有一个尝试使用 Spark 解决的用例。用例是我必须调用一个需要batchSizetoken 的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


    【解决方案1】:

    这种设计不可扩展,会给驱动程序带来很大压力,因此预计会崩溃。此外,在写入 s3 之前,内存中会累积大量数据。

    我会推荐你​​使用 Spark 流从 API 中读取数据。这样,许多执行器将完成工作,并且解决方案将具有很大的可扩展性。这是一个例子 - RestAPI service call from Spark Streaming

    在这些执行器中,您可以平衡地累积 API 响应,例如累积 20,000 条记录但不等待 5M 条记录。在说 20,000 之后以“追加”模式将它们写入 S3。 “附加”模式将帮助多个进程协同工作,而不是相互踩踏。

    【讨论】:

    • 感谢您的建议。在我的情况下,对 REST API 的第二次调用取决于第一次的响应。火花流如何出现?因为这个实现需要多个执行器根据 API 令牌获取不同的数据
    • 您可以在第一次回复后立即拨打第二次电话吗?如果是这样,则进行 2 次调用,然后将两个响应都写入 s3。如果您不能在第 1 次之后立即进行第 2 次调用,则编写 2 进程 - 一个进行第一次调用并写入响应,另一个从 s3 读取响应并进行第二次调用
    猜你喜欢
    • 1970-01-01
    • 2015-11-27
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-02-24
    • 2017-03-17
    • 2017-09-06
    • 2020-01-11
    相关资源
    最近更新 更多