【问题标题】:Stateful udfs in spark sql, or how to obtain mapPartitions performance benefit in spark sql?spark sql 中的 stateful udfs,或者如何在 spark sql 中获得 mapPartitions 性能优势?
【发布时间】:2018-09-08 12:49:12
【问题描述】:

在转换导致创建或加载昂贵的资源(例如 - 对外部服务进行身份验证或创建数据库连接)的情况下,使用 map over map partitions 可以显着提升性能。

mapPartition 允许我们为每个分区初始化一次昂贵的资源,而不是像标准 map 那样每行初始化一次。

但是,如果我使用数据帧,我应用自定义转换的方式是指定用户定义的函数,这些函数逐行操作 - 所以我失去了使用 mapPartitions 对每个块执行一次繁重的工作的能力。

在 spark-sql/dataframe 中有解决方法吗?

更具体

我需要对一堆文档进行特征提取。我有一个输入文档并输出向量的函数。

计算本身涉及初始化与外部服务的连接。我不想或不需要为每个文档初始化它。这在规模上具有不小的开销。

【问题讨论】:

    标签: apache-spark optimization pyspark user-defined-functions


    【解决方案1】:

    一般来说你有三个选择:

    • DataFrame转换为RDD并直接应用mapPartitions。由于您使用 Python udf,因此您已经破坏了某些优化并支付了 serde 成本,并且使用 RDD 平均不会使情况变得更糟。
    • 懒惰initialize required resources(另见How to run a function on all Spark workers before processing data in PySpark?)。
    • 如果可以使用 Arrow 序列化数据,请使用矢量化 pandas_udf(Spark 2.3 及更高版本)。不幸的是,您不能直接将它与VectorUDT 一起使用,因此您必须扩展向量并稍后折叠,因此这里的限制因素是向量的大小。此外,您必须小心控制分区的大小。

    请注意,使用UserDefinedFunctions 可能需要promoting objects to non-deterministic 变体。

    【讨论】:

    • 我担心没有直接的方法可以做到这一点。 Lazy init 和 pandas udf 听起来很有趣。我去看看。
    • 顺便说一句,在单例对象中维护状态怎么样 - 你为什么不提这个呢?
    • 主要是因为我发现在 Python 中使用单例对象并不是一个好的做法。此外,这个问题与范围更相关。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2015-01-16
    • 1970-01-01
    • 2019-11-23
    • 2019-12-05
    • 1970-01-01
    • 1970-01-01
    • 2021-05-03
    相关资源
    最近更新 更多