Spark DataFrame按地区统计:获取Top3高频犯罪类型及月度犯罪数中位数
Spark分组统计:各地区Top3犯罪类型及月度犯罪中位数
原始数据构造
df = spark.createDataFrame( [('D1', 'ROBBERY', '2024-02-01', 2), ('D1', 'ROBBERY', '2024-02-01', 2), ('D1', 'DRUGS', '2024-03-05', 3), ('D1', 'FRAUD', '2024-03-05', 3), ('D1', 'AUTO THEFT', '2024-01-09',1), ('D1', 'AUTO THEFT', '2024-01-03', 1), ('D2', 'MURDER', '2024-05-04', 5), ('D2', 'MURDER', '2024-06-01', 6), ('D2', 'RAPE', '2024-07-02', 7)], ['district', 'crime_type', 'date', 'month'])
实现步骤
1. 计算各地区Top3犯罪类型
先按district和crime_type分组统计频次,再对每个地区内的犯罪类型按频次降序排序,取前3个后用逗号拼接成字符串:
from pyspark.sql import Window import pyspark.sql.functions as F # 统计每个地区各犯罪类型的出现次数 crime_count = df.groupBy("district", "crime_type")\ .count()\ .withColumnRenamed("count", "crime_freq") # 按地区分组排序,取前3犯罪类型 window_top3 = Window.partitionBy("district").orderBy(F.desc("crime_freq")) top3_crimes = crime_count.withColumn("rank", F.row_number().over(window_top3))\ .filter(F.col("rank") <= 3)\ .groupBy("district")\ .agg(F.concat_ws(", ", F.collect_list("crime_type")).alias("top_3_crime_types"))
2. 计算各地区月度犯罪数量的中位数
先按district和month分组统计每月犯罪数,再对每个地区的月度犯罪数求中位数:
# 统计每个地区每月的犯罪数量 monthly_crime = df.groupBy("district", "month")\ .count()\ .withColumnRenamed("count", "monthly_crime_num") # 计算月度犯罪数的中位数 median_result = monthly_crime.groupBy("district")\ .agg(F.expr("percentile_approx(monthly_crime_num, 0.5)").alias("median_crimes_monthly"))
3. 合并结果
将Top3犯罪类型和中位数结果按district关联,得到最终DataFrame:
final_df = top3_crimes.join(median_result, on="district", how="inner") final_df.show(truncate=False)
预期输出
| district | top_3_crime_types | median_crimes_monthly |
|---|---|---|
| D1 | ROBBERY, AUTO THEFT, DRUGS | 2 |
| D2 | MURDER, RAPE | 1 |
内容的提问来源于stack exchange,提问作者Dliv
相关产品推荐
相关产品推荐

