【问题标题】:How Spark Structured Streaming handles backpressure?Spark Structured Streaming 如何处理背压?
【发布时间】:2023-03-31 19:25:01
【问题描述】:

我正在分析 Spark Structured Streaming 的背压功能。有谁知道细节?是否可以通过代码调整处理传入记录? 谢谢

【问题讨论】:

  • 你如何定义背压?
  • 我的意思是,动态管理记录摄取率的功能。如果您使用 Kafka,则可以在 Spark Streaming 上激活并且可以在 kafka.maxRatePerPartition 上工作。那么结构化流媒体呢?它在内部是如何运作的?程序员可以管理吗?

标签: apache-spark spark-structured-streaming backpressure


【解决方案1】:

如果您的意思是动态更改结构化流中每个内部批次的大小,那么否。结构化流中没有基于接收器的源,因此完全没有必要。从另一个角度来看,Structured Streaming 不能做真正的背压,因为比如 Spark 不能告诉其他应用程序减慢将数据推送到 Kafka 的速度。

一般而言,结构化流式处理会在默认情况下尝试尽可能快地处理数据。每个源中都有允许控制处理速率的选项,例如文件源中的maxFilesPerTrigger 和Kafka 源中的maxOffsetsPerTrigger。阅读以下链接了解更多详情:

http://spark.apache.org/docs/latest/structured-streaming-programming-guide.html#input-sources http://spark.apache.org/docs/latest/structured-streaming-kafka-integration.html

【讨论】:

    【解决方案2】:

    只有基于推送的机制才需要处理背压。 Kafka 消费者是基于拉取的,只有在当前批次完成处理和保存后,Spark 才会拉取下一批记录。如果处理和保存在 spark 中延迟,它不会拉新一批记录,因此不需要背压处理。

    maxOffsetsPerTrigger 可以更改每个 spark 批处理集处理的记录数,backpressure.enabled 更改接收率,但这与您去告诉源减速的背压不同。

    【讨论】:

    猜你喜欢
    • 2021-05-22
    • 1970-01-01
    • 2020-01-05
    • 2019-06-25
    • 2020-09-03
    • 1970-01-01
    • 2020-03-19
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多