You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.13 08:36:22