☰
基于Hadoop+Spark+Django的网购用户购买力差异分析实战项目全解析
2026/10/2 10:24:36 网站建设 项目流程

你说一个用户有没有购买力,不能只看他一个月花多少钱。同样是月消费5000元,一个只买日用品和零食的用户,和一个经常买手机、相机、电脑的用户,商业价值完全不一样。这种朴素认知,放到大数据场景里就是“用户购买力差异分析”要解决的问题。

我最近完整跑通了一套基于Hadoop+Spark+Django的网购平台用户购买力差异分析项目,从原始订单数据生成、HDFS存储、Spark特征工程与聚类分群,到Django提供API、可视化大屏展示,全程源码可用、文档齐全、还能一步步断点调试。如果你正在做大数据方向的课程设计、毕业设计,或者想找一个“从数据底层架构到前端展示”的完整横向项目,这篇文章基本能让你少走一半弯路。我会把实战中容易卡住的地方都摊开讲,包括几个当时调试到半夜才解决的问题。

1. 为什么要用Hadoop+Spark+Django来做购买力差异分析

首先把这个项目要解决的事情说明白。网购平台用户的购买力差异,不是一个单一指标能描述的。常见定义有累计消费金额、消费频次、最近一次消费时间、平均客单价、购买品类宽度、是否偏好高价品类等。不同业务方对“购买力”的诉求不同,运营希望找到高价值用户做定向召回,风控希望识别异常偏高或偏低的消费行为。这就需要一个能从海量订单里提取用户级特征、再对用户进行分群的技术方案。

既然是“海量”,就意味着不能用Excel、不能用单机Pandas硬扛。这也是课程设计和毕业设计选这套架构的天然理由。Hadoop负责分布式存储,把订单数据放到HDFS上;Spark负责分布式计算,尤其适合K-Means这种迭代式算法;Django负责把分析结果沉淀成可查询的Web服务,最后前端大屏把差异直观呈现出来。三个组件各管一段,链路清晰,答辩的时候也很好讲。

我在环境选择上踩过一些坑,先给你一个已经在多台机器上验证过的组合:

组件推荐版本主要职责
Hadoop3.3.4HDFS存储原始订单与中间结果
Spark3.3.2特征提取、聚类分析、结果聚合
Python3.8+生成模拟数据、Django开发
Django3.2 LTS 或 4.1Web框架、REST API
ECharts5.x大屏可视化
MySQL5.7 / 8.0最终结果库(可选,SQLite也能跑)

基本上8G内存的笔记本就能跑通全流程,前提是数据量控制在200万条订单以内。数据再多,伪分布式调度器的压力就上来了,跑一次分析要等很久,不利于调试。这个限制不是架构不行,而是课程设计场景下的合理取舍。真正的大规模部署需要上YARN集群,但这个项目模型保持一致,换集群只是配置文件的事。

这里还涉及一个很关键的选型理由:为什么计算框架用Spark而不是MapReduce?因为购买力差异分析要经历“聚合特征→标准化→聚类→解读簇中心”几个阶段,尤其是K-Means,每一轮迭代都要遍历全量数据。MapReduce为了一个迭代就要写多个Job,每个Job都要落盘一次;Spark把中间结果放内存,迭代计算快一到两个数量级。而且pyspark的DataFrame API对Python用户友好,代码看起来像Pandas,学生上手快。如果你在课程设计里用MapReduce硬写K-Means,代码量和排查难度都会直线上升。

Django在这个链条里的作用经常被低估。很多人把Spark算完,输出一个CSV,就认为项目结束了。但完整课题往往要求“分析结果可查询、可展示”,这就需要Web服务层。选Django而不是Flask的理由是:Django自带ORM、Admin后台和项目结构规范,写一个带模型和接口的小型服务非常自然。而且Django的REST接口配合前端大屏做定时请求,比Flask多出来的那一点重量完全值得。

再补一句,这个项目的数据流可以概括成:Python脚本生成订单明细CSV → 上传HDFS → Spark读取并清洗 → 构建用户特征宽表 → KMeans分群 → 结果写回数据库 → Django读取并封装API → 可视化大屏定时拉取。这样一条链路,实际上覆盖了大数据项目从存储、计算到应用层的所有核心环节。下一步我们逐段把它拆开。

