
我最早拿到这个选题的时候第一反应是“这套技术栈是不是有点重了”——二手房爬虫而已Python加个Requests加个Pandas不就能跑吗但真正把数据量铺开、把“爬取-清洗-分析-展示”整条链路跑完我才意识到HadoopSparkDjango这套组合不是炫技它解决的恰恰是单机脚本处理不了的问题海量房源数据去重、区域价格聚合计算、以及多批次增量数据的稳定存储。这篇文章就把我整个项目的落地过程写清楚包括环境搭建、爬虫设计、Spark分析任务、Django后端和可视化大屏的完整链路以及每一层踩过的坑和取舍逻辑。适合正在做大数据课设、或想用真实业务数据串起Hadoop和Spark技术点的读者。1. 为什么是HadoopSparkDjango整条链路的选型逻辑1.1 爬虫数据量级带来的第一个分岔路很多人做二手房爬虫爬到一两万条数据就开始做分析了。这个量级确实用不到Hadoop一台8G内存的笔记本跑Pandas完全没压力。但真实的二手房源信息有三个特点字段多、重复率高、时间维度长。字段多到什么程度一个小区的房源详情页包含标题、价格、单价、面积、户型、楼层、朝向、装修、建筑年代、小区名、区域、经纬度、挂牌时间、带看次数等三十多个字段。我用Scrapy跑了一个月的增量采集单城市累计数据很快到了几十万条。这时候再做全量去重、区域聚合Pandas的内存压力就很明显了——我一度在内存16G的机器上跑GroupBy直接卡死。Hadoop在这里解决的是“存得下、并得了”的问题。HDFS把数据分散到多节点哪怕是伪分布式的多目录Spark用内存计算替代Pandas的单机DataFrame操作几十万条数据的聚合任务从原来的分钟级降到秒级。这个体验落差是选型的直接动力。1.2 Django在链路里扮演的角色Django在整个项目里不负责爬虫也不负责大数据计算它只做一件事把Spark算好的结果通过REST接口稳定地喂给前端大屏。相当于大数据分析链路和用户展示层之间的“翻译官”。为什么选Django而不是Flask或FastAPI因为项目里需要后台定时任务比如每小时拉取一次Spark输出目录的最新结果、需要简单的用户登录鉴权防止大屏接口被乱刷Django的ORM和Admin后台能省掉很多重复造轮子的工作。另外后期如果想把爬虫采集的原始数据也接入后台管理Django的Migration机制管理表结构变更非常省心。1.3 模块拆分的最终形态整个项目拆成四个模块各自独立又通过文件和接口串联爬虫模块Scrapy Requests采集房源详情输出JSON和CSV到本地目录存储模块数据上传至HDFS按日期分区计算模块Spark SQL DataFrame完成清洗、去重、聚合统计展示模块Django REST Framework提供API可视化大屏通过ECharts拉取渲染这个拆分的好处是每层都能独立调试。爬虫跑挂了不影响展示Spark任务失败了大屏显示的还是上一轮结果Django服务崩了Spark照样能把计算结果写到HDFS。对于课设和毕设项目来说这种松耦合结构写文档时也好交代。2. 爬虫层的落地细节字段设计、去重策略与增量采集2.1 目标网站分析与爬虫框架选择二手房源信息采集我当时锁定的目标是某大型房产平台的公开挂牌数据。这类网站的反爬策略主要有三块请求频率限制、部分字段动态加载、User-Agent和Cookie校验。我最终选了Scrapy框架不是因为它比Requests快而是因为它的信号机制和Pipeline非常适合做“解析-清洗-存储”的分层处理。爬虫项目里我建了三个关键文件items.py定义房源字段结构pipelines.py负责去重、清洗、格式化spiders/ershoufang.py核心爬虫逻辑Items定义大概长这样import scrapy class HouseItem(scrapy.Item): title scrapy.Field() # 标题 total_price scrapy.Field() # 总价万元 unit_price scrapy.Field() # 单价元/平 area scrapy.Field() # 面积 layout scrapy.Field() # 户型如 3室1厅 floor scrapy.Field() # 楼层信息 orientation scrapy.Field() # 朝向 building_type scrapy.Field() # 建筑类型 district scrapy.Field() # 行政区 bizcircle scrapy.Field() # 商圈 community scrapy.Field() # 小区名 listing_time scrapy.Field() # 挂牌时间 agent_name scrapy.Field() # 维护人 agent_phone scrapy.Field() # 联系方式脱敏处理字段定这么细是有原因的。后面Spark做分析时“户型”“朝向”“区域”都是天然的分组键如果爬虫阶段不把这些字段单独拆出来后面的数据清洗会非常痛苦。我见过很多爬虫项目把所有信息都塞到一整个详情字符串里分析阶段再重新解析绕了一大圈。2.2 去重策略从Redis到HDFS的两级去重爬虫阶段必须做一次去重不然重复数据传到HDFS只是浪费存储。我的方案是两级去重第一级在Scrapy的Pipeline里用Redis的Set做指纹去重。指纹怎么算把“小区名户型面积总价”拼接后做MD5。这四个字段联合起来基本能确定唯一房源。Redis的SISMEMBER操作是O(1)复杂度几十万条数据也能轻松扛住。import redis import hashlib r redis.Redis(hostlocalhost, port6379, db0) DEDUP_KEY house_fingerprints def build_fingerprint(item): raw f{item[community]}|{item[layout]}|{item[area]}|{item[total_price]} return hashlib.md5(raw.encode(utf-8)).hexdigest() def is_duplicate(fp): return r.sismember(DEDUP_KEY, fp) def add_fingerprint(fp): r.sadd(DEDUP_KEY, fp)第二级去重放在Spark任务里用DataFrame的dropDuplicates指定相同的联合字段。为什么搞两道因为爬虫可能会分多轮跑而且第一轮可能没爬到完整字段就中断了Redis里存的指纹不完整。Spark这级去重兜底保证最终HDFS上的数据是干净的。2.3 Request频率控制与异常重试爬虫最忌讳的就是一口气猛爬。我当时设置的策略是每请求一个详情页之前随机sleep 1到3秒并发数控制在4每个IP的请求数超过阈值就歇一会儿。Scrapy的下载中间件里可以写个简单的频率控制器import random import time from scrapy import signals class RandomDelayMiddleware: def process_request(self, request, spider): delay random.uniform(1, 3) time.sleep(delay) return None另外要处理反爬的封禁问题。我的经验是请求头里除了User-AgentReferer也很重要很多网站会校验Referer是否合法。我在爬虫里专门维护了一个User-Agent池每次请求随机挑选同时把Referer伪装成目标网站首页。2.4 异常数据清洗的一个细节爬虫采集的数据不可能完全干净。我遇到最多的三类问题总价字段带着“万”字或者逗号比如“1,200万”单价字段为空部分房源信息未挂牌单价户型字段不统一比如“3室1厅”和“3室1厨1厅”逻辑上应该是同类这些清洗工作我放在了Scrapy的Pipeline里做初洗然后在Spark任务里做二次清洗。清洗规则写在配置文件里方便调整def clean_price(value): if not value: return None return float(value.replace(万, ).replace(,, )) def clean_unit_price(value): if not value or 元/平 not in value: return None return float(value.replace(元/平, ).replace(,, ))这里要提醒一句清洗规则一定不能只做一次。爬虫数据源变了或者网站页面结构调整了清洗逻辑很容易失效。所以我给每个清洗函数都加了日志输出清洗前后数据条数能对得上异常数据能定位到原始记录。3. Hadoop伪分布式搭建与Spark分析任务的实战拆解3.1 伪分布式环境搭建一台机器怎么跑出“大数据感”课设和毕设场景下大部分人的处理器只有一台电脑。Hadoop伪分布式模式是唯一的现实选择——它把NameNode、DataNode、ResourceManager、NodeManager这些角色都跑在单机上虽然只有一个节点但完整保留了HDFS的存储机制和YARN的资源调度逻辑。我用的版本是Hadoop 3.3.4JDK 1.8。这里有个坑Hadoop 3.x对JDK版本要求比2.x严格JDK 11也能跑但有些老代码会报JAXB相关错误建议直接用JDK 8省事。配置文件需要改五个core-site.xml设置NameNode地址hdfs-site.xml设置HDFS副本数和NameNode目录yarn-site.xml配置YARN资源管理mapred-site.xml配置MapReduce运行模式workers或slaves设置从节点我当时最关键的配置如下!-- core-site.xml -- configuration property namefs.defaultFS/name valuehdfs://localhost:9000/value /property /configuration!-- hdfs-site.xml -- configuration property namedfs.replication/name value1/value /property property namedfs.namenode.name.dir/name value/home/user/hadoop_data/namenode/value /property property namedfs.datanode.data.dir/name value/home/user/hadoop_data/datanode/value /property /configuration启动顺序也很讲究必须先格式化NameNode第一次启动时再依次启动HDFS和YARNhdfs namenode -format start-dfs.sh start-yarn.sh jpsjps这个命令会列出Java进程正常能看到NameNode、DataNode、SecondaryNameNode、ResourceManager、NodeManager五个进程。少任何一个都要回去看日志日志在$HADOOP_HOME/logs/目录下。3.2 数据的HDFS目录规划与上传流程HDFS上不是随便扔数据的。我的目录规划是按日期分层/house_data/ raw/2025/01/06/ part-00001.json clean/2025/01/06/ part-00001.parquet analysis/ price_by_district/ layout_distribution/为什么要分raw和clean两层因为原始爬虫数据是半结构化的JSONSpark处理时可能会发现之前没注意到的脏数据如果直接覆盖原始数据就再也找不回原始信息了。我都是先存raw清洗后再写clean层这样能追溯每一步的数据变化。上传命令很简单hdfs dfs -mkdir -p /house_data/raw/2025/01/06 hdfs dfs -put ./output/ershoufang_20250106.json /house_data/raw/2025/01/06/3.3 Spark任务的核心逻辑读JSON、清洗、聚合Spark程序我用的Spark SQL的DataFrame API。一方面是因为它比RDD API更接近SQL思维写起来直观另一方面是DataFrame底层有Catalyst优化器做谓词下推和列裁剪比手动RDD转换高效很多。完整的Spark分析任务我用的是PySpark整个分析流程分五步第一步读取HDFS上的JSON数据from pyspark.sql import SparkSession spark SparkSession.builder \ .appName(HousePriceAnalysis) \ .master(yarn) \ .config(spark.sql.shuffle.partitions, 10) \ .getOrCreate() df spark.read.json(hdfs://localhost:9000/house_data/raw/2025/01/06/*.json)注意这里master(yarn)是指跑在YARN集群上而不是本地模式。如果只想调试可以临时改成local[*]快速验证逻辑。但要记住提交到YARN时PySpark的环境变量和依赖包路径都要提前配好否则容易报Python环境找不到的错误。第二步清洗数据。比如过滤掉空值、格式修正from pyspark.sql.functions import col, when df_clean df.filter(col(total_price).isNotNull()) \ .filter(col(area) 0) \ .withColumn(unit_price_final, when(col(unit_price).isNull(), (col(total_price) * 10000 / col(area)) ).otherwise(col(unit_price)))这里有个逻辑单价字段为空时可以用总价乘以10000万元转元除以面积估算出单价。这个兜底逻辑在真实业务里很常见。第三步区域价格聚合。计算每个行政区的平均单价、平均总价、房源数量price_by_district df_clean.groupBy(district) \ .agg( avg(unit_price_final).alias(avg_unit_price), avg(total_price).alias(avg_total_price), count(*).alias(house_count) ) \ .orderBy(col(house_count).desc())第四步户型分布统计layout_distribution df_clean.groupBy(layout) \ .count() \ .orderBy(col(count).desc())第五步将聚合结果写回HDFS的analysis目录price_by_district.write.mode(overwrite) \ .parquet(hdfs://localhost:9000/house_data/analysis/price_by_district) layout_distribution.write.mode(overwrite) \ .parquet(hdfs://localhost:9000/house_data/analysis/layout_distribution)3.4 Spark任务提交的两种方式PySpark任务有两种提交方式我在这上面纠结了一阵子。直接python3 xxx.py跑走的是local模式不适合生产环境用spark-submit提交能指定集群资源和依赖spark-submit \ --master yarn \ --deploy-mode cluster \ --num-executors 2 \ --executor-memory 2G \ --driver-memory 1G \ analysis.py我在实际调试时遇到过一个非常隐蔽的问题spark-submit的--deploy-mode cluster模式下Python脚本内的相对路径引用会失效因为Driver跑在不同节点上。所以代码里所有涉及路径的地方都必须写HDFS的绝对路径。这个坑我花了半天才排查出来具体细节放在后面的调试章节细说。4. 可视化大屏背后的数据通道Django如何优雅对接Spark结果4.1 Django项目的目录设计Django在项目里的职责是API服务不是渲染大屏HTML——当然它也完全可以渲染只是我的整个前端是独立的静态页面用ECharts画图Django只负责把数据通过JSON接口喂给前端。项目结构大概是django_backend/ manage.py house_api/ settings.py urls.py wsgi.py dashboard/ migrations/ models.py views.py serializers.py urls.py static/ bigscreen/ index.html echarts.min.js新建Django项目和App后我建议第一时间把CORS配置好。因为我的静态大屏页面是用Live Server直接打开的跟Django后端的端口8000不一样跨域是必然的。我用的是django-cors-headers这个库# settings.py INSTALLED_APPS [ ... corsheaders, rest_framework, ] MIDDLEWARE [ corsheaders.middleware.CorsMiddleware, ... ] CORS_ALLOW_ALL_ORIGINS True # 开发阶段用上线后改成白名单4.2 数据表模型用Django ORM管理分析结果虽然Spark把聚合结果算好了但Django不能直接读HDFS上的Parquet文件。最稳妥的方式是Spark任务跑完后把聚合结果导出成最通用的JSON或PostgreSQL能直接读取的格式然后Django通过ORM管理这些数据。我用的是SQLite加Django ORM表结构就三张from django.db import models class PriceByDistrict(models.Model): district models.CharField(max_length50, uniqueTrue) avg_unit_price models.FloatField() avg_total_price models.FloatField() house_count models.IntegerField() updated_at models.DateTimeField(auto_nowTrue) class LayoutDistribution(models.Model): layout models.CharField(max_length50) count models.IntegerField() updated_at models.DateTimeField(auto_nowTrue) class HouseTrend(models.Model): date models.DateField() avg_price models.FloatField() listing_count models.IntegerField()Spark任务每次跑完后生成一份CSV然后通过Django的管理命令导入这些CSV到数据库。管理命令我用manage.py import_data这样的自定义Command实现# dashboard/management/commands/import_data.py from django.core.management.base import BaseCommand import csv from dashboard.models import PriceByDistrict class Command(BaseCommand): help Import Spark analysis results into database def add_arguments(self, parser): parser.add_argument(--file, requiredTrue) def handle(self, *args, **options): with open(options[file], r) as f: reader csv.DictReader(f) for row in reader: obj, created PriceByDistrict.objects.update_or_create( districtrow[district], defaults{ avg_unit_price: row[avg_unit_price], avg_total_price: row[avg_total_price], house_count: row[house_count], } ) self.stdout.write(self.style.SUCCESS(Import finished))为什么用update_or_create而不是直接create因为Spark任务可能重复跑重复导入会导致数据翻倍。update_or_create按district这个唯一键匹配有则更新、无则新增完美解决幂等问题。4.3 REST API与前端大屏的数据对接API逻辑写在views.py里用DRF的APIView或者直接JsonResponse都行。我为了保证接口返回格式统一用了DRF的ModelSerializer# serializers.py from rest_framework import serializers from dashboard.models import PriceByDistrict class PriceByDistrictSerializer(serializers.ModelSerializer): class Meta: model PriceByDistrict fields [district, avg_unit_price, avg_total_price, house_count]# views.py from rest_framework.response import Response from rest_framework.views import APIView from dashboard.models import PriceByDistrict from dashboard.serializers import PriceByDistrictSerializer class DistrictPriceAPI(APIView): def get(self, request): queryset PriceByDistrict.objects.order_by(-house_count)[:20] serializer PriceByDistrictSerializer(queryset, manyTrue) return Response({ code: 0, data: serializer.data, msg: success })前端大屏页面我用的是ECharts数据拉到之后直接setOption换数据// bigscreen/index.html 中的关键片段 async function loadDistrictData() { const res await fetch(http://127.0.0.1:8000/api/district-price/); const json await res.json(); const data json.data; myChart.setOption({ series: [{ type: bar, data: data.map(item item.avg_unit_price) }], xAxis: { data: data.map(item item.district) } }); }大屏上我做了五个核心图表区域均价柱状图、总价分布饼图、户型占比环形图、挂牌数量趋势折线图、单价Top10排行榜。数据源分别是Spark算好的五张聚合表。4.4 定期刷新机制让大屏数据“活”起来大屏不能只显示一次数据它需要周期性地刷新。我的做法是在Django里写一个定时任务每小时重新从Spark输出目录拉取最新CSV并更新数据库。这里用的是django-crontab或者APScheduler我最后选了逻辑更清晰的APScheduler# dashboard/jobs.py from apscheduler.schedulers.background import BackgroundScheduler from django.core.management import call_command def refresh_data_job(): # 先拉取最新的Spark聚合结果到本地临时目录 # 然后通过管理命令导入数据库 call_command(import_data, --file, /tmp/latest_price_by_district.csv) scheduler BackgroundScheduler() scheduler.add_job(refresh_data_job, interval, hours1) scheduler.start()有一点必须注意APScheduler在Django的debug模式下会重复启动两次因为runserver会reload。解决方案是在settings.py里加一个判断或者用--noreload参数启动。5. 调试阶段的排雷经历伪分布式搭建与Spark内存问题复盘5.1 Hadoop伪分布式搭建最常见的三个报错整个项目耗时最多的不是写代码而是调环境。我把Hadoop伪分布式搭建过程中遇到的三个典型报错完整复盘一下。第一个报错启动HDFS时DataNode起不来日志里报Incompatible clusterIDs。这个问题的根源是格式化NameNode后DataNode目录里残留了上次运行的clusterID两边的ID对不上就拒绝启动。解决方法是把DataNode和NameNode的数据目录全部清空重新格式化rm -rf /home/user/hadoop_data/namenode/* rm -rf /home/user/hadoop_data/datanode/* hdfs namenode -format start-dfs.sh第二个报错localhost:9000: Call From host/192.168.x.x to localhost:9000 failed on connection exception。原因有迹可循多半是core-site.xml里用了IP地址而机器的hostname解析没有对应关系。我的解决方式是把fs.defaultFS从IP改成hdfs://localhost:9000同时确保/etc/hosts里有127.0.0.1 localhost这一行。第三个报错YARN的ResourceManager起了但任务提交时一直停留在ACCEPTED状态不往下走。这个问题十有八九是内存配置不够。默认的yarn-site.xml里分配的内存比较大一个NodeManager需要的内存超过了机器实际可用内存。解决方法是在yarn-site.xml显式调小property nameyarn.nodemanager.resource.memory-mb/name value4096/value /property property nameyarn.scheduler.maximum-allocation-mb/name value4096/value /property property nameyarn.scheduler.minimum-allocation-mb/name value512/value /property5.2 Spark任务在YARN上运行失败Python in worker has different version这个报错我印象太深了。本地跑PySpark一切正常一旦用spark-submit --master yarn提交就报这个有时候还会附带着org.apache.spark.SparkException: Python worker failed to connect back。排查过程是这样的先看报错的完整堆栈确认是Executors上的Python环境与Driver不一致。原因在于我的服务器上装了多个Python版本——系统自带的Python 3.6、我手动装的Python 3.10。YARN的NodeManager节点上Spark默认找的python命令对应的是老版本。解决方案有两种。第一种是临时指定Python路径spark-submit \ --master yarn \ --conf spark.pyspark.python/usr/bin/python3.10 \ --conf spark.pyspark.driver.python/usr/bin/python3.10 \ analysis.py第二种更彻底在所有节点上统一Python版本并把路径写入/etc/profile.d/spark_env.sh确保Spark启动时能找到正确的PYSPARK_PYTHON环境变量。这个方案需要在每个worker节点都配置如果只有一个节点的伪分布式模式第一种方法就够了。5.3 shuffle分区过大导致的内存溢出还有一次Spark任务执行时报的是java.lang.OutOfMemoryError: Java heap space发生在GroupBy聚合阶段。问题是shuffle产生的分区数太多了。Spark默认的spark.sql.shuffle.partitions是200对于几十万条数据来说完全没必要。调整方法很简单spark SparkSession.builder \ .config(spark.sql.shuffle.partitions, 10) \ .getOrCreate()从200调到10之后内存压力小了很多任务也跑完了。这个参数不是越小越好要看数据量和集群资源一般经验值是分区数大概等于核心数的2到3倍。伪分布式单机只有几个核10到20就够用。5.4 Django端调试时的一个隐蔽BugAPScheduler重复注册Django的debug模式会自动reloadAPScheduler如果写在App的ready方法里会被执行两次导致定时任务重复注册。我调试时发现import_data任务一个小时执行了两次数据库里的updated_at时间差了不到一秒。解决方案是在启动前检查是否已有实例from datetime import datetime from apscheduler.schedulers.background import BackgroundScheduler _scheduler None def get_scheduler(): global _scheduler if _scheduler is None: _scheduler BackgroundScheduler() return _scheduler或者干脆在settings里根据DEBUG决定是否挂载定时任务调度器if not DEBUG: scheduler get_scheduler() scheduler.add_job(...) scheduler.start()这个处理方式其实是个双保险——开发时定时任务不生效但数据刷新由job的运行逻辑兜底。6. 从单机脚本到集群部署的演进这套架构还能怎么扩展6.1 数据量再翻一倍时的升级路径我的项目在伪分布式的单机上跑几十万条数据是没问题的。但如果数据量再往上走比如要采集多个城市、甚至全国的二手房源整套架构有几个地方需要升级。HDFS层面从伪分布式到真集群只需要把workers文件里加上多台机器的hostname并确保各节点SSH免密登录配好。代码逻辑完全不用改Spark的hdfs://localhost:9000改成hdfs://namenode_host:9000即可。Spark层面资源不够时可以增加executor数量和内存spark-submit \ --master yarn \ --num-executors 4 \ --executor-memory 4G \ --driver-memory 2G \ analysis.py爬虫层面单机Scrapy的并发是有限的如果数据量成倍增长可以考虑用Scrapy-Redis做分布式爬虫让多个爬虫节点共享请求队列、共享去重指纹。Redis的去重逻辑我已经在项目里预留了接口升级时只需要把Dedup Pipeline改为从Redis读取统一指纹集合。6.2 计算结果的时效性优化目前的数据刷新链路是“Sparking算完 → 导出CSV → Django定时导入”。一小时刷新一次对于这种房源分析大屏来说够了。但如果将来业务需要分钟级实时性更好的方案是引入消息队列Spark Streaming或Structured Streaming监听Kafka里的房源变更事件计算完直接写入MySQL或ClickHouseDjango只负责读库展示。这块我虽然没在原始项目里完全落地但架构上预留了接口——Django的模型层完全不感知数据从哪里来只负责读。所以将来把“定时导入CSV”换成“直连MySQL”时前端代码一行都不用改。6.3 关于“课设/毕设项目深度”的一点建议最后聊聊项目深度的呈现逻辑。很多同学做大屏项目最容易被追问的一个问题是“你用了Hadoop但你的数据量根本不需要Hadoop为什么要用”我建议在项目文档里明确写清楚不是“用Hadoop”本身有价值而是“在大数据量场景下传统单机方案遇到了什么瓶颈HadoopSpark如何解决”。哪怕演示时用的数据量只有十万条也要把“如果要处理百万条以上我的架构将如何横向扩展”这个思路写明白。这种设计思路比单纯堆技术栈更能体现工程能力。我在最终的项目文档里放了一张“数据量-技术选型”对照表把千条、万条、十万条、百万条不同量级的技术方案差异讲清楚导师和评委看到这种思考深度通常不会抓着“你是不是大材小用”不放。6.4 一个小技巧用Airflow统一调度整条链路如果你想把这个项目做得再完整一点可以引入Airflow把爬虫、Spark任务、Django数据导入三步串成一条DAG。每一步都做成一个PythonOperator上游失败时自动跳过下游。我当时因为时间原因只在文档里规划了这个方案但现在已经把Airflow的DAG文件写好了排错效率比手动跑三个命令高非常多——尤其是Spark任务失败后能自动重试不用半夜爬起来看日志。Airflow的DAG定义代码也很简短from airflow import DAG from airflow.operators.bash_operator import BashOperator from datetime import datetime, timedelta default_args { owner: me, depends_on_past: False, start_date: datetime(2025, 1, 1), retries: 2, retry_delay: timedelta(minutes5) } dag DAG(house_analysis_pipeline, default_argsdefault_args, schedule_intervaldaily) task1 BashOperator(task_idrun_spider, bash_commandscrapy crawl ershoufang, dagdag) task2 BashOperator(task_idrun_spark, bash_commandspark-submit --master yarn analysis.py, dagdag) task3 BashOperator(task_idimport_django, bash_commandpython manage.py import_data --file /tmp/latest_data.csv, dagdag) task1 task2 task37. 写在最后复盘整个项目的得与失整个项目从环境搭建到可视化大屏完整跑通前后花了我大概三周时间。现在回头看最耗时间的反而不是写代码而是Hadoop伪分布式环境的一次次格式化重启、Spark调参的一次次失败重试。但这些坑恰恰是项目最值钱的部分——如果所有东西一次就通了那这个项目学到的只是“调API”的皮毛毛。我个人最大的体会是把“爬虫-大数据-后端可视化”整条链路串起来比单独精通其中某一环收获大得多。你会在调试中被迫理解HDFS的存储原理被迫搞懂Spark Driver和Executor之间的通信机制被迫想清楚Django ORM的性能边界。这些知识单看教程是一回事真踩过坑是另一回事。如果你正准备做类似的项目我的建议是先把整条链路的最小版本跑通——哪怕先用十条规定数据、本地模式跑Spark、Django返回静态JSON——然后再逐层替换成Hadoop、YARN模式、Spark真正读HDFS。一次性把大象装进冰箱只会让你在环境问题里耗尽信心。分层递进的节奏才是这类项目最稳的推进方式。最后再提醒一句HDFS上的数据目录权限容易被人忽略如果上传文件时遇到Permission denied直接看一下hdfs dfs -chmod -R 777 /house_data这个命令能帮你省下至少二十分钟的排查时间。