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