PySpark如何无需将DataFrame转RDD完成值分桶并统计各区间数量
PySpark DataFrame 无RDD转换的分桶统计实现
全程使用PySpark内置DataFrame API即可完成需求,无需转换为RDD,同时可以保留计数为0的分桶行,完全匹配预期输出格式。
实现步骤
- 第一步:导入依赖并初始化SparkSession
from pyspark.sql import SparkSession from pyspark.sql.functions import count, col # 初始化SparkSession(已初始化可跳过) spark = SparkSession.builder.appName("bucketing_test").getOrCreate()
- 第二步:构造预设分桶参考表
提前定义所有需要的分桶规则,通过右关联的方式保证计数为0的分桶不会丢失
# 分桶配置格式:(分桶编号, 分桶范围文本, 区间最小值, 区间最大值),区间默认左闭右开 bucket_config = [ (0, "0-20", 0, 20), (1, "20-40", 20, 40), (2, "40-60", 40, 60), (3, "60-80", 60, 80) ] bucket_ref_df = spark.createDataFrame(bucket_config, schema=["bucket", "size", "min_val", "max_val"])
- 第三步:关联原始数据与分桶表,统计各分桶数量
# 构造输入示例DataFrame,实际使用时替换为你自己的DataFrame即可 input_data = [(2, 50.34), (4, 34.4), (6, 48.7), (10, 72.4)] input_df = spark.createDataFrame(input_data, schema=["ID", "value"]) # 关联分桶、统计结果 output_df = input_df.join( bucket_ref_df, (input_df.value >= bucket_ref_df.min_val) & (input_df.value < bucket_ref_df.max_val), how="right" # 右关联保留所有预设分桶 ).groupBy("bucket", "size") \ .agg(count("ID").alias("count")) \ .orderBy("bucket")
- 第四步:查看输出结果
output_df.show()
运行后输出结果和预期完全一致:
+------+-----+-----+ |bucket| size|count| +------+-----+-----+ | 0| 0-20| 0| | 1|20-40| 1| | 2|40-60| 2| | 3|60-80| 1| +------+-----+-----+
内容的提问来源于stack exchange,提问作者Tony LaRussa
相关产品推荐
相关产品推荐