2. 订单数据入湖:从原始表到HDFS再到Spark DataFrame

2.1 先造一份“像样的”网购订单数据

真实电商数据拿不到,课程设计通常用参数化脚本生成模拟数据。别小看这一步,模拟数据的质量直接决定后面分析好不好看。我写生成脚本的时候,不是随机往表里灌数字,而是尽量模拟真实业务规律。

核心字段我建议这样设计:

  • user_id:用户ID,生成30000~50000个不同用户
  • order_id:全局唯一订单号
  • item_category:商品一级品类,包括“手机数码”“家用电器”“服饰鞋包”“食品生鲜”“美妆个护”
  • item_price / quantity:单价和数量,单价按品类分布设置不同区间。比如手机数码价格偏高、购买频次低;食品生鲜价格偏低、购买频次高
  • pay_time:付款时间,跨过去12个月,并让近期数据占比略高,模拟平台增长
  • region:用户所在省份,按人口权重分布
  • is_paid:是否支付成功,生成时留3%~5%的未支付记录,供后续清洗演示

生成量级可以控制在20万到200万行。我实际测试常用100万行,Spark在local模式跑特征提取大概一两分钟,聚类几十秒,整体体验好又不至于太假。

2.2 上传HDFS之前的目录规划与权限检查

伪分布式启动好之后,第一件事不是急着put,而是规划HDFS目录。养成好习惯,不然到后面中间结果越堆越多,自己都找不到。我习惯这样建目录:

# 在HDFS根目录下建业务目录 hdfs dfs -mkdir -p /user/hadoop/order/raw hdfs dfs -mkdir -p /user/hadoop/order/spark_output hdfs dfs -mkdir -p /user/hadoop/order/checkpoint

然后上传原始数据:

hdfs dfs -put ./data/order_data.csv /user/hadoop/order/raw/

上传后一定要验证数据真的落对了。用hdfs dfs -ls查看大小,用hdfs fsck查文件块分布,至少确认文件不是0字节、副本数符合预期。伪分布式下为了省空间,可以在hdfs-site.xml里把dfs.replication设为1,否则默认3副本会让100万行的文件白白占三倍空间。

Spark里读取这个文件的地址,在local模式下可以写:

df = spark.read \ .option("header", "true") \ .schema(custom_schema) \ .csv("hdfs://localhost:9000/user/hadoop/order/raw/order_data.csv")

这里的hostname和端口要和core-site.xml里fs.defaultFS保持一致。有个容易踩的坑:如果你用的是hdfs://master:9000,但/etc/hosts没有把master映射到127.0.0.1,Spark客户端会解析失败。最简单是直接配一个hosts映射,或者统一用localhost。

2.3 Spark读取HDFS数据时的schema设计

我见过太多人直接用spark.read.csv不指定schema,然后解析出来全是string类型,后面写聚合的时候各种cast。正确做法是第一步就定义好结构。下面是我这段代码的核心部分:

from pyspark.sql.types import StructType, StructField, IntegerType, StringType, DoubleType, TimestampType custom_schema = StructType([ StructField("user_id", IntegerType(), True), StructField("order_id", StringType(), True), StructField("item_category", StringType(), True), StructField("item_price", DoubleType(), True), StructField("quantity", IntegerType(), True), StructField("pay_time", StringType(), True), StructField("region", StringType(), True), StructField("is_paid", IntegerType(), True) ]) df = spark.read.option("header", "true").schema(custom_schema).csv(...)

pay_time先用StringType读进来,后面统一转时间。直接设成TimestampType也行,但CSV里格式稍微不对就会整列变null,所以我习惯先当字符串读,清洗阶段再统一转换。

数据清洗这一段,要覆盖三个经典问题:

  1. is_paid = 0 的记录直接过滤,模拟交易流水里的失败订单。
  2. item_price <= 0 或 quantity <= 0 的数据视为脏数据删除。
  3. order_id 有重复的,按pay_time保留最新一条。
