如何用SQL/Spark SQL统计指定信用卡近30天交易频次?
统计指定信用卡30天交易频次的SQL与Spark Java实现
前提说明
基于transactions表结构(含credit_card_number信用卡号字段、transaction_timestamp交易时间戳字段),以下实现针对2019年1月1日至2020年12月31日的交易数据,统计指定信用卡的30天滚动交易频次,并按年月维度聚合展示。
1. 标准SQL查询实现
场景1:按年月展示每月内的30天滚动交易频次(每笔交易对应过去30天的累计次数)
-- 替换 'XXXX-XXXX-XXXX-XXXX' 为目标信用卡号 SELECT DATE_FORMAT(transaction_timestamp, 'yyyy-MM') AS year_month, COUNT(*) AS monthly_transaction_total, -- 计算当前交易日期往前30天内的累计交易次数 COUNT(*) OVER ( PARTITION BY credit_card_number ORDER BY transaction_timestamp RANGE BETWEEN INTERVAL 30 DAY PRECEDING AND CURRENT ROW ) AS rolling_30d_transaction_count FROM transactions WHERE credit_card_number = 'XXXX-XXXX-XXXX-XXXX' AND transaction_timestamp BETWEEN '2019-01-01 00:00:00' AND '2020-12-31 23:59:59' GROUP BY year_month, transaction_timestamp, credit_card_number ORDER BY year_month;
场景2:按年月统计每月末往前30天的交易总频次
-- 替换 'XXXX-XXXX-XXXX-XXXX' 为目标信用卡号 SELECT year_month, COUNT(*) AS rolling_30d_total FROM ( SELECT DATE_FORMAT(transaction_timestamp, 'yyyy-MM') AS year_month, transaction_timestamp FROM transactions WHERE credit_card_number = 'XXXX-XXXX-XXXX-XXXX' AND transaction_timestamp BETWEEN '2019-01-01 00:00:00' AND '2020-12-31 23:59:59' -- 筛选当月最后一天往前30天内的交易 AND transaction_timestamp >= LAST_DAY(transaction_timestamp) - INTERVAL 30 DAY ) t GROUP BY year_month ORDER BY year_month;
2. Java Spark SQL 实现
以下代码通过Spark计算指定信用卡的30天滚动交易频次,按年月聚合展示月度总交易数、最大30天频次、平均30天频次:
import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Row; import org.apache.spark.sql.SparkSession; import org.apache.spark.sql.Window; import org.apache.spark.sql.WindowSpec; import static org.apache.spark.sql.functions.*; public class CardTransactionStats { public static void main(String[] args) { // 初始化SparkSession(生产环境移除master配置) SparkSession spark = SparkSession.builder() .appName("Card30dTransactionFrequency") .master("local[*]") .getOrCreate(); // 读取transactions表,根据实际存储格式调整(如CSV、JDBC、Parquet) Dataset<Row> transactionsDF = spark.read() .format("parquet") .load("/path/to/transactions-data"); // 目标信用卡号 String targetCard = "XXXX-XXXX-XXXX-XXXX"; // 定义窗口:按信用卡号分区,按交易时间戳排序,取过去30天范围 WindowSpec windowSpec = Window.partitionBy("credit_card_number") .orderBy(col("transaction_timestamp").cast("long")) .rangeBetween(-30 * 24 * 60 * 60, 0); // 30天转换为秒数 // 计算并聚合结果 Dataset<Row> resultDF = transactionsDF .filter(col("credit_card_number").equalTo(targetCard)) .filter(col("transaction_timestamp").between("2019-01-01 00:00:00", "2020-12-31 23:59:59")) .withColumn("year_month", date_format(col("transaction_timestamp"), "yyyy-MM")) .withColumn("rolling_30d_count", count("*").over(windowSpec)) .groupBy("year_month") .agg( count("*").alias("monthly_total_transactions"), max("rolling_30d_count").alias("max_30d_frequency"), avg("rolling_30d_count").alias("avg_30d_frequency") ) .orderBy("year_month"); // 输出结果 resultDF.show(); // 关闭SparkSession spark.stop(); } }
注意事项
- 若
transaction_timestamp为日期类型而非时间戳,需调整日期函数语法(如使用date_add替代时间戳范围计算)。 - Spark SQL中部分版本支持直接用
interval 30 day作为窗口范围,可替换秒数计算逻辑以提升可读性。 - 若从数据库读取数据,需将
read()部分替换为JDBC连接配置。
内容的提问来源于stack exchange,提问作者Gumaniuc Alexandru
相关产品推荐
相关产品推荐

