【问题标题】:efficient symmertic computation in sparkspark中的高效对称计算
【发布时间】:2021-01-26 04:48:08
【问题描述】:

我在包含对称性的算法中看到的一个常见结构是

for (int i = 0; i < n ; i++) {
    for (int j = i+1; j < n ; j++) {
        [compute x]
        objects[i][j] += x;
        objects[j][i] -= x;
    }
}

(虽然仍然具有 O(n^2) 复杂度)通过利用对称性减少了所需的计算量。你能告诉我在 pyspark 代码中引入这种优化的方法是什么吗?

例如,我编写的代码根据公式(其中r 是位置)计算作用于系统中每个粒子的每单位质量的力:

         N    m_j*(r_i - r_j)
F = -G * Σ   -----------------
        i!=j   |r_i - r_j|^3

在其中,我首先将我的数据框与自身进行叉积以获得每对交互,然后通过 id 将它们全部聚合以获得作用在每个粒子上的总力:

def calc_F(df_clust, G=1):

    # cartesian product of the dataframe with itself
    renameCols = [f"`{col}` as `{col}_other`" for col in df_clust.columns]
    df_cart = df_clust.crossJoin(df_clust.selectExpr(renameCols))
    df_clust_cartesian = df_cart.filter("id != id_other")

    df_F_cartesian = df_clust_cartesian.selectExpr("id", "id_other", "m_other",
                                                   "`x` - `x_other` as `diff(x)`",
                                                   "`y` - `y_other` as `diff(y)`",
                                                   "`z` - `z_other` as `diff(z)`"
                                                   )
    df_F_cartesian = df_F_cartesian.selectExpr("id", "id_other",
                                               "`diff(x)` * `m_other` as `num(x)`",
                                               "`diff(y)` * `m_other` as `num(y)`",
                                               "`diff(z)` * `m_other` as `num(z)`",
                                               "sqrt(`diff(x)` * `diff(x)` + `diff(y)`"
                                               "* `diff(y)` + `diff(z)` * `diff(z)`) as `denom`",
                                               )
    df_F_cartesian = df_F_cartesian.selectExpr("id", "id_other",
                                               "`num(x)` / pow(`denom`, 3) as `Fx`",
                                               "`num(y)` / pow(`denom`, 3) as `Fy`",
                                               "`num(z)` / pow(`denom`, 3) as `Fz`",
                                               )
    # squish back to inital particles
    sumCols = ["Fx", "Fy", "Fz"]
    df_agg = df_F_cartesian.groupBy("id").sum(*sumCols)
    renameCols = [f"`sum({col})` as `{col}`" for col in sumCols]
    df_F = df_agg.selectExpr("id", *renameCols)

    df_F = df_F.selectExpr("id",
                           f"`Fx` * {-G} as Fx",
                           f"`Fy` * {-G} as Fy",
                           f"`Fz` * {-G} as Fz")

    return df_F

但我知道两个粒子之间的力是对称的——F_ij = -F_ji(我假设所有质量都是相等的)——所以在这里我计算了两次力的数量,而不是重复使用它们。因此,在这种特殊情况下,例如,我想将df_clust_cartesian = df_cart.filter("id != id_other") 转换为df_clust_cartesian = df_cart.filter("id &lt; id_other"),并在计算函数第二部分的总力时以某种方式重用这些力。 (当然理想情况下我想学一般的做)

这种情况下的示例输入是

a = sc.parallelize([
    [0.48593906,-0.52435857,-0.53198230,0.46153894,-0.33775792E-01,-0.32276499,0.15625001E-04,1],
    [-0.65960690E-01,0.80844238E-01,-0.27603051,-0.57578009,1.1078150,-0.29340765,0.15625001E-04,2],
    [-0.34809157E-01,0.76795481E-01,-0.39087987,-0.55399138,-0.17386098,0.59250806E-01,0.15625001E-04,3]
])                                                           

from pyspark.sql.types import * 
 
clust_input = StructType([ 
    StructField('x',  DoubleType(), False), 
    StructField('y',  DoubleType(), False), 
    StructField('z',  DoubleType(), False), 
    StructField('vx', DoubleType(), False), 
    StructField('vy', DoubleType(), False), 
    StructField('vz', DoubleType(), False), 
    StructField('m',  DoubleType(), False), 
    StructField('id', IntegerType(), False) 
])    

df_clust = a.toDF(schema=clust_input) 

【问题讨论】:

    标签: python apache-spark pyspark


    【解决方案1】:

    基本上,您只想计算id &lt; other_id 时的公式,并使用该结果通过对称性生成id &gt; other_id 的所有元素。

    你只需要通过这个来修改你的过滤器

    df_clust_cartesian = df_cart.filter("id < id_other")
    

    然后,一旦你有了数据框df_F_cartesian,你就有了每对(id, id_other) 一行。您可以使用该行生成与(id_other, id) 对应的行,并将减号添加到Fx, FyFz

    这可以通过在聚合之前添加以下步骤来完成:

    from pyspark.sql import functions as F
    
    sumCols = ["Fx", "Fy", "Fz"]
    oppositeSums = [ (-F.col(c)).alias(c) for c in sumCols]
    df_F_cartesian = df_F_cartesian.select(F.explode(F.array(
        F.struct(F.col("id"), *sumCols),
        F.struct(F.col("id_other").alias("id"), *oppositeSums)
    )).alias("s")).select("s.*")
    

    【讨论】:

    • 我尝试了这种方法,不幸的是它比原来慢了近 4 倍;也许我会尝试一些不同的方式来生成对称数据
    • 所以目标是优化工作。在 spark 中,将计算量除以 2 并不是什么大问题(但这就是问题所在),真正重要的是 shuffle 和数据分区的方式。你有多少个不同的 id,有多少个执行者/核心,你的工作需要多少时间?
    • 是的,我开始认为将计算除以 2 可能实际上并不会使其明显更快。我有 64000 个 ID。当在 8 个线程 cpu + 默认并行度 16 上本地运行时,大约需要 5 分钟。在 17 个执行器上远程运行时,每个执行器 2 个核心 + 默认并行度 64,大约需要一分钟。然后我需要重复数千次(算法本身效率非常低)
    猜你喜欢
    • 1970-01-01
    • 2020-01-29
    • 2019-10-09
    • 2018-05-28
    • 1970-01-01
    • 2014-01-02
    • 2015-05-04
    • 2021-09-05
    • 2018-03-30
    相关资源
    最近更新 更多