简介:作为一份山西大学研究生项目设计报告的整理版本,这份《大数据实例:网站用户行为分析》以网站用户行为分析为实战场景,完整串起大数据处理中数据预处理、存储、查询和可视化分析的核心链路。文档按五个实验步骤推进:先完成实验环境准备,逐步搭建操作系统、关系型数据库、分布式计算框架、列族数据库、数据仓库、数据传输工具、统计绘图环境等组件;再将本地数据集经过下载、解压和预处理后上传至数据仓库,利用查询分析完成业务指标提取;随后通过数据同步工具完成数据仓库、关系型数据库与列族数据库之间的相互导入导出;最后使用统计语言读取关系型数据库中的结果,绘制可视化图表。除了完整流程,文档还给出了各类组件安装时的配置思路,以及编码导致的中文乱码、连接权限等常见问题的处理办法。对于需要完成大数据课程设计、实验报告或入门大数据生态的读者,是一份步骤清晰、可照做的项目式参考。资源包内为1个docx文档,大小36KB,内容精炼,便于阅读和二次整理;目前平台上已有1045人学习下载,适合边看边练、逐步还原实验环境。
1. 从一份“网站用户行为分析”文档说起,先弄清楚要算哪些指标
第一次拿到“网站用户行为分析”这类文档,很容易被“大数据”三个字带偏方向,以为重点在Spark参数和集群规模,真正决定交付质量的却是埋点字段和session口径。用户从落地页进入、搜索、加购、下单,这一串行为如果采集不全、时间戳对不齐、登录前后身份没打通,后面任何分析都算不出可信指标。下面按一条我自己跑过多遍的链路来讲:埋点与采集、数仓分层、漏斗留存路径三类指标计算,以及结果上大屏前的校验。适合正在搭用户行为数据平台的数据工程师、数仓开发和分析师,照着章节顺序就能把离线分析链路跑通。
2. 从埋点到Kafka:网站用户行为日志的采集链路设计
用户行为分析的第一步不是建表,而是定义“一次行为”到底长什么样。网站场景下有两类常见的采集方式:前端埋点通过JS SDK上报页面浏览、按钮点击、停留时长;后端埋点通过Nginx access log记录每个HTTP请求。很多团队最后是两者混用,前端上报视角事件,后端记录服务端确认的搜索、加购、支付等关键动作。我一般建议把两端事件统一成一张事件表,至少包含下面要讲的公共字段。
2.1 行为事件公共字段与业务字段的分层设计
公共字段是所有事件都有的,它们的取值质量直接决定后面的去重和session划分是否靠谱。下面这张表是简化后的字段清单,实际项目中还会加app_id、channel等维度字段,但核心是这几个。
| 字段名 | 含义 | 说明 |
|---|---|---|
| event_id | 事件唯一ID | 前端SDK生成UUID,后端事件复用同一规则 |
| user_id | 登录用户ID | 未登录时为NULL,配合visitor_id做识别 |
| visitor_id | 访客ID | 由Cookie或浏览器指纹生成,跨会话识别关键 |
| session_id | 会话ID | 可在采集端生成,也可在数仓按规则回填 |
| page_url | 页面URL | 注意URL中query参数是否需要保留 |
| referer_url | 来源页 | 用于来源渠道分析 |
| event_time | 事件时间 | 统一为yyyy-MM-dd HH:mm:ss,时区用UTC+8 |
| event_type | 事件类型 | page_view / click / search / add_cart / order |
| device_info | 设备信息 | UA字符串,后续按正则解析出os、browser |
这张表里最容易被忽视的是visitor_id。如果只用user_id,未登录用户的浏览行为全丢;如果只用visitor_id,登录前后的行为无法贯通。常见做法是在数仓计算时生成一个分析主键:user_id不为空时优先取user_id,否则取visitor_id。这个逻辑放DWD层做,后面所有分析都基于这个analysis_id。
2.1.1 前端埋点与后端日志的字段对齐
前端埋点字段自由,但依赖浏览器端定时或页面卸载时上报,有一定丢失率。后端Nginx日志字段固定,丢失率低,但看不到页面内交互。把两边字段对齐时,我一般先列出要从Nginx access log里拿到的原始字段,再写解析逻辑统一转成上面表格里的字段命名。最省事的方案是直接从Nginx配置这层就输出JSON,省掉一层正则解析。
2.2 Nginx access log的JSON化采集与Kafka接入
采集端常见的套路是Filebeat监控Nginx日志文件,增量发送到Kafka。Nginx日志格式一般这样定义,直接把JSON序列化做在log_format里:
log_format user_behavior escape=json '{' '"remote_addr":"$remote_addr",' '"time_local":"$time_local",' '"request_uri":"$request_uri",' '"status":"$status",' '"http_referer":"$http_referer",' '"http_user_agent":"$http_user_agent",' '"cookie_visitor":"$cookie_visitor"' '}'; server { listen 80; server_name example.com; access_log /var/log/nginx/user_behavior.log user_behavior; location / { proxy_pass http://backend; } }关键点是escape=json这个参数,Nginx 1.11.0之后才支持,它会把双引号和转义字符自动处理成合法JSON,采集端直接按JSON解析,比正则切字段稳得多。$cookie_visitor是前置JS种在Cookie里的访客ID,需要业务页面里提前执行setCookie逻辑。如果还没做前端SDK,最快速的起点就是先把后端access log的JSON格式定好。
日志文件有了,采集配置用Filebeat,下面是filebeat.yml的关键段:
filebeat.inputs: - type: filestream enabled: true paths: - /var/log/nginx/user_behavior.log parsers: - ndjson: ~ fields: log_topic: user_behavior_log output.kafka: hosts: ["kafka-1:9092", "kafka-2:9092"] topic: user_behavior_log partition.round_robin: reachable_only: true required_acks: 1 compression: gzipparsers里的ndjson告诉Filebeat每行按JSON解析,fields里的log_topic是自定义字段,下游消费者根据它的值判断数据来源。output.kafka里的required_acks设成1,表示Kafka leader收到消息就算成功,吞吐比all高一截,日志类数据丢一条可接受。compression用gzip是因为行为日志重复字段多,压缩比通常有60%以上,能显著降低Kafka带宽压力。
2.3 事件流进入Kafka后的质量校验
日志进入Kafka后通常还要做一层轻量校验,这层可以用Flink SQL消费原始topic,过滤后再写入clean topic,下面是一段能放进作业的Flink SQL:
CREATE TABLE kafka_source ( event_id STRING, visitor_id STRING, user_id STRING, event_type STRING, event_time STRING, pt STRING ) WITH ( 'connector' = 'kafka', 'topic' = 'user_behavior_log', 'properties.bootstrap.servers' = 'kafka-1:9092', 'format' = 'json', 'scan.startup.mode' = 'latest-offset' ); CREATE TABLE kafka_clean ( event_id STRING, visitor_id STRING, user_id STRING, event_type STRING, event_time STRING, pt STRING ) WITH ( 'connector' = 'kafka', 'topic' = 'user_behavior_log_clean', 'properties.bootstrap.servers' = 'kafka-1:9092', 'format' = 'json' ); INSERT INTO kafka_clean SELECT event_id, visitor_id, user_id, event_type, event_time, pt FROM kafka_source WHERE event_id IS NOT NULL AND event_time IS NOT NULL AND event_type IN ('page_view','click','search','add_cart','order');这一段做的事情是丢弃event_id或event_time为空的记录,同时把事件类型限定在业务关心的五种里。像爬虫流量和内部测试流量,我建议在这一层先保留,到DWD层再按UA特征和来源IP过滤,因为过滤规则经常变,放在清洗层比采集层好维护。
3. 用户行为数仓分层:Hive建表与DWD层ETL落地
埋点和采集解决“有没有数据”,数仓解决“数据能不能直接被分析师使用”。网站用户行为分析一般建三层:ODS层原样存放Kafka过来的事件日志,DWD层完成清洗、维度补全、会话划分,DWS层按天聚合各类指标。第2章停在了Kafka的clean topic,这一章从Kafka往Hive落ODS开始讲。
3.1 ODS层建表与分区策略
ODS表最重要的策略是分区。用户行为日志量大,按天分区最常见,单日超过亿级事件的项目还要考虑按小时分区。分区字段用pt,值为yyyy-MM-dd或yyyy-MM-dd-HH。下面给出ODS建表SQL:
CREATE EXTERNAL TABLE ods_user_event_log ( event_id STRING COMMENT '事件UUID', user_id STRING COMMENT '用户ID,可为空', visitor_id STRING COMMENT '访客ID', session_id STRING COMMENT '采集端预置会话ID', page_url STRING COMMENT '页面URL', referer_url STRING COMMENT '来源URL', event_type STRING COMMENT '事件类型', event_time STRING COMMENT '事件时间', device_info STRING COMMENT 'UA信息' ) PARTITIONED BY (pt STRING COMMENT '分区字段,yyyy-MM-dd') STORED AS parquet TBLPROPERTIES ('parquet.compression'='snappy');这里有几个容易被问到的点。第一,为什么用外部表:日志数据是采集系统先写到HDFS,再导入分区,外部表删表不影响原始文件。第二,为什么用parquet而不存text:行为日志按列读取,分析时往往只取部分字段,parquet的列裁切能少读大量数据。第三,event_time在ODS层保留string,合法性判断放到DWD层解析,ODS只负责原样落盘。
分层设计可以先用这张表说明,后面讨论时方便对齐:
| 分层 | 表名 | 粒度 | 说明 |
|---|---|---|---|
| ODS | ods_user_event_log | 一条事件一行 | 原样接入,按天分区 |
| DWD | dwd_user_event_log | 一条事件一行 | 清洗后事件,带analysis_id和session_seq |
| DWS | dws_user_event_daily | 一个用户一天一行 | 汇总PV、会话数、加购数、下单数 |
3.1.1 定时调度与分区幂等写入方式
工程落地上,很少手动执行LOAD命令,一般用DolphinScheduler或Airflow定时调一个Shell脚本消费Kafka写入分区。这里有一个关键点:如果任务失败重跑,要保证幂等。常见做法是先写临时目录,再对目标分区执行INSERT OVERWRITE,让第二次运行完全覆盖第一次结果,而不是在分区里追加文件。
3.2 DWD层ETL:时间清洗与30分钟会话切分
DWD层要做三件事:把ODS里的字符串event_time转成timestamp,过滤掉预发布环境的无效事件,然后给每个事件切出会话编号。会话切分的口径,业内用得最多的是“同一个analysis_id相邻两条事件间隔超过30分钟,就切一个新会话”。
INSERT OVERWRITE TABLE dwd_user_event_log PARTITION (pt = '2026-01-12') WITH base AS ( SELECT event_id, COALESCE(user_id, visitor_id) AS analysis_id, visitor_id, page_url, referer_url, event_type, FROM_UNIXTIME(UNIX_TIMESTAMP(event_time, 'yyyy-MM-dd HH:mm:ss')) AS event_ts, UNIX_TIMESTAMP(event_time, 'yyyy-MM-dd HH:mm:ss') AS ts FROM ods_user_event_log WHERE pt = '2026-01-12' AND page_url IS NOT NULL AND event_time IS NOT NULL ) SELECT event_id, analysis_id, visitor_id, page_url, referer_url, event_type, event_ts, SUM(IF(ts - LAG(ts, 1, ts) OVER w > 1800, 1, 0)) OVER w + 1 AS session_seq FROM base WINDOW w AS ( PARTITION BY analysis_id ORDER BY ts ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW );这段SQL的思路:先用LAG取上一条事件的时间戳,如果当前事件与上一条间隔超过1800秒,就记1个会话边界;再用SUM开窗把边界数累加起来,得到这个用户当前事件属于第几个会话。ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW表示累加范围从该用户第一行到当前行。最终的session_seq从1开始递增,用analysis_id加event_ts加session_seq就能拼出唯一会话ID,日常分析里直接拿session_seq做计数就行。
这里有个性能注意点:ORDER BY ts要求同一分析主键的数据能在一个reduce里处理,Hive会为这个窗口触发shuffle。用户量大时这个shuffle代价不低,可以先用DISTRIBUTE BY analysis_id SORT BY ts把数据重新组织一遍,再跑窗口函数,能减少单次查询的中间结果膨胀。另外,如果未来要支持实时会话划分,那就不能等离线跑批,得用Flink的keyed state在流上切session,离线跑批和实时job的切分口径必须保持一致,否则两个数据源对不上。
3.3 DWS层指标聚合与用户维度关联
DWD算完后,分析指标落在DWS层。通常按analysis_id、stat_date聚合,计算每天的PV、会话数、业务事件数。
INSERT OVERWRITE TABLE dws_user_event_daily PARTITION (pt = '2026-01-12') SELECT analysis_id, pt AS stat_date, COUNT(*) AS pv, COUNT(DISTINCT session_seq) AS session_cnt, COUNT(IF(event_type = 'add_cart', 1, NULL)) AS add_cart_cnt, COUNT(IF(event_type = 'order', 1, NULL)) AS order_cnt FROM dwd_user_event_log WHERE pt = '2026-01-12' GROUP BY analysis_id;COUNT(DISTINCT session_seq)对本案例合理,因为session_seq只在同一个analysis_id内部递增。COUNT(IF(...))是Hive里条件计数最常见的写法,不用CASE WHEN包一层。DWS层核心是“维度加指标”,维度可以是用户、页面、来源渠道,指标是PV、UV、会话数、各事件计数。数据量中等时这层放Hive就够,一旦要支撑秒级延迟就得换成Flink流式聚合,产物落到Doris或ClickHouse。对离线分析场景,ODS、DWD、DWS三层结构已经能覆盖绝大多数报表需求。
4. 漏斗、留存、路径三类分析的实现与口径对齐
有了数仓分层,接下来回答业务问题:用户从进入网站到下单,流失在哪一步?第二天还回不回来?不同来源渠道入口的行为路径差异在哪?这三类问题分别由漏斗分析、留存分析、路径分析回答,也是网站用户行为分析里几乎必写的三个章节。这章直接给可复现的SQL和Spark实现。
4.1 漏斗分析:从“事件独立计数”到“步骤有序判定”
漏斗分析要回答的是“每一步有多少人走到这里”。业务上常见误区是直接按event_type做group by,那样算出来的只是每个事件的独立人数,不是漏斗。正确做法是先把每个用户的事件按时间排序,再判断其是否依次经过关键步骤。下面的SQL用ROW_NUMBER给用户事件序列编号,再按是否存在关键事件判定是否进入对应步骤。
| 步骤 | 事件类型 | 含义 |
|---|---|---|
| step1 | page_view | 用户进入落地页 |
| step2 | search | 使用了站内搜索 |
| step3 | add_cart | 将商品加入购物车 |
| step4 | order | 生成订单 |
WITH ordered AS ( SELECT analysis_id, event_type, event_ts, ROW_NUMBER() OVER ( PARTITION BY analysis_id ORDER BY event_ts ) AS rn FROM dwd_user_event_log WHERE pt >= '2026-01-12' AND pt <= '2026-01-18' ), funnel AS ( SELECT analysis_id, MAX(IF(event_type = 'page_view', rn, NULL)) AS step1_rn, MAX(IF(event_type = 'search', rn, NULL)) AS step2_rn, MAX(IF(event_type = 'add_cart', rn, NULL)) AS step3_rn, MAX(IF(event_type = 'order', rn, NULL)) AS step4_rn FROM ordered GROUP BY analysis_id ) SELECT COUNT(*) AS base_users, COUNT(IF(step1_rn IS NOT NULL, 1, NULL)) AS step1_users, COUNT(IF(step2_rn IS NOT NULL, 1, NULL)) AS step2_users, COUNT(IF(step3_rn IS NOT NULL, 1, NULL)) AS step3_users, COUNT(IF(step4_rn IS NOT NULL, 1, NULL)) AS step4_users FROM funnel;这里的口径是“存在即算”:只要该用户全量事件里出现过目标事件,就算进入了对应步骤。大多数日常漏斗用这个口径能说明问题,但严格的业务漏斗要求步骤之间有时序关系。比如“先搜索后加购”和“先加购后搜索”在存在即算口径下都算通过,实际上业务只认前者。这时需要把漏斗SQL改成:在ranking层面就施加时序约束,让stepN只在stepN-1出现之后才被记录。实现方式可以按下面的方案做。
4.1.1 用路径字符串做严格有序漏斗
我在实际项目中处理“必须按顺序”的需求时,常用一个更直观的思路:把每个用户的事件类型按时间顺序拼成字符串,再用正则匹配判断是否存在“浏览→搜索→加购→下单”的子序列。下面是具体的SQL:
WITH user_path AS ( SELECT analysis_id, CONCAT_WS('>', COLLECT_LIST(event_type) OVER ( PARTITION BY analysis_id ORDER BY event_ts ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW )) AS event_path FROM dwd_user_event_log WHERE pt = '2026-01-12' ) SELECT COUNT(DISTINCT analysis_id) AS matched_users FROM user_path WHERE event_path REGEXP 'page_view(>.*)?>search(>.*)?>add_cart(>.*)?>order';COLLECT_LIST开窗函数按event_ts把用户事件拼成类似“page_view>search>add_cart>order”的字符串,ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW保证拼接范围从该用户第一条事件到当前行。REGEXP里的(>.*)?表示两个目标步骤之间可以插入任意其他事件。字符串拼接在单用户事件量很大时会占用较多内存,实际使用可以加个限制,只保留前100个事件再做路径匹配。
4.2 留存分析:按首次活跃日分组算次日与7日留存
留存分析按“首次活跃日”分组。1月12日第一次访问的用户里,1月13日还有行为的就算次日留存命中。Hive里用min(pt)取首次活跃日,再和每日活跃明细做join就能算出来:
WITH first_active AS ( SELECT analysis_id, MIN(pt) AS first_day FROM dwd_user_event_log WHERE pt >= '2026-01-01' AND pt <= '2026-01-31' GROUP BY analysis_id ), daily_active AS ( SELECT DISTINCT analysis_id, pt FROM dwd_user_event_log WHERE pt >= '2026-01-01' AND pt <= '2026-01-31' ) SELECT fa.first_day, COUNT(DISTINCT fa.analysis_id) AS new_users, COUNT(DISTINCT IF(DATEDIFF(da.pt, fa.first_day) = 1, fa.analysis_id, NULL)) AS day1_retained, COUNT(DISTINCT IF(DATEDIFF(da.pt, fa.first_day) = 7, fa.analysis_id, NULL)) AS day7_retained FROM first_active fa LEFT JOIN daily_active da ON fa.analysis_id = da.analysis_id GROUP BY fa.first_day;留存分析最常见的错误是口径偏移:有人用“DAU/新增用户数”来近似留存,会把一个在1月13日和1月20日各活跃一次的用户在两个窗口都算进去,拉高留存率。上面的写法用DATEDIFF精确限定第1天和第7天,避免这类偏差。另外注意,first_active基于“有行为事件”的用户,和用户注册表里的注册用户数不一定一致,如果业务要按注册口径算留存,需要把注册表join进来再聚合。
4.3 路径分析:用Spark统计用户行为转移Top序列
路径分析的目标是拿到用户从某个事件开始,最常见的下一步事件分布。最简方案是Spark读取DWD数据,按用户和时间排序后用滞后函数取下一事件类型,再groupBy统计频次:
import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions._ val windowSpec = Window .partitionBy("analysis_id") .orderBy("event_ts") val df = spark.table("dwd_user_event_log") .filter("pt = '2026-01-12'") val pairDf = df .withColumn("next_event", lag("event_type", -1).over(windowSpec)) .filter("next_event IS NOT NULL") .groupBy("event_type", "next_event") .count() .orderBy(desc("count")) pairDf.show(20)这段代码用lag函数的offset为-1取出按时间排序后的下一事件,得到“当前事件→下一事件”的二元组。Spark 3.x里也可以用lead函数直接取下一行,语义更直观。结果中频率最高的转移对通常是page_view→search和search→add_cart,它们直接反映了用户从进入到转化的主要路径。如果想看三步以上的序列,可以把连续n个事件用concat_ws拼成长串再groupBy,或者用mapWithState在流上实时计算转移矩阵。路径分析的结果对推荐策略和页面布局优化有直接参考价值,但离线版本只反映历史行为,要在新版页面灰度期间看实时转移变化,还是要接Flink。
5. 结果落ClickHouse与可视化大屏的固化交付
前三章算出的结果最终要让人看得懂。多数团队的交付形态是Hive里的DWS结果同步到ClickHouse,后端接口查询,前端用ECharts画可视化大屏,数据同步工具可选DataX或ClickHouse自带的JDBC写入。这一章不展开后端工程,只讲结果表结构和最容易踩的坑:同环比口径。
5.1 ClickHouse结果表结构与查询参数
同步到ClickHouse时,表结构按查询维度建。下面的DDL适合按用户加日期查询的场景:
CREATE TABLE ch_user_event_daily ( stat_date Date, analysis_id String, pv UInt64, uv UInt64, session_cnt UInt64, add_cart_cnt UInt64, order_cnt UInt64 ) ENGINE = MergeTree PARTITION BY stat_date ORDER BY (stat_date, analysis_id);MergeTree的ORDER BY决定了索引裁剪粒度,这里的查询条件固定为stat_date加analysis_id,按这两个字段排序能保证点查走主键索引。如果大屏接口还要支持按来源渠道筛选,建议把channel字段加进ORDER BY,否则筛选会退化成全表扫描。
5.2 三个可落地的结果校验技巧
结果上大屏前,我建议先做三个校验。第一,幂等校验:用同一个分区跑两次批任务,PV总数必须完全相等,若不等说明脚本里有非确定性计算,比如用了未固定种的随机函数或依赖了系统时间。第二,漏斗单调性:漏斗从step1到step4的转化人数必须单调递减,一旦出现中间层人数比上一层多,基本可以断定事件定义重复或去重键不一致。第三,明细抽样核对:随机抽100个用户在DWD层手工核对session_seq,确认没有漏掉超时边界和跨天边界。
其中第三点是最容易出问题的地方。跨天会话场景是这样的:用户在23:50进入网站,操作持续到次日00:10,按天分区会被拆成两个分区,如果session切分逻辑只按天独立计算,就会把同一会话算成两次。解决思路是ODS层记录timestamp类型的event_ts,切分session时按analysis_id和event_ts全局排序,再根据相邻事件间隔判断是否跨天合并。这套逻辑放在离线DWD层能做,但效率不如在Flink实时阶段直接切好session再落DWD,能省掉每小时调度对齐的时间损耗。大屏本身只是展示层,真正值钱的是口径一致、可复算的结果表,校验这一步建议做成定时任务,每天跑批后自动告警,比上线前手动检查一次可靠得多。
本文还有配套的精品资源,点击获取