【问题标题】:PySpark broadcast variables from local functionsPySpark 广播来自本地函数的变量
【发布时间】:2015-01-13 13:40:35
【问题描述】:

我正在尝试从 Python 方法中创建广播变量(尝试抽象我正在创建的一些依赖于分布式操作的实用程序方法)。但是,我似乎无法从 Spark 工作人员中访问广播变量。

假设我有这个设置:

def main():
    sc = SparkContext()
    SomeMethod(sc)

def SomeMethod(sc):
    someValue = rand()
    V = sc.broadcast(someValue)
    A = sc.parallelize().map(worker)

def worker(element):
    element *= V.value  ### NameError: global name 'V' is not defined ###

但是,如果我改为消除 SomeMethod() 中间人,它就可以正常工作。

def main():
    sc = SparkContext()
    someValue = rand()
    V = sc.broadcast(someValue)
    A = sc.parallelize().map(worker)

def worker(element):
    element *= V.value   # works just fine

如果可以的话,我宁愿不必将所有 Spark 逻辑都放在 main 方法中。有什么方法可以从本地函数中广播变量并使它们对 Spark 工作人员全局可见?

或者,对于这种情况,什么是好的设计模式——例如,我想为 Spark 编写一个专门的方法,它是自包含的并执行我想重用的特定功能?

【问题讨论】:

    标签: python apache-spark pyspark


    【解决方案1】:

    我不确定我是否完全理解了这个问题,但是,如果您需要工作函数中的 V 对象,那么您绝对应该将它作为参数传递,否则该方法并不是真正独立的:

    def worker(V, element):
        element *= V.value
    

    现在为了在 map 函数中使用它,你需要使用一个 partial,这样 map 只能看到一个 1 参数函数:

    from functools import partial
    
    def SomeMethod(sc):
        someValue = rand()
        V = sc.broadcast(someValue)
        A = sc.parallelize().map(partial(worker, V=V))
    

    【讨论】:

    • 像这样传递广播变量是否会对性能产生影响?举例来说,我在 map() 函数中依赖于数万(或更多)行的广播变量。 def transform(row): return broadcast_variable.value[row[0]] 之类的东西,然后在 map() 函数中使用,例如 rdd.map(transform)
    • 感谢这个解决方案帮助我避免使用 global 作为广播变量。请注意,您应该替换工作方法参数的顺序,以便“元素”参数(由 Spark 框架填充)将是第一个。否则将无法正常工作。
    • @elyase 非常感谢。你真的拯救了我的一天。我对 spark 有不同的问题,但你在我的场景中的建议就像一个魅力。非常感谢!
    猜你喜欢
    • 2015-06-27
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多