據(jù)挖掘項目實戰(zhàn):從寬表建模到RFM標(biāo)簽計算)
簡介這份源碼資源面向大數(shù)據(jù)與電商方向的開發(fā)者、數(shù)據(jù)挖掘?qū)W習(xí)者提供一套基于Spark的電商用戶畫像完整實現(xiàn)可用于理解用戶行為分析、標(biāo)簽體系構(gòu)建與個性化推薦的技術(shù)落地路徑。壓縮包共462個文件約13.45MB以296個class編譯文件、70個scala源文件、20個java源文件為核心輔以properties、xml、json等配置與數(shù)據(jù)交換文件以及js、css、html構(gòu)成的前端展示層另有jar包、字體與圖標(biāo)等資源整體結(jié)構(gòu)完整。目錄按tags-model、tags-web、tags_ml、tags-etl等模塊劃分覆蓋數(shù)據(jù)抽取轉(zhuǎn)換、模型訓(xùn)練到前端呈現(xiàn)的全流程便于按模塊研讀。目前已有339人學(xué)習(xí)下載適合希望掌握Spark分布式數(shù)據(jù)處理與用戶畫像建模的讀者參考借鑒。1. 從一份「基于Spark的電商用戶畫像數(shù)據(jù)挖掘項目源碼」說起它到底能跑出什么電商后臺每天沉淀的原始數(shù)據(jù)其實很樸素訂單表、商品表、用戶表、行為埋點日志。真正讓運營團(tuán)隊頭疼的不是數(shù)據(jù)量而是「同一個用戶在不同表里長得不一樣」——訂單里他是收貨手機號埋點里他是設(shè)備 ID注冊表里他又成了會員編號。所謂電商用戶畫像本質(zhì)就是把這些散落的身份線索收斂成一張寬表再在這張寬表上算出 RFM、品類偏好、價格敏感度、活躍分層這些標(biāo)簽。而 Spark 在這里的價值是把原本要跑一整夜的 Hive 批處理壓縮到幾十分鐘并且用 DataFrame / Spark SQL 把清洗、關(guān)聯(lián)、聚合、標(biāo)簽計算串成一條可調(diào)度的流水線。這份「基于Spark的電商用戶畫像數(shù)據(jù)挖掘項目源碼」適合兩類人一類是剛接觸 Spark、想找一個完整鏈路練手的數(shù)據(jù)開發(fā)新手另一類是手里有真實電商數(shù)據(jù)、想搭一套可落地標(biāo)簽體系的工程師。它不解決推薦算法本身也不做實時流核心是把離線畫像的工程骨架講清楚。下面我按自己實際搭過的順序從環(huán)境、數(shù)據(jù)建模、標(biāo)簽計算一路講到踩過的坑能抄的地方直接給代碼。2. 環(huán)境與數(shù)據(jù)底座Spark 集群怎么搭、數(shù)據(jù)從哪來2.1 單機偽分布式先跑通再談集群很多人一上來就想搞 Spark 集群搭建結(jié)果卡在 YARN 資源隊列上三天沒跑通一條 SQL。我的建議是先用 local 模式把邏輯跑對再遷移到 standalone 或 on YARN。偽分布式最小依賴只有 JDK、Scala、Spark 三樣Hadoop 可以后補。# 以 Spark 3.x 為例解壓后配置環(huán)境變量 tar -zxvf spark-3.x-bin-hadoop3.tgz -C /opt/ echo export SPARK_HOME/opt/spark-3.x-bin-hadoop3 ~/.bashrc echo export PATH$SPARK_HOME/bin:$PATH ~/.bashrc source ~/.bashrc # local 模式驗證注意 master 用 local[*] 吃滿本機核 spark-shell --master local[*]邏輯說明local[*]表示用本機所有可用核跑一個 driver 加多個 executor 線程適合開發(fā)調(diào)試。參數(shù)上真正影響性能的是spark.executor.memory和spark.sql.shuffle.partitions后者默認(rèn) 200小數(shù)據(jù)量下會產(chǎn)生大量空任務(wù)本地調(diào)試建議改成 8 或 16。生產(chǎn)集群則相反shuffle 分區(qū)數(shù)要按數(shù)據(jù)量放大否則單分區(qū)數(shù)據(jù)傾斜會拖垮整個 stage。提示本地調(diào)試時把spark.sql.shuffle.partitions調(diào)小能顯著減少小文件和小任務(wù)開銷上集群前記得改回去。2.2 電商數(shù)據(jù)的三張核心表與埋點日志畫像項目的數(shù)據(jù)源通常分四塊用戶注冊表user_id、注冊時間、渠道、訂單表order_id、user_id、金額、下單時間、商品類目、商品表item_id、類目、價格帶、行為日志曝光、點擊、加購、收藏。行為日志一般是 JSON 格式Spark 讀取 JSON 是高頻操作也是熱搜里常被問到的點。from pyspark.sql import SparkSession from pyspark.sql.functions import col, from_json from pyspark.sql.types import StructType, StringType, LongType spark SparkSession.builder \ .appName(ecommerce_profile) \ .config(spark.sql.shuffle.partitions, 16) \ .getOrCreate() # 行為日志 schema顯式聲明比 inferSchema 快且穩(wěn) log_schema StructType() \ .add(user_id, StringType()) \ .add(item_id, StringType()) \ .add(event, StringType()) \ .add(ts, LongType()) logs spark.read.schema(log_schema).json(hdfs:///data/behavior/*.json) logs.createOrReplaceTempView(behavior_log)邏輯說明顯式 schema 避免了inferSchemaTrue觸發(fā)的全量掃描在日志量大時差距非常明顯。ts用 LongType 存毫秒時間戳后續(xù)做時間窗口聚合比字符串解析快。參數(shù)上讀取路徑用通配符*.json讓 Spark 按文件切分 task單文件別太小否則會掉進(jìn)小文件陷阱。2.3 寬表建模把用戶身份收斂成一行畫像的第一步是「用戶對齊」。訂單表用 user_id行為日志可能只有設(shè)備號注冊表又有手機號。常見做法是維護(hù)一張映射表把設(shè)備號、手機號、會員號統(tǒng)一映射到 user_id。這一步做不干凈后面所有標(biāo)簽都是錯的。-- 用注冊表作為主表左連接訂單和行為收斂到 user_id 粒度 CREATE OR REPLACE TEMP VIEW user_base AS SELECT u.user_id, u.register_time, u.channel, COUNT(DISTINCT o.order_id) AS order_cnt, SUM(o.amount) AS total_amount, MAX(o.order_time) AS last_order_time FROM user_register u LEFT JOIN orders o ON u.user_id o.user_id GROUP BY u.user_id, u.register_time, u.channel;邏輯說明以注冊表為主表保證每個用戶至少有一行左連接避免丟用戶。COUNT(DISTINCT)在數(shù)據(jù)傾斜時是性能殺手如果訂單表里同一 user_id 重復(fù)度極高可以先按 user_id 預(yù)聚合再關(guān)聯(lián)。參數(shù)上GROUP BY的字段越少 shuffle 數(shù)據(jù)量越小但維度丟了后面補不回來這里保留渠道是為了算渠道質(zhì)量標(biāo)簽。3. 標(biāo)簽計算RFM、偏好與分層的 Spark 實現(xiàn)3.1 RFM 三個指標(biāo)怎么算才不翻車RFM 是畫像里最經(jīng)典也最容易算錯的標(biāo)簽。R最近一次消費、F消費頻次、M消費金額看著簡單坑在于時間基準(zhǔn)和統(tǒng)計窗口。我一般用「當(dāng)前日期減去最后下單日期」算 R用近 90 天窗口算 F 和 M而不是全歷史否則老用戶會被歷史大額訂單永久拉高。from pyspark.sql.functions import datediff, current_date, count, sum, when rfm spark.sql( SELECT user_id, MAX(order_time) AS last_order_time, COUNT(order_id) AS freq_90d, SUM(amount) AS amount_90d FROM orders WHERE order_time date_sub(current_date(), 90) GROUP BY user_id ) rfm rfm.withColumn(recency, datediff(current_date(), col(last_order_time))) \ .withColumn(r_score, when(col(recency) 7, 5) .when(col(recency) 30, 4) .when(col(recency) 60, 3) .when(col(recency) 90, 2).otherwise(1))邏輯說明datediff返回天數(shù)差比手寫時間戳相減可讀性好。分檔閾值 7/30/60/90 是經(jīng)驗值不同品類要調(diào)快消品可以壓到 3/7/15/30。參數(shù)上窗口用date_sub(current_date(), 90)而不是寫死日期保證每天調(diào)度時自動滾動。注意current_date()依賴集群時區(qū)跨時區(qū)業(yè)務(wù)要顯式指定。3.2 品類偏好標(biāo)簽用 explode 打散再聚合用戶偏好哪個類目不能只看訂單還要結(jié)合行為日志的點擊和加購。常見做法是把行為日志里的 item_id 關(guān)聯(lián)商品表拿到類目再按 user_id 類目聚合打分最后取 top1 或 top3。from pyspark.sql.functions import explode, split, collect_list, struct, desc # 行為日志按事件加權(quán)點擊1分加購3分下單5分 weighted spark.sql( SELECT b.user_id, i.category, SUM(CASE b.event WHEN click THEN 1 WHEN cart THEN 3 WHEN order THEN 5 ELSE 0 END) AS score FROM behavior_log b JOIN items i ON b.item_id i.item_id GROUP BY b.user_id, i.category ) # 取每個用戶得分最高的前3個類目 pref weighted.groupBy(user_id) \ .agg(collect_list(struct(category, score)).alias(cats)) \ .withColumn(top_cats, expr(slice(array_sort(cats, (l, r) - r.score - l.score), 1, 3)))邏輯說明加權(quán)打分把不同行為的重要性區(qū)分開比單純計數(shù)更貼近真實偏好。array_sort配合 lambda 按 score 降序排slice取前三。參數(shù)上權(quán)重 1/3/5 是常見起點如果加購轉(zhuǎn)化率低可以調(diào)高 cart 權(quán)重。注意collect_list在單用戶類目極多時會撐爆內(nèi)存必要時先過濾低分項。3.3 用戶分層把標(biāo)簽落成可運營的群體標(biāo)簽算完要能落到運營動作上否則就是一堆數(shù)字。常見分層是「高價值活躍」「高價值流失」「低價值活躍」「沉睡」四象限用 R 和 M 交叉即可。layered rfm.withColumn(segment, when((col(r_score) 4) (col(amount_90d) 1000), 高價值活躍) .when((col(r_score) 2) (col(amount_90d) 1000), 高價值流失) .when((col(r_score) 4) (col(amount_90d) 1000), 低價值活躍) .otherwise(沉睡用戶)) layered.write.mode(overwrite).parquet(hdfs:///data/user_profile/segment)邏輯說明分層規(guī)則用when/otherwise鏈?zhǔn)奖磉_(dá)清晰且易改。金額閾值 1000 要按業(yè)務(wù)客單價定不能照搬。寫出用 parquet 列式存儲后續(xù) BI 查詢只讀需要的列。參數(shù)上mode(overwrite)適合全量重算增量場景應(yīng)改成按分區(qū)覆蓋。4. 避坑與排查畫像項目里最容易翻車的五件事4.1 數(shù)據(jù)傾斜導(dǎo)致個別 task 跑幾小時現(xiàn)象Spark UI 里某個 stage 的少數(shù) task 耗時遠(yuǎn)超其他shuffle read 數(shù)據(jù)量差幾十倍。原因熱門商品或大 V 用戶的行為日志集中在少數(shù) key 上。解決先對熱點 key 加隨機前綴打散聚合后再去掉前綴或者對傾斜 key 單獨用 broadcast join 處理。4.2 JSON 解析出 null 卻不報錯現(xiàn)象行為日志讀進(jìn)來大量字段為 null任務(wù)正常結(jié)束但結(jié)果全空。原因schema 和實際 JSON 字段名或類型不匹配Spark 默認(rèn)把解析失敗置 null 而不拋異常。解決讀取時加modePERMISSIVE并配合columnNameOfCorruptRecord把壞數(shù)據(jù)單獨落盤排查別讓它靜默丟失。4.3 shuffle 分區(qū)數(shù)沒調(diào)產(chǎn)出上千小文件現(xiàn)象寫 parquet 后目錄里幾千個幾十 KB 的小文件下游 Hive 查詢慢。原因spark.sql.shuffle.partitions默認(rèn) 200數(shù)據(jù)量小的時候每個分區(qū)只寫一點點。解決按數(shù)據(jù)量估算分區(qū)數(shù)寫出前用coalesce或repartition收斂單文件控制在 128MB 左右。4.4 時間窗口用錯時區(qū)R 值集體偏一天現(xiàn)象凌晨調(diào)度時算出的 recency 比預(yù)期多 1。原因current_date()取的是集群默認(rèn)時區(qū)和業(yè)務(wù)時區(qū)不一致。解決統(tǒng)一在 SQL 里用from_utc_timestamp轉(zhuǎn)換或者調(diào)度參數(shù)里顯式傳入業(yè)務(wù)日期別依賴隱式時區(qū)。4.5 全量重算沒做冪等重跑產(chǎn)生重復(fù)數(shù)據(jù)現(xiàn)象任務(wù)失敗重跑后畫像表里同一 user_id 出現(xiàn)多行。原因?qū)懗鲇?append 模式且沒有按分區(qū)覆蓋。解決分區(qū)表用insert overwrite指定分區(qū)或者寫出前按主鍵去重保證重跑結(jié)果一致。5. 進(jìn)階技巧讓畫像任務(wù)從「能跑」到「跑得省」真正把畫像項目跑進(jìn)生產(chǎn)后你會發(fā)現(xiàn)瓶頸往往不在算法而在資源調(diào)度和數(shù)據(jù)組織。分享幾個我踩坑后固定下來的習(xí)慣。第一個是緩存復(fù)用。RFM 和品類偏好都要讀訂單表如果中間結(jié)果會被多次引用果斷cache()但用完記得unpersist()否則 executor 內(nèi)存被占滿后續(xù) stage 頻繁 spill。判斷標(biāo)準(zhǔn)很簡單看 Spark UI 里某個 DataFrame 是否被多個 action 觸發(fā)。第二個是廣播小表。商品表通常幾萬行關(guān)聯(lián)行為日志時用broadcast()提示能避免一次大 shuffle。參數(shù)上spark.sql.autoBroadcastJoinThreshold默認(rèn) 10MB商品表超過這個值就手動 broadcast別硬等自動判斷。第三個是分區(qū)裁剪。畫像表按日期分區(qū)查詢時一定帶上分區(qū)條件否則全表掃描。我見過有人寫WHERE dt 2024-01-01卻因為格式不匹配導(dǎo)致分區(qū)失效血淚經(jīng)驗是分區(qū)字段類型和查詢字面量必須一致。調(diào)優(yōu)項默認(rèn)值建議值適用場景shuffle.partitions200數(shù)據(jù)量GB×2中大集群executor.memory1g4g~8g聚合密集autoBroadcastJoinThreshold10MB30MB小維表關(guān)聯(lián)serializerJavaKryo全場景最后說驗證方法。畫像結(jié)果不能只看任務(wù)成功要抽樣核對隨機抽 100 個 user_id手工從訂單表算一遍 RFM和產(chǎn)出表比對。差異超過 5% 就說明邏輯或數(shù)據(jù)有問題。這個習(xí)慣幫我攔下過好幾次「任務(wù)綠了但結(jié)果是錯的」的翻車。我自己現(xiàn)在做任何畫像任務(wù)第一件事不是寫 SQL而是先把數(shù)據(jù)源的口徑和時區(qū)確認(rèn)清楚再動手。希望幫到你。本文還有配套的精品資源點擊獲取