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

将含键与字典的温度数据RDD转换为DataFrame(PySpark实现)

处理多区域温度数据:从RDD按键合并到DataFrame转换

我来帮你把这段代码补全并梳理清楚逻辑,针对你手里的多时间、多大洲温度测量数据,我们可以按下面的步骤来实现从RDD到按键合并再转DataFrame的完整流程:

第一步:定义聚合所需的辅助函数

这三个函数是为了配合RDD的aggregateByKey操作,实现按键(这里就是大洲名称)合并所有对应的数据记录:

# 将单个元素转为列表,作为聚合的初始值
def to_list(a):
    return [a]

# 分区内合并:把新元素追加到已有列表中
def append(a, b):
    a.append(b)
    return a

# 分区间合并:把两个列表合并成一个
def extend(a, b):
    a.extend(b)
    return a

第二步:完整的处理流程代码

接下来我们把整个流程补全,包括初始化Spark环境、加载数据、聚合、转换为DataFrame:

from pyspark.sql import SparkSession
from pyspark.sql import Row

def to_list(a): return [a]
def append(a, b): 
    a.append(b) 
    return a
def extend(a, b): 
    a.extend(b) 
    return a

def main():
    # 初始化SparkSession(替代单独的SparkContext,更适合DataFrame操作)
    spark = SparkSession.builder.appName("TemperatureDataProcessing").getOrCreate()
    sc = spark.sparkContext

    # 模拟你的温度测量数据RDD(和你提供的样例结构一致)
    data = [
        ('Africa', {'time': '1', 'temp': '2'}),
        ('Africa', {'time': '2', 'temp': '3'}),
        ('America', {'time': '1', 'temp': '18'}),
        ('America', {'time': '2', 'temp': '20'}),
        ('Asia', {'time': '1', 'temp': '15'}),
        ('Asia', {'time': '2', 'temp': '17'})
    ]
    temp_rdd = sc.parallelize(data)

    # 按键(大洲)聚合,把同一个大洲的所有记录汇总成列表
    aggregated_rdd = temp_rdd.aggregateByKey([], to_list, append, extend)

    # 方式1:直接将聚合后的RDD转为DataFrame(每个大洲对应一个记录列表)
    df_aggregated = aggregated_rdd.toDF(["continent", "temperature_records"])
    print("聚合后的DataFrame(按大洲分组):")
    df_aggregated.show(truncate=False)

    # 方式2:展开记录列表,每条温度记录单独成行(更常用的结构化格式)
    flattened_rdd = aggregated_rdd.flatMapValues(lambda records: records) \
                                 .map(lambda x: Row(continent=x[0], time=x[1]['time'], temp=x[1]['temp']))
    df_flattened = spark.createDataFrame(flattened_rdd)
    print("\n展开后的结构化DataFrame:")
    df_flattened.show()

    # 停止Spark会话
    spark.stop()

if __name__ == "__main__":
    main()

关键步骤解释

  • aggregateByKey操作:这里我们用它来按大洲名称分组,把每个大洲的所有温度记录收集到一个列表里。初始值是空列表,分区内用append添加新记录,分区间用extend合并不同分区的列表。
  • DataFrame转换:提供了两种转换方式:
    1. 保留按大洲分组的结构,每个大洲对应一个记录列表,适合查看分组后的整体数据;
    2. 展开每条记录为单独的行,转换成标准的结构化DataFrame,方便后续进行数据分析、过滤、统计等操作。

内容的提问来源于stack exchange,提问作者Yaniv Irony

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 11:47:42