求SQL与Spark代码:识别60天内下单3次及以上的客户及批量订单
识别并标记60天内下单3次及以上客户的批量订单(SQL & Spark实现)
要实现识别任意时段内60天内下单3次及以上的客户,并仅将该时段内对应订单标记为批量订单,以下是具体实现方案:
SQL实现
假设数据表名为orders,包含customer_name、order_date、order_id字段。使用滑动窗口统计每个客户在当前订单日期往前60天内的订单数量,再根据数量标记批量订单:
SELECT customer_name, order_date, order_id, CASE WHEN order_count >= 3 THEN '是' ELSE '否' END AS is_batch_order FROM ( SELECT customer_name, order_date, order_id, -- 按客户分组,统计当前订单日期前60天内的订单总数 COUNT(order_id) OVER ( PARTITION BY customer_name ORDER BY order_date RANGE BETWEEN INTERVAL '60' DAY PRECEDING AND CURRENT ROW ) AS order_count FROM orders -- 可选:添加时段过滤条件,比如限定2023年的订单 -- WHERE order_date BETWEEN '2023-01-01' AND '2023-12-31' ) AS order_counts
注意:如果使用的SQL方言不支持RANGE BETWEEN INTERVAL语法(如低版本MySQL),可以改用日期函数结合关联查询实现:
SELECT o1.customer_name, o1.order_date, o1.order_id, CASE WHEN COUNT(o2.order_id) >= 3 THEN '是' ELSE '否' END AS is_batch_order FROM orders o1 LEFT JOIN orders o2 ON o1.customer_name = o2.customer_name AND o2.order_date BETWEEN DATE_SUB(o1.order_date, INTERVAL 60 DAY) AND o1.order_date -- 可选:添加时段过滤 -- WHERE o1.order_date BETWEEN '2023-01-01' AND '2023-12-31' GROUP BY o1.customer_name, o1.order_date, o1.order_id
Spark实现
Scala版本
基于Spark DataFrame,使用窗口函数实现相同逻辑,支持分布式计算场景:
import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions._ // 加载订单数据(假设已通过Spark表或文件加载) val ordersDF = spark.read.table("orders") // 定义滑动窗口:按客户分组,按订单日期排序,窗口范围为当前日期往前60天 val windowSpec = Window .partitionBy("customer_name") .orderBy(col("order_date").cast("timestamp")) // 转成timestamp以支持秒级范围计算 .rangeBetween(-60 * 24 * 3600, 0) // 60天对应的秒数,负数表示往前推的时间 // 计算订单数并标记批量订单 val batchOrdersDF = ordersDF .withColumn("order_count", count("order_id").over(windowSpec)) .withColumn("is_batch_order", when(col("order_count") >= 3, "是").otherwise("否")) .select("customer_name", "order_date", "order_id", "is_batch_order") // 可选:过滤特定时段的订单 // val filteredBatchOrdersDF = batchOrdersDF.where(col("order_date").between("2023-01-01", "2023-12-31")) // 查看结果 batchOrdersDF.show(10, truncate = false)
Python版本
语法逻辑与Scala一致,仅调整Python风格的写法:
from pyspark.sql import Window from pyspark.sql.functions import col, count, when ordersDF = spark.read.table("orders") windowSpec = Window \ .partitionBy("customer_name") \ .orderBy(col("order_date").cast("timestamp")) \ .rangeBetween(-60*24*3600, 0) batchOrdersDF = ordersDF \ .withColumn("order_count", count("order_id").over(windowSpec)) \ .withColumn("is_batch_order", when(col("order_count") >= 3, "是").otherwise("否")) \ .select("customer_name", "order_date", "order_id", "is_batch_order") batchOrdersDF.show(10, truncate=False)
内容的提问来源于stack exchange,提问作者Kim
相关产品推荐
相关产品推荐

