PySpark中如何根据多行列条件统计同时使用两种支付方式的ID数量
PySpark实现统计同时使用两种支付方式的ID总数
实现思路
- 首先对原始数据按
ID分组,统计每个ID的不同支付方式数量 - 筛选出同时覆盖
Cash和Credit Card两种支付方式的ID - 对符合条件的ID直接计数即可得到结果
代码实现
首先导入依赖包,初始化SparkSession并构造示例数据:
from pyspark.sql import SparkSession from pyspark.sql.functions import countDistinct # 初始化Spark会话 spark = SparkSession.builder.appName("payment_stat").getOrCreate() # 构造测试数据集 data = [ (1, "Cash", 200, "2020-01-01"), (1, "Credit Card", 500, "2020-01-06"), (2, "Cash", 300, "2020-02-01"), (3, "Credit Card", 400, "2020-02-02"), (3, "Credit Card", 500, "2020-01-03"), (3, "Cash", 200, "2020-01-04") ] df = spark.createDataFrame(data, schema=["ID", "Payment", "Value", "Date"])
核心统计逻辑有两种常用写法,都可以实现需求:
写法1:先过滤再分组统计
# 先过滤出只涉及目标两种支付方式的记录,按ID分组统计不同支付方式的数量 # 筛选出支付方式数量为2的ID,计数即为最终结果 res = df.filter(df.Payment.isin("Cash", "Credit Card")) \ .groupBy("ID") \ .agg(countDistinct("Payment").alias("pay_type_count")) \ .filter("pay_type_count = 2") \ .count() print(res) # 输出结果为2,符合示例预期
写法2:收集支付方式集合后匹配
from pyspark.sql.functions import collect_set, array_contains # 按ID分组收集去重后的支付方式集合,筛选同时包含两种支付方式的ID计数 res = df.groupBy("ID") \ .agg(collect_set("Payment").alias("pay_types")) \ .filter(array_contains("pay_types", "Cash") & array_contains("pay_types", "Credit Card")) \ .count() print(res) # 输出结果同样为2
内容的提问来源于stack exchange,提问作者Alf
相关产品推荐
相关产品推荐

