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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.01 02:48:03