from pyspark.sql import functions as F clean_df = df.filter("is_paid = 1") \ .filter("item_price > 0 AND quantity > 0") \ .dropDuplicates(["order_id"]) \ .withColumn("pay_time_ts", F.to_timestamp("pay_time", "yyyy-MM-dd HH:mm:ss")) \ .filter("pay_time_ts is not null")

清洗完我习惯先count一下,看过滤掉多少比例。如果脏数据比例异常大,回头检查生成脚本,问题一般在模拟数据阶段而不是代码阶段。这个判断很重要,可以节省很多调试时间。

关于编码:全链路统一UTF-8。CSV文件生成时用utf-8-sig还是utf-8要注意,前者适合Excel打开不乱码,但Spark读的时候BOM头可能混进第一列列名。我在项目里直接用utf-8生成,如果评审老师想用Excel看数据,再单独输出一份utf-8-sig副本。这个细节虽然小,但实际体验差别很大。

3. 购买力画像的Spark实现:指标设计、聚类分组与结果输出

3.1 从“订单流水”到“用户特征宽表”

很多初学者拿到清洗好的订单表,直接group by user_id算一个总金额,就当购买力了。实际上购买力是一个多维概念,最少也要从金额、频次、时间、品类四个维度去刻画。我在这个项目里构建了六个特征,做成用户级宽表:

字段名业务含义计算逻辑
total_amount累计消费金额sum(item_price * quantity)
order_cnt下单次数count(order_id)
avg_order_amount平均客单价total_amount / order_cnt
item_category_cnt购买品类数量approx_count_distinct(item_category)
recent_gap_days最近一次购买距今天数datediff(now, max(pay_time_ts))
max_single_amount最大单笔订单金额max(item_price * quantity)

代码上用一个groupBy加多个agg就能完成:

from pyspark.sql import functions as F user_features = clean_df \ .groupBy("user_id") \ .agg( F.sum("item_price * quantity").alias("total_amount"), F.count("order_id").alias("order_cnt"), (F.sum("item_price * quantity") / F.count("order_id")).alias("avg_order_amount"), F.approx_count_distinct("item_category").alias("item_category_cnt"), F.max("pay_time_ts").alias("last_pay_time"), F.max("item_price * quantity").alias("max_single_amount") ) \ .withColumn("recent_gap_days", F.datediff(F.current_date(), F.col("last_pay_time"))) \ .drop("last_pay_time")

需要注意两个细节。第一,approx_count_distinct比count distinct在数据量大时性能好很多,误差在可接受范围。第二,如果直接用sum/count,遇到除数为0会报错,好在前面已经过滤了无效订单,每个用户必然有至少一条记录。如果担心数据质量问题,可以加when保护。

3.2 为什么选择K-Means做分层,以及如何确定K

用户购买力“差异分析”不能只靠拍脑袋切阈值。比如有人拍“消费超过5000就是高购买力”,可如果整体消费水平低,5000可能是天花板,整个人群全被归为低购买力,看不出差异。这时候无监督聚类更合适。K-Means实现简单、Spark MLlib原生支持、结果可解释,是课程设计里的最佳入门选择。

选K要先做手肘法或者轮廓系数。我不会在这个项目里画特别复杂的图,但会跑一组K值对比。简单来说,k从2到6,每个k训练一次,记录计算代价,也就是各点到簇中心的平方距离和。当k增加时代价下降变缓的那个拐点,就是合适的K。实际用默认k=3也能跑,但为了把“高中低”三档说清楚,我最后固定用k=3。

聚类前必须做标准化。这个非常关键,因为total_amount动辄几千上万,order_cnt只有个位数或几十,如果直接用原始特征算欧式距离,金额会把其他特征完全淹没,聚类结果约等于按金额排序切分,那就失去多维分析的意义了。

from pyspark.ml.feature import VectorAssembler, StandardScaler from pyspark.ml.clustering import KMeans feature_cols = ["total_amount", "order_cnt", "avg_order_amount", "item_category_cnt", "recent_gap_days", "max_single_amount"] assembler = VectorAssembler(inputCols=feature_cols, outputCol="features_raw") scaled = StandardScaler(inputCol="features_raw", outputCol="features", withStd=True, withMean=True) kmeans = KMeans().setK(3).setSeed(42).setFeaturesCol("features").setPredictionCol("cluster")

