【发布时间】: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