【问题标题】:Find nearest stores given single store location + what is max broadcast variable size in pyspark?在给定单个商店位置的情况下查找最近的商店 + pyspark 中的最大广播变量大小是多少?
【发布时间】:2018-02-08 00:57:03
【问题描述】:

我正在尝试将包含所有杂货店的文件读入 RDD。目标是找出给定杂货店最近的 500 家商店并进行一些处理。 例如,如果一家商店被占用,我必须找到该城市的所有商店。如果我确实在 RDD 上进行映射,我会获得转换函数的单一存储。我如何获得该功能中的所有商店。

片段

def getNearestStores(store):
    stores_city = stores.filter("city="+store.city)
    return (store.id,stores_city.count()) 
stores = sc.textFile("stores.json").map(getNearestStores).count()

这是简单的代码sn-p。 Stores.json 文件很大

1) 如何在 getNearestStores 函数中获取最近的 500 家商店,大概是通过使用 stores.json?

2) PySpark 中的最大广播变量大小是多少?

【问题讨论】:

  • 如果我从代码中理解了任务,它会为每家商店计算同一城市中所有商店的数量......如果这是正确的,我不会使用广播变量......我会在 RDD[Store] 上执行类似操作以获得 RDD[(store_id, count stores in same city)]...rdd.keyBy(_.city).groupByKey.flatMap{ case (city, iter_stores) => iter_stores.map(one_store => (one_store.id, iter_stores.size)) }
  • @kmh 它不仅仅是计数。需要执行额外的处理。这是完整的用例:对于每个商店,获取同一城市附近的 500 家商店,并对 500 家商店进行一些处理。最终结果将是(商店,[500 个附近的已处理商店])
  • 我仍然会做一个 keyBy、groupByKey、flatMap 来到达那里。我还会更改您问题的标题,因为这实际上与广播变量最大大小无关。

标签: apache-spark pyspark pyspark-sql


【解决方案1】:

~2GB 是广播变量的最大大小,因为任何广播变量在序列化期间都会变成 java 字节数组,而 java 数组的最大大小为 Integer.MAX_VALUE。

这些是与 PySpark 2.x 相关的较旧(类似)问题

Is there any limit on size of a spark broadcast variable?

Evaluate the max size for a spark broadcast variable

用于跟踪此问题的 JIRA:https://issues.apache.org/jira/browse/SPARK-6235

编辑: 使用 SparkSession + DataFrame(因为这是 PySpark - 大概是 2.x)进行如下连接。

from pyspark.sql import SparkSession

spark = SparkSession \
    .builder \
    .appName("Python Spark SQL basic example") \
    .config("spark.some.config.option", "some-value") \
    .getOrCreate()
df = spark.read.json("stores.json")
nearest_stores_df = df.join(df, "city") # self-join

【讨论】:

  • 好的。那讲得通。但是您知道我该如何解决上述问题吗?
  • 是的。 1)使用 SparkSession + DataFrame 或 Dataset,2)为商店创建一个表并在城市上进行连接,而不是广播的高成本
【解决方案2】:

回复您的上述评论,详细说明您真正想要完成的工作。

我不在 spark 中使用 python,但在 scala 中我会这样开始。大多数重要的东西都是特定于 spark 的,应该很容易转换为 python。

商店类:

class Store(val id: String, 
            val city: String, 
            val lat: Double, 
            val lon: Double) extends Haversine {

  // distance in km 
  def distance(that: Store) = {
    val radius = 6371 // earth in km 
    val dlat = toRadians(that.lat - this.lat)
    val dlon = toRadians(that.lon - this.lon)
    val a = sin(dlat/2) * sin(dlat/2) +
            cos(toRadians(this.lat)) * cos(toRadians(that.lat)) * sin(dlon/2) * sin(dlon/2)
    val c = 2 * atan2(sqrt(a), sqrt(1-a))
    val d = radius * c 
    d 
  }
}

既然你仍然在使用 python,我会留给你把你的 JSON 解析到你的类中。

// read JSON into class Store
val rdd_store: RDD[Store] = ???

这是重要的部分...在我调试之前我的看起来会像这样...但是如果我打算将难以阅读的部分移动到辅助函数或类中真实的。

// result
val rdd2: RDD[(Store, List[Store])] = {
  rdd_store
  .keyBy(_.city)                               // RDD[(city, store)]               (one row per store)
  .groupByKey                                  // RDD[(city, all stores in city)]  (one row per city)
  .flatMap { case (city, iter_stores) =>
    iter_stores.map(one_store => {
      // up-to 500 closest stores in same city as 'one_store'
      val top500: List[Store] = {
        iter_stores
        .toList
        .map { s => 
          val dist = one_store distance s
          (dist, s)
        }
      }.sorted.take(500).map{ case (dist, that_store) => that_store }

      // result
      (one_store, top500)
    })
  }                                           // RDD[(store, list of top 500 stores in same city)] (one row per store)
}

【讨论】:

  • 好的。这真的很有帮助。非常感谢。请提供另一个建议:如果该城市没有 500 家商店,我还有一个文件说明要选择哪些城市才能达到 500 家商店。我怎么能这样做。
猜你喜欢
  • 2023-03-05
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-10-22
  • 1970-01-01
  • 2020-02-12
  • 1970-01-01
相关资源
最近更新 更多