PySpark RDD中合并多字典:按grid_id聚合骑手经纬度数据
PySpark RDD: Merge Data into Dictionary Grouped by
grid_id with rider_id as Keys 我来帮你解决这个问题!你当前代码的核心问题是**dict.update()方法的返回值是None**——它会原地修改字典,但不会返回修改后的字典,所以reduceByKey每次合并时都会把None传递下去,最终得到(grid_id, None)的结果。
下面给你两种在RDD层面实现需求的可行方案:
方案1:修改reduceByKey的合并逻辑
既然update()返回None,我们可以改用字典解包的方式创建新字典来合并,这样就能返回正确的合并结果:
from pyspark.sql import SparkSession sqlContext = SparkSession.builder.appName("test").enableHiveSupport().getOrCreate() data = [(1,2,0.1,0.3),(1,2,0.1,0.3),(1,3,0.1,0.3),(1,3,0.1,0.3), (11, 12, 0.1, 0.3),(11,12,0.1,0.3),(11,13,0.1,0.3),(11,13,0.1,0.3)] trajectory_df = sqlContext.createDataFrame(data, schema=['grid_id','rider_id','lng','lat']) # 先转换成((grid_id, rider_id), coords_list)的格式,和你之前的步骤一致 trajectory_rdd = trajectory_df.rdd.map(lambda row: ((row.grid_id, row.rider_id), [row.lng, row.lat]))\ .groupByKey().mapValues(list) # 转换为(grid_id, {rider_id: coords_list}),然后用字典解包合并 result_rdd = trajectory_rdd.map(lambda x: (x[0][0], {x[0][1]: x[1]}))\ .reduceByKey(lambda x, y: {**x, **y}) print(result_rdd.take(10))
输出结果:
[(1, {2: [[0.1, 0.3], [0.1, 0.3]], 3: [[0.1, 0.3], [0.1, 0.3]]}), (11, {12: [[0.1, 0.3], [0.1, 0.3]], 13: [[0.1, 0.3], [0.1, 0.3]]})]
方案2:先按grid_id分组,再构建字典
这种方式更直观,先把所有数据按grid_id分组,然后在每个分组内部直接构建目标字典:
from pyspark.sql import SparkSession sqlContext = SparkSession.builder.appName("test").enableHiveSupport().getOrCreate() data = [(1,2,0.1,0.3),(1,2,0.1,0.3),(1,3,0.1,0.3),(1,3,0.1,0.3), (11, 12, 0.1, 0.3),(11,12,0.1,0.3),(11,13,0.1,0.3),(11,13,0.1,0.3)] trajectory_df = sqlContext.createDataFrame(data, schema=['grid_id','rider_id','lng','lat']) # 转换为(grid_id, (rider_id, coords))格式 rdd = trajectory_df.rdd.map(lambda row: (row.grid_id, (row.rider_id, [row.lng, row.lat]))) # 按grid_id分组,然后在每个分组内构建字典:key是rider_id,value是对应的coords列表 result_rdd = rdd.groupByKey().mapValues(lambda group: { rider: [coord for (r, coord) in group if r == rider] for rider in set(r for r, coord in group) }) print(result_rdd.take(10))
这个方案的输出和方案1完全一致,逻辑上更容易理解,适合处理更复杂的分组场景。
为什么你的原代码会返回None?
再强调一下:dict.update(y)会把y的键值对添加到x里,但这个操作是原地修改x,函数本身的返回值是None。所以reduceByKey(lambda x,y:x.update(y))中,每次合并都会返回None,最终每个grid_id对应的结果就是None。
内容的提问来源于stack exchange,提问作者giser_yugang
相关产品推荐
相关产品推荐

