Spark DataFrame中基于季度聚合的Distinct计数改造方案咨询
按指定季度区间统计Spark DataFrame聚合值的解决方案
要实现你要求的按自定义季度区间统计每个x_id+product下的去重h_id数,我们可以分两步处理:给每条数据打上对应的区间标签,再按区间标签分组聚合。下面提供两种实用的实现方案:
方案一:直接用case when映射区间(适合固定区间场景)
如果你的季度区间不会频繁变动,直接在DataFrame中通过case when逻辑给每条数据分配对应的区间标识是最直接的方式:
from pyspark.sql import functions as sf start = '2018-11-01' end = '2021-02-19' # 第一步:加载数据并添加季度区间标识 t1_with_period = ( spark.table('schema.t1') .where(sf.col('yyyy_mm_dd').between(start, end)) .select('yyyy_mm_dd', 'x_id', 'h_id', 'app', 'kpi') .withColumn('quarter_period', sf.when(sf.col('yyyy_mm_dd').between('2018-11-01', '2018-12-31'), '2018Q4(11-12)') .when(sf.col('yyyy_mm_dd').between('2019-01-01', '2019-03-31'), '2019Q1') .when(sf.col('yyyy_mm_dd').between('2019-04-01', '2019-06-30'), '2019Q2') .when(sf.col('yyyy_mm_dd').between('2019-07-01', '2019-09-30'), '2019Q3') .when(sf.col('yyyy_mm_dd').between('2019-10-01', '2019-12-31'), '2019Q4') .when(sf.col('yyyy_mm_dd').between('2020-01-01', '2020-03-31'), '2020Q1') .when(sf.col('yyyy_mm_dd').between('2020-04-01', '2020-06-30'), '2020Q2') .when(sf.col('yyyy_mm_dd').between('2020-07-01', '2020-09-30'), '2020Q3') .when(sf.col('yyyy_mm_dd').between('2020-10-01', '2020-12-31'), '2020Q4') .when(sf.col('yyyy_mm_dd').between('2021-01-01', '2021-02-19'), '2021Q1(1-2)') .otherwise('其他') # 过滤后的数据不会触发此分支,仅作兜底 ) ) # 第二步:关联产品表并按区间聚合统计 aggregate = ( t1_with_period .join(t2, on=['app', 'kpi'], how='left') .groupby('x_id', 'product', 'quarter_period') # 新增quarter_period作为分组维度 .agg( sf.countDistinct('h_id').alias('count_ever') ) .orderBy('x_id', 'product', 'quarter_period') # 可选,让结果按顺序展示 )
方案二:用区间映射表关联(适合区间需灵活调整的场景)
如果后续可能需要修改季度区间,建议创建一个独立的区间映射表,通过关联的方式给数据打标签,更易于维护:
from pyspark.sql import functions as sf start = '2018-11-01' end = '2021-02-19' # 第一步:创建自定义季度区间映射表 periods_data = [ ('2018Q4(11-12)', '2018-11-01', '2018-12-31'), ('2019Q1', '2019-01-01', '2019-03-31'), ('2019Q2', '2019-04-01', '2019-06-30'), ('2019Q3', '2019-07-01', '2019-09-30'), ('2019Q4', '2019-10-01', '2019-12-31'), ('2020Q1', '2020-01-01', '2020-03-31'), ('2020Q2', '2020-04-01', '2020-06-30'), ('2020Q3', '2020-07-01', '2020-09-30'), ('2020Q4', '2020-10-01', '2020-12-31'), ('2021Q1(1-2)', '2021-01-01', '2021-02-19'), ] periods_df = spark.createDataFrame(periods_data, ['quarter_period', 'period_start', 'period_end']) # 确保日期字段类型为date periods_df = periods_df.withColumn('period_start', sf.to_date('period_start')) periods_df = periods_df.withColumn('period_end', sf.to_date('period_end')) # 第二步:加载业务数据并关联区间映射表 t1_with_period = ( spark.table('schema.t1') .where(sf.col('yyyy_mm_dd').between(start, end)) .select('yyyy_mm_dd', 'x_id', 'h_id', 'app', 'kpi') .join( periods_df, sf.col('yyyy_mm_dd').between(periods_df['period_start'], periods_df['period_end']), how='left' ) ) # 第三步:关联产品表并聚合统计 aggregate = ( t1_with_period .join(t2, on=['app', 'kpi'], how='left') .groupby('x_id', 'product', 'quarter_period') .agg( sf.countDistinct('h_id').alias('count_ever') ) .orderBy('x_id', 'product', 'quarter_period') )
关键说明
- 区间标识命名:你可以根据业务需求修改
quarter_period的取值(比如用2018-11至2018-12这种更直观的格式),只要能清晰区分每个区间即可。 - 数据完整性:由于我们已经过滤了
yyyy_mm_dd在start和end之间的数据,且区间映射表完全覆盖了该时间范围,所以每条业务数据都会匹配到唯一的区间,不会出现quarter_period为null的情况。 - 性能优化:如果数据量极大,
countDistinct可能会较慢,你可以考虑用sf.approx_count_distinct('h_id', rsd=0.05)做近似统计(误差可通过rsd参数调整),或者先按区间+h_id去重再聚合,提升效率。
内容的提问来源于stack exchange,提问作者Someguywhocodes
相关产品推荐
相关产品推荐