注意:StandardScaler这里我开了withMean=True,等价于先中心化再缩放到单位方差。对于KMeans来说中心化不是必须的,但可以让迭代更稳,后面解读簇中心也更直观。我实际跑下来,标准化后三个簇的边界比直接用原始特征清晰得多,高购买力簇的特征均值和其他簇拉开明显差距。总的执行时间约30秒,体验很好。

3.3 聚类结果解读与落地

模型跑完之后,Pipeline给出每个用户的cluster_id。我建议先用groupBy看每个簇的人数占比,再join原始特征看簇内均值。这时候往往会出现一个很有意思的现象:金额最高的一群人,不一定是最“健康”的购买力人群。比如高金额簇可能全是“低频高客单”的用户,一年只买一两次大件,而次高簇才是高频稳定复购的用户。购买力的差异分析到这里才算真的有发现。

结果落地有两种方式。我在这套项目里同时写了两个分支:

  1. 写回HDFS:
result.write.mode("overwrite") \ .option("header", "true") \ .csv("hdfs://localhost:9000/user/hadoop/order/spark_output/user_power_result")

Spark写CSV会生成一个目录,里面有part-xxxxx文件,不是单文件。如果想合并成单文件,可以用coalesce(1)再写,或者后续在本地用glob把part文件拼起来。这个细节在答辩演示时很加分,因为你不用每次打开目录翻找part文件。

  1. 写回MySQL,Django后面直接读:
result.select("user_id", "total_amount", "order_cnt", "avg_order_amount", "item_category_cnt", "recent_gap_days", "max_single_amount", "cluster") \ .write.format("jdbc") \ .option("url", "jdbc:mysql://localhost:3306/ecommerce") \ .option("dbtable", "user_power_result") \ .option("user", "root").option("password", "****").save()

用JDBC写数据库需要提前建好表结构,字段类型要和DataFrame对齐,否则会报Data truncation。还有一个坑:MySQL驱动jar要放到Spark的jars目录里,不然报ClassNotFound。如果不想折腾驱动,单机场景下完全可以把结果输出CSV,再用Django脚本导入数据库。前期推荐后者,简单可控。

4. Django服务层设计:把分析结果变成货架上的API

4.1 结果从Spark到Django数据库的三种搬运方式

上一章末尾提到两种落地方式,这一章具体展开。Spark和Django都是Python生态,但一个是分布式计算环境,一个是Web进程,两者默认不共享内存,结果必须通过“某个中间介质”搬运。我总结有三种方式:

方式优点缺点适用场景
Spark写CSV,Django导入无额外依赖,调试直观多一步手动/脚本导入课程设计、数据量小
Spark写MySQL,Django读MySQL贴近生产,链路自动要配JDBC驱动和建表毕设展示、稍大数据量
Spark直接调Django API也能实现,前端变实时每批数据要分片请求,效率低实时增量场景

实际项目里,第一阶段用方式1,跑通全流程后切换方式2,再把驱动的配置记录下来写进文档。这样从简单到进阶,逻辑完整。

4.2 Django模型与业务口径

Django项目里,核心模型是两个:一个存用户购买力结果表,一个存面向大屏的聚合统计表。模型设计要和前面Spark的输出字段完全对应,不然接口返回的数据对不上。

# apps/user_power/models.py from django.db import models class UserPower(models.Model): user_id = models.IntegerField(unique=True) total_amount = models.DecimalField(max_digits=12, decimal_places=2) order_cnt = models.IntegerField() avg_order_amount = models.DecimalField(max_digits=10, decimal_places=2) item_category_cnt = models.IntegerField() recent_gap_days = models.IntegerField() max_single_amount = models.DecimalField(max_digits=10, decimal_places=2) cluster = models.IntegerField() class Meta: db_table = "user_power_result" class DashboardSummary(models.Model): stat_date = models.DateField(auto_now_add=True) total_users = models.IntegerField() total_amount = models.DecimalField(max_digits=16, decimal_places=2) high_power_size = models.IntegerField() middle_power_size = models.IntegerField() low_power_size = models.IntegerField() average_order_amount = models.DecimalField(max_digits=10, decimal_places=2)

