ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Spark电商用户行为分析实战:漏斗、复购率与RFM分群

Spark电商用户行为分析实战:漏斗、复购率与RFM分群 简介这是一套面向计算机专业本科生的电商用户行为分析实战项目适用于Java课程设计、毕业设计及期末大作业场景聚焦Spark实时计算与用户行为路径挖掘核心能力训练。资源包含273个文件以58个Scala核心业务逻辑文件如UserSessionAnalysisFunction2、AreaTop3ProductFunc等为主体辅以208个XML配置与依赖管理文件、2个properties环境配置文件及README.md等说明文档整体压缩包仅177KB轻量易部署。已有63人学习下载项目经本地编译验证可直接运行评审得分98分内容由助教审定难度适中且结构完整——涵盖数据接入、会话分析、区域热榜、行为漏斗等典型模块代码规范、注释清晰并附详细文档说明便于理解Spark Streaming与RDD协同处理逻辑及电商分析指标设计思路。1. 为什么电商团队现在必须用 Spark 做用户行为分析而不是 Hive 或 MySQL某中型电商平台上线半年后运营发现用户从首页点击商品到最终下单的路径中有 37% 的会话在「加入购物车」环节中断但传统 BI 工具查不出是前端卡顿、库存同步延迟还是推荐策略失效。他们用 Spark 重写了行为日志处理链路——不是因为 Spark 更“酷”而是因为原始日志是每秒 20 万条的 JSON 流含设备 ID、页面停留时长、滚动深度、按钮点击坐标单次会话跨度可达 47 分钟且需关联用户画像表千万级、商品类目表百万级、促销活动表动态更新。Hive 批处理跑一次全量路径还原要 6 小时MySQL 直接扛不住写入吞吐。Spark 的结构化流处理Structured Streaming DataFrame API 外部 shuffle 服务如 Spark on K8s with external shuffle service让这个场景真正可落地既能按 session window 实时聚合用户动作序列又能用 broadcast join 快速关联维度表还能通过spark.sql.adaptive.enabledtrue自动优化倾斜 join。本文不讲 Spark 安装或 Scala 语法只聚焦「如何把一份真实电商用户行为日志用 Spark 跑出复购率、跳失率、路径转化漏斗、RFM 分群这四类业务指标」——所有代码基于 Spark 3.3适配 HDFS/S3/OSS 三种存储参数配置来自生产集群调优实测数据。2. 搭建最小可行分析环境本地伪分布式 Spark 模拟电商日志生成器2.1 为什么选 Spark 3.3 而非 2.x 或 4.xSpark 3.3 是当前电商数仓最稳定的 LTS 版本它原生支持 Iceberg 0.14解决小文件合并问题、内置 AQEAdaptive Query Execution对 skew join 的自动拆分比 3.2 提升 40%、DataFrameReader 支持option(multiline, true)直接解析嵌套 JSON 日志。而 Spark 4.x 尚未被主流云厂商如阿里云 EMR、腾讯 EMR全面适配Spark 2.4 的 Catalyst 优化器无法识别window函数中的range between语义导致会话超时计算不准。验证版本命令# 下载官方二进制包非源码编译 wget https://archive.apache.org/dist/spark/spark-3.3.4/spark-3.3.4-bin-hadoop3.tgz tar -xzf spark-3.3.4-bin-hadoop3.tgz export SPARK_HOME$(pwd)/spark-3.3.4-bin-hadoop3 export PATH$SPARK_HOME/bin:$PATH spark-submit --version # 输出应为Welcome to # ____ __ # / __/__ ___ _____/ /__ # _\ \/ _ \/ _ / __/ _/ # /___/ .__/\_,_/_/ /_/\_\ version 3.3.4 # /_/提示不要用spark-shell交互式调试行为分析逻辑——它默认内存仅 1G且无法复现生产中spark.sql.adaptive.enabled等关键参数生效路径。所有测试必须用spark-submit提交脚本。2.2 用 Python 生成符合真实电商场景的模拟日志真实日志字段必须包含event_time(ISO8601)、user_id(MD5 hash)、event_type(view/click/add_cart/buy)、page_url(含 utm 参数)、item_id(空字符串表示非商品页)、session_id(15分钟无操作则新会话)、device_type(mobile/pc/tablet)、referrer(来源渠道)。以下脚本生成 10 万行带时间序列依赖的日志避免随机时间戳导致会话断裂# generate_log.py import json import time import random from datetime import datetime, timedelta def gen_session_events(): start_time datetime.now() - timedelta(hours24) user_ids [fuser_{i:06d} for i in range(1000)] items [fitem_{i:08d} for i in range(5000)] pages [home, search, category, product, cart, checkout, pay_success] logs [] for _ in range(100000): user random.choice(user_ids) session_start start_time timedelta(secondsrandom.randint(0, 86400)) # 会话内事件时间递增间隔 1-30 秒 event_time session_start session_id fsess_{int(event_time.timestamp())}_{user} # 模拟典型路径home → search → product → add_cart → buy path random.choices( [[home, search, product, add_cart, buy], [home, category, product, buy], [home, product, buy]], weights[0.5, 0.3, 0.2] )[0] for i, page in enumerate(path): event_time timedelta(secondsrandom.randint(1, 30)) log { event_time: event_time.isoformat(), user_id: user, event_type: view if i 0 else (click if page ! product else view), page_url: fhttps://shop.com/{page}?utm_sourcedirectutm_medium{random.choice([app, wechat, baidu])}, item_id: random.choice(items) if page product else , session_id: session_id, device_type: random.choice([mobile, pc, tablet]), referrer: https://google.com if i 0 else } # 在 product 页加 click 和 add_cart 事件 if page product: log[event_type] view logs.append(log.copy()) log[event_type] click log[page_url] log[page_url].replace(product, product_detail) logs.append(log.copy()) if random.random() 0.7: log[event_type] add_cart log[item_id] log[item_id] logs.append(log.copy()) else: logs.append(log) return logs if __name__ __main__: logs gen_session_events() with open(simulated_logs.json, w) as f: for log in logs: f.write(json.dumps(log, ensure_asciiFalse) \n)运行后生成simulated_logs.json每行一个 JSON 对象——这是 Spark Structured Streaming 的标准输入格式也是后续所有分析的原始数据源。2.3 用 spark-submit 运行第一个 ETL 任务清洗并写入 Parquet 分区表电商日志必须按天分区dt20240520且需过滤掉event_type为空或user_id为测试账号test_*的数据。以下 Scala 脚本完成三件事1读取 JSON 日志2解析event_time为timestamp类型并提取dt分区字段3写入本地./data/ods_user_event目录// etl_job.scala import org.apache.spark.sql.{SparkSession, DataFrame} import org.apache.spark.sql.functions._ object OdsEventETL { def main(args: Array[String]): Unit { val spark SparkSession.builder() .appName(OdsUserEventETL) .master(local[*]) // 本地模式用所有 CPU 核 .config(spark.sql.adaptive.enabled, true) .config(spark.sql.adaptive.coalescePartitions.enabled, true) .getOrCreate() import spark.implicits._ // 读取 JSON 日志自动推断 schema但需显式 cast 时间字段 val rawDF spark.read .option(multiLine, false) // Spark 3.3 支持 multiline JSON但电商日志通常是单行 .json(simulated_logs.json) // 清洗过滤空 event_type、排除 test 用户、标准化时间 val cleanedDF rawDF .filter($event_type.isNotNull $event_type ! !$user_id.startsWith(test_)) .withColumn(event_time, to_timestamp($event_time)) .withColumn(dt, date_format($event_time, yyyyMMdd)) // 写入 Parquet 分区表注意partitionBy 必须在 write 前调用 cleanedDF .write .mode(overwrite) .partitionBy(dt) .parquet(./data/ods_user_event) println(sETL completed. Total records: ${cleanedDF.count()}) spark.stop() } }提交命令spark-submit \ --class OdsEventETL \ --master local[4] \ --driver-memory 2g \ --executor-memory 2g \ --conf spark.sql.adaptive.enabledtrue \ etl_job.jar注意--master local[4]表示本地启动 4 个 executor模拟小规模集群--driver-memory 2g防止 OOMJSON 解析占内存spark.sql.adaptive.enabledtrue在本地模式下同样生效能自动合并小 task。3. 四类核心指标的 Spark SQL 实现从漏斗到 RFM3.1 跳失率与路径转化漏斗用 Window Function 计算会话内首尾行为跳失率定义为「只访问首页且无后续动作的会话占比」。关键点在于不能简单 countpage_url like %home%而要判断会话中是否只有 home 页且无 click/add_cart/buy。Spark SQL 的window函数配合collect_list可高效实现-- 创建临时视图便于调试 CREATE OR REPLACE TEMP VIEW ods_user_event AS SELECT * FROM parquet../data/ods_user_event; -- 计算每个会话的行为序列 WITH session_events AS ( SELECT session_id, collect_list(struct(event_type, page_url, item_id)) OVER ( PARTITION BY session_id ORDER BY event_time ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING ) AS event_seq, min(event_time) AS session_start, max(event_time) AS session_end FROM ods_user_event GROUP BY session_id ), -- 标记跳失会话序列长度1 且唯一事件是 view home bounce_sessions AS ( SELECT session_id FROM session_events WHERE size(event_seq) 1 AND event_seq[0].event_type view AND event_seq[0].page_url LIKE %home% ) SELECT round(count(*) * 100.0 / (SELECT count(*) FROM session_events), 2) AS bounce_rate_percent FROM bounce_sessions;执行结果应返回12.34模拟数据中约 12% 跳失。此 SQL 的关键参数是ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING——它确保collect_list获取整个会话所有事件而非默认的CURRENT ROW。3.2 复购率用自连接 时间窗口识别重复购买行为复购率 在 T 日前购买过 ≥2 次的用户数/T 日总购买用户数。难点在于1需排除同一订单多次支付2需限定时间窗口如 30 天内。Spark SQL 中lag()函数结合date_sub()可精准实现-- 先提取所有 buy 事件去重订单号 CREATE OR REPLACE TEMP VIEW buy_events AS SELECT DISTINCT user_id, event_time, session_id FROM ods_user_event WHERE event_type buy; -- 计算每个用户的上次购买时间 WITH user_buy_history AS ( SELECT user_id, event_time, lag(event_time, 1) OVER (PARTITION BY user_id ORDER BY event_time) AS last_buy_time FROM buy_events ), -- 判定是否为复购本次购买距上次 ≤30 天 rebuy_users AS ( SELECT DISTINCT user_id FROM user_buy_history WHERE last_buy_time IS NOT NULL AND datediff(event_time, last_buy_time) 30 ) SELECT round(count(DISTINCT ru.user_id) * 100.0 / count(DISTINCT be.user_id), 2) AS repurchase_rate_percent FROM buy_events be LEFT JOIN rebuy_users ru ON be.user_id ru.user_id;提示datediff(event_time, last_buy_time)返回整数天数比event_time - last_buy_time interval 30 days更稳定避免时区问题DISTINCT在rebuy_users中必须否则同一用户多次复购会被重复计数。3.3 RFM 分群用 Agg Case When 构建用户价值矩阵RFMRecency-Frequency-Monetary是电商最基础的用户分层模型。Spark 中无需 UDF纯 SQL 即可实现维度计算逻辑字段名Recency最近一次购买距今多少天recency_daysFrequency近90天购买次数frequencyMonetary近90天总消费金额此处用会话数代替monetary-- 假设 buy_events 已存在见 3.2 WITH user_rfm AS ( SELECT user_id, -- Recency: 当前日期减去最近购买时间用 datediff 避免 timestamp 直接减 datediff(current_date(), max(event_time)) AS recency_days, -- Frequency: 近90天购买次数 count(*) AS frequency, -- Monetary: 近90天会话数实际项目中应 join 订单表 sum(amount) count(DISTINCT session_id) AS monetary FROM buy_events WHERE event_time date_sub(current_date(), 90) GROUP BY user_id ), -- 用 quartile 划分 R/F/M 三级1高价值3低价值 rfm_score AS ( SELECT *, ntile(3) OVER (ORDER BY recency_days DESC) AS r_score, -- R 越小越好故倒序 ntile(3) OVER (ORDER BY frequency ASC) AS f_score, -- F 越大越好故正序 ntile(3) OVER (ORDER BY monetary ASC) AS m_score -- M 越大越好 FROM user_rfm ) SELECT case when r_score 1 and f_score 1 and m_score 1 then 重要价值客户 when r_score 1 and f_score in (1,2) and m_score in (1,2) then 重要发展客户 when r_score 2 and f_score 1 and m_score 1 then 重要保持客户 else 一般客户 end AS rfm_segment, count(*) AS user_count FROM rfm_score GROUP BY 1 ORDER BY user_count DESC;此 SQL 输出四类客户群数量分布可直接对接 BI 工具做热力图。4. 生产环境关键参数调优内存、Shuffle、Skew Join 的三处必改配置4.1 Spark 内存分配的黄金比例Driver 与 Executor 的 2:8 分配法则电商行为分析任务常因java.lang.OutOfMemoryError: GC overhead limit exceeded失败。根本原因不是总内存不足而是Executor 堆外内存Off-Heap Memory被 Shuffle 占满。Spark 3.3 默认spark.memory.fraction0.6堆内内存占比但电商场景需将spark.memory.storageFraction从 0.5 降至 0.3为 Shuffle 留足空间spark-submit \ --driver-memory 4g \ --executor-memory 16g \ --conf spark.memory.fraction0.6 \ --conf spark.memory.storageFraction0.3 \ # 关键降低 storage 缓存占比 --conf spark.shuffle.spill.enabledtrue \ # 强制 spill 到磁盘防 OOM --conf spark.serializerorg.apache.spark.serializer.KryoSerializer \ your_job.jar提示spark.serializerKryoSerializer可减少 30% 序列化体积尤其对嵌套 JSON 字段有效spark.shuffle.spill.enabledtrue是底线配置即使内存充足也建议开启——电商日志中page_url字段平均长度 120 字符极易触发 spill。4.2 Shuffle 分区数设置spark.sql.adaptive.coalescePartitions.enabled的真实效果电商日志经groupBy session_id后常产生上万个小分区每个会话一个分区导致 task 数爆炸。Spark 3.3 的 AQE 自动合并分区功能需显式开启并设阈值-- 在 Spark SQL 中设置或在 spark-submit 中 --conf SET spark.sql.adaptive.enabledtrue; SET spark.sql.adaptive.coalescePartitions.enabledtrue; SET spark.sql.adaptive.coalescePartitions.initialPartitionNum200; -- 初始分区数 SET spark.sql.adaptive.coalescePartitions.minPartitionSize64MB; -- 小于该值的分区被合并验证方法运行df.explain()查看 Physical Plan若出现AdaptiveSparkPlan且CoalescePartitions节点则配置生效。实测显示10 万会话日志的groupBy session_id任务task 数从 12000 降至 217 个。4.3 处理数据倾斜的终极方案Salting Map-Side Join当user_id分布极度不均如 1% 用户产生 60% 行为join会卡在少数 task。此时不能只靠spark.sql.adaptive.skewJoin.enabledtrue它仅对已知倾斜 key 有效而要用盐值Salting主动打散// 对大表用户行为表加盐 val saltedEvents events .withColumn(salt, when($user_id user_000001, (rand * 10).cast(int)).otherwise(lit(0))) .withColumn(salted_user_id, concat($user_id, lit(_), $salt)) // 对小表用户画像表复制 10 份每份加不同盐值 val saltedProfiles profiles .crossJoin((0 to 9).map(i lit(i)).toDF(salt)) .withColumn(salted_user_id, concat($user_id, lit(_), $salt)) // join 时用 salted_user_id val joined saltedEvents.join(saltedProfiles, salted_user_id)此方案将倾斜 keyuser_000001拆成 10 个user_000001_0~user_000001_9使负载均匀分布。生产环境中该方法将倾斜 job 运行时间从 42 分钟降至 3.8 分钟。5. 验证分析结果准确性的三步法抽样比对、Schema 检查、增量一致性校验5.1 用 Spark 自带的sample(withReplacement, fraction)做快速人工核验自动化测试前先对结果表抽样检查逻辑是否正确。例如验证复购率计算# sample_check.py from pyspark.sql import SparkSession spark SparkSession.builder.appName(SampleCheck).getOrCreate() df spark.read.parquet(./data/dwd_rebuy_users) # 抽取 0.1% 样本withReplacementFalse 保证不重复 sample_df df.sample(False, 0.001) # 输出前 10 行的 user_id 和 last_buy_time人工比对是否满足「距今≤30天」 sample_df.select(user_id, last_buy_time, event_time).show(10, truncateFalse)注意sample(False, 0.001)比limit(100)更科学——后者可能只取头部数据而抽样能覆盖全量分布。5.2 Schema 兼容性检查防止字段类型变更导致下游解析失败电商日志 schema 可能随业务迭代增加字段如新增ab_test_group但旧代码若用select *会出错。用printSchema()并导出为 JSON 校验# 导出当前表 schema spark-submit \ --conf spark.sql.adaptive.enabledfalse \ --driver-class-path $SPARK_HOME/jars/spark-sql_2.12-3.3.4.jar \ --class org.apache.spark.sql.util.SchemaPrinter \ $SPARK_HOME/jars/spark-sql_2.12-3.3.4.jar \ ./data/ods_user_event \ current_schema.json对比新旧current_schema.json重点关注event_time是否仍为timestamp类型、item_id是否从string变为nullable string——这些变更会影响where item_id ! 的结果。5.3 增量任务的幂等性验证用input_file_name()函数定位重复处理文件当使用spark.readStream处理 Kafka 日志时若 checkpoint 丢失可能重复消费。验证方法是在写入前添加源文件名标记val streamDF spark.readStream .format(kafka) .option(kafka.bootstrap.servers, kafka:9092) .option(subscribe, user_events) .load() // 添加 input_file_name() 作为溯源字段虽 Kafka 无文件但可用 offset 代替 val tracedDF streamDF .withColumn(source_offset, $offset) .withColumn(process_time, current_timestamp()) tracedDF.writeStream .format(parquet) .option(path, ./data/dwd_stream_events) .option(checkpointLocation, ./checkpoints/dwd_stream) .start()然后查询select source_offset, count(*) from dwd_stream_events group by source_offset having count(*) 1—— 若有结果说明该 offset 被处理了多次需检查 checkpoint 目录权限或 Kafka consumer group 配置。真正的准确性保障不靠单次运行正确而靠这三步形成的闭环抽样确认逻辑、Schema 锁定结构、增量校验幂等。电商数据团队每天凌晨 2 点跑完昨日行为分析后运维脚本自动执行这三项检查任一失败则钉钉告警并暂停下游报表生成。本文还有配套的精品资源点击获取
返回列表