我特别想提醒一个“业务口径”问题:不同统计口径得到的数值不同。比如total_amount,是包含历史所有订单还是只看最近一年?order_cnt是支付成功的订单还是包含退款?这些在Spark阶段定了什么口径,Django模型和前端大屏就必须保持一致。我建议把口径说明写进文档的第一页,答辩的时候老师问起来你直接给答案,印象分会高很多。

导入数据时,写一个management command,比如:

python manage.py import_user_power --csv-path=./data/user_power_result.csv

脚本内部用Django ORM的bulk_create批量导入,几十万行也就十几秒。不要一行一条save,慢到怀疑人生。

4.3 接口与URL设计

大屏前端只需要几个接口,多了反而乱。我最终保留这四个:

  • GET /api/overview/:返回核心指标卡数据,包括总用户数、总销售额、高/中/低购买力用户占比、平均客单价
  • GET /api/user_power/labels/:返回三个簇的人数分布和特征均值,用于饼图和雷达图
  • GET /api/top_users/?cluster=high&limit=10:返回指定簇Top N用户,用于排行榜
  • GET /api/trend/:返回按月销售额趋势,用于折线图

Django代码直接用JsonResponse,不需要引入DRF。这样依赖少,项目干净,逻辑也足够清晰。

# apps/user_power/views.py from django.http import JsonResponse from .models import UserPower, DashboardSummary def overview(request): summary = DashboardSummary.objects.order_by("-stat_date").first() return JsonResponse({ "total_users": summary.total_users, "total_amount": float(summary.total_amount), "high_power_size": summary.high_power_size, "middle_power_size": summary.middle_power_size, "low_power_size": summary.low_power_size, "average_order_amount": float(summary.average_order_amount), }) def top_users(request): cluster = request.GET.get("cluster", "high") limit = int(request.GET.get("limit", 10)) users = UserPower.objects.filter(cluster=cluster) \ .order_by("-total_amount")[:limit] data = [{ "user_id": u.user_id, "total_amount": float(u.total_amount), "order_cnt": u.order_cnt, "avg_order_amount": float(u.avg_order_amount), "recent_gap_days": u.recent_gap_days, } for u in users] return JsonResponse({"code": 0, "data": data})

这里有个小经验:JsonResponse返回的Decimal类型必须转成float,否则json.dumps会直接报错。我在实际调试中碰到过一次,后来约定所有金额字段在接口层统一转float,前端也不用单独处理字符串。

另外,如果打算用HTTP请求调试,记得在Django的settings.py里加django-cors-headers,并在MIDDLEWARE里加上CorsMiddleware,配置CORS_ALLOW_ALL_ORIGINS = True。本地开发时大屏和Django通常在不同端口,不配跨域,前端fetch会被浏览器拦得死死的。这是可视化大屏环节最常见的“前端怎么没数据”的原因之一。

5. 可视化大屏:购买力差异的最终呈现方案

5.1 大屏布局与信息层级

可视化大屏的项目属性很特殊,它不只是展示,更要回答“购买力差异到底差异在哪”。布局上我采用经典的三段式:

  • 顶部:全局关键指标。三个数字卡:总用户数、总销售额、平均客单价,外加高/中/低购买力用户占比的小圆环。
  • 中部:核心分析区。左半区放“购买力分层占比”的饼图和“三簇特征均值雷达图”,中间放“用户购买力分布热力地图”,按省份聚集总消费;右半区放“高购买力Top10用户”排行列表。
  • 底部:横向条形图,展示不同品类在高中低购买力簇中的销售额占比;再放一个折线图,展示过去12个月销售额趋势。

这样从上到下,是一个“总体规模→结构差异→地理分布→头部玩家→品类偏好”的递进逻辑。老师看大屏,第一眼看到全局,第二眼看到分析,第三眼才能看到细节。信息层级清晰,比一股脑堆十几个图表更能说明白问题。

5.2 ECharts与Django接口对接

前端我坚持用原生HTML+ECharts,不引Vue。原因是这个项目的核心是大数据链路,不是前端工程化,引入Vue还要npm构建,反而增加调试负担。ECharts的CDN引入方式足够:

<script src="https://cdn.jsdelivr.net/npm/echarts@5.4.3/dist/echarts.min.js"></script>

所有数据都通过fetch拿。以购买力分层饼图为例:

async function loadUserPowerLabels() { const resp = await fetch('/api/user_power/labels/'); const json = await resp.json(); const pieOption = { tooltip: { trigger: 'item' }, series: [{ type: 'pie', data: json.labels, // [{name: '高购买力', value: 2340}, ...] }] }; userPowerChart.setOption(pieOption); } loadUserPowerLabels(); setInterval(loadUserPowerLabels, 60000);

定时器用60秒刷新一次,这是大屏的常规节奏,避免高频请求Django接口造成压力。大屏样式用深色背景、暖色数据,视觉上更有“大屏感”。具体配色不关键,关键的是“数据对比明显”。

5.3 数据刷新的姿势:预聚合+缓存接口

这里有个容易踩的坑:大屏如果每个图表都直接请求明细数据,Django不仅响应慢,数据库压力也大。正确做法是在Django侧缓存接口结果。因为Spark分析的结果是离线批处理产出,不是实时数据,一个小时内的查询结果完全一样,没必要每次都count和sum。

我用的最简单方案:在Django视图层加一个内存缓存,或者用django.core.cache。项目体量小,直接用cache_page或缓存装饰器就够了。

from django.views.decorators.cache import cache_page @cache_page(60 * 5) # 缓存5分钟 def overview(request): ...

这样大屏每分钟刷新并发起请求,Django只需要5分钟算一次,响应从几百毫秒降到几十毫秒。如果数据更新频率高,就把缓存时间调短或者上Redis。课程设计里用本地内存缓存,完全够。

还有一件事容易被忽略:前端单位。Spark端算出的金额单位是元,前端KPI卡显示“总销售额¥1,234,567”。要在接口里统一好命名和单位,别出现一个字段是万元、一个是元,前端直接用会闹出百倍误差。我习惯在接口字段名里加_unit后缀,比如total_amount_rmb,语义清楚,联调时不互相甩锅。

6. 调试实录:六个最容易让项目翻车的点

6.1 Hadoop伪分布式:Datanode起不来的经典原因

伪分布式搭建里最常见的诡异现象是:jps能看到NameNode,但Datanode进程总是在启动后消失或显示为0个节点。我第一次遇到时一度怀疑是端口问题,最后发现根子是多次执行bin/hdfs namenode -format导致的clusterID不一致。格式化会重新生成NameNode的clusterID,但DataNode目录里的clusterID还是旧的,启动时校验失败,进程直接退出。

解决链路是固定的:

# 1. 停掉所有Hadoop进程 stop-all.sh # 2. 清理临时目录 rm -rf /tmp/hadoop-* rm -rf /usr/local/hadoop/tmp/dfs # 3. 重新格式化NameNode hdfs namenode -format # 4. 启动并检查 start-dfs.sh jps hdfs dfsadmin -report

另外,要把core-site.xml里的hadoop.tmp.dir配成一个明确的目录,不要用系统默认的/tmp,否则系统清理临时文件时又会出类似问题。这个坑我在项目文档里专门写了一页,基本每个实训小组都遇到过,照着处理就好。

6.2 Spark版本与Hadoop版本强绑定的坑

Spark发布时针对不同Hadoop版本做了预编译包,下载前必须选对。如果你下载的Spark是hadoop2.7版,却连Hadoop3.3集群,运行时大概率报NoClassDefFoundError或者ClassNotFound。

我建议的组合是Spark 3.3.2配Hadoop 3.3.4。如果遇到报错,第一反应不是去搜“为什么报错”,而是先确认版本矩阵对不对:

Spark版本预编译Hadoop版本可连的Hadoop集群版本
Spark 3.0.xhadoop2.7 / hadoop3.2Hadoop 2.7~3.2
Spark 3.2.xhadoop3.2Hadoop 3.2/3.3
Spark 3.4.xhadoop3Hadoop 3.3+

顺便说一句,除非你真的要访问HDFS上的文件,否则local模式Spark跑起来并不太依赖Hadoop集群。但本项目的卖点就是“基于Hadoop”,所以还是要让Spark从HDFS读取,不能只用本地CSV。

6.3 中文乱码:从CSV到控制台到前端一条线

中文乱码是个全链路问题,必须一次解决。常见表现:生成CSV时用了GBK,Spark读出来乱码;Django接口返回中文,前端页面乱码;终端打印日志中文乱码。

逐段排查方法:

  1. CSV统一用UTF-8,生成脚本用encoding="utf-8",不要用utf-8-sig。
  2. Spark读取时指定编码:.option("encoding", "UTF-8")。
  3. Django设置DEFAULT_CHARSET="utf-8",前端HTML加 。
  4. 终端如果还乱,检查系统locale是否为UTF-8,通常Linux下不必改。

注意:如果Spark把中文写进CSV后,用Excel打开乱码,那只是Excel对UTF-8无BOM的读取问题,不代表数据有问题。答辩前可以另出一份utf-8-sig文件给老师看,代码层面保持utf-8不动。

6.4 Django接口500:日志里藏着答案

有一次前端大屏整个没数据,Network面板里所有接口都是500。打开后端日志,发现是Decimal转换问题——接口里直接JsonResponse返回Decimal对象,json.dumps不支持,直接抛异常。这个问题前面已经说过了,统一转float就能解决。

还有一个典型场景是MySQL驱动版本不对。如果用的MySQL 8,pymysql版本太老可能报Authentication plugin 'caching_sha2_password' cannot be loaded。升级pymysql到2.0以上,或者建库时指定mysql_native_password。这类错误信息和代码逻辑无关,排查方向错了会浪费很多时间。

6.5 大屏加载慢与Spark内存设置

在local模式跑100万行数据,如果spark.executor.memory设置得太小,作业会频繁GC甚至OOM。我一般这样配置spark-submit的参数:

spark-submit \ --master local[4] \ --driver-memory 4g \ --executor-memory 4g \ analyse_user_power.py

local[4]表示用4个线程跑,充分利用笔记本多核。如果你电脑只有8G内存,driver和executor各设2g,同时别开太多浏览器标签页。还有,Spark的checkpoint目录如果设到HDFS上,任务结束后会留下临时文件,记得用--conf spark.cleaner.referenceTracking.cleanCheckpoint=true,保持目录干净。

6.6 调试方法论:链路每一环都做人工抽样

最后这条不是具体bug,而是我反复吃亏后的方法论。整个项目链路很长,一旦最终大屏数据不对,你根本不知道问题出在清洗、聚类、导入还是API。我的做法是:在每个阶段末尾做人工抽样。

  • HDFS阶段:hdfs dfs -cat文件前几行,确认原始数据没问题。
  • Spark清洗后:打印clean_df.sample(0.01).show(),人工看几条。
  • 聚类后:打印每个簇count和均值。
  • 导入Django后:在数据库里select几条,核对数量。
  • API阶段:浏览器直接访问接口,看JSON结构。

如果用这五步走一遍,发现问题基本能控制在单个环节内。如果没有抽样习惯,你会在“Spark代码对不对”和“Django代码对不对”之间反复横跳,效率极低。

做完这个项目之后我最大的体会是,购买力差异分析这种课题,真正的难点不在某个单一算法,而在于把数据、计算、服务、展示四层串成一条完整的链路。遇到问题先分段定位,别急着改代码;文档里把链路图和字段口径随手记下来,后面调试和答辩都能少很多痛苦。尤其是源码和文档的配合,我习惯在文档开头写一段“30分钟跑通指南”,把所有启动命令按顺序贴上去,这样过一个星期你再看自己的代码,也不会一脸懵。如果你正卡在某个环节,可以按我上面说的链路抽样法,一截一截排查,绝大多数问题都会浮出水面。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询