ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

Spark本地模式实战:Django+Scrapy+Redis构建实时电商分析系统

2026/10/3 4:41:00 拓冰建站 浏览量
Spark本地模式实战:Django+Scrapy+Redis构建实时电商分析系统 简介这是一套面向Python初学者与毕业设计学生的全栈大数据可视化项目资源聚焦电子产品信息采集、分析与Web展示全流程。项目融合Django后端、Vue前端、Spark分布式计算及自研爬虫技术解决多源电商数据获取、清洗、存储与交互式图表呈现的实际问题适用于毕设开发、课程设计或工程实训。压缩包含507个文件以70个Vue组件构建响应式界面、47个Python脚本实现爬虫与Spark任务调度、42个JPG/PNG素材支撑UI展示、41个JS逻辑文件增强交互辅以SQL建表语句、BAT一键部署脚本及.bak备份文件便于调试溯源整体大小20.13MB。已有87人学习下载资源提供可直接运行的完整源码、MySQL 5.7兼容数据库结构、清晰的安装/运行指引及典型页面如IndexMain.vue、update-password.vue的模块化组织助学习者快速理解前后端协同机制与大数据处理链路。1. 这不是又一个“爬虫Django图表”的缝合怪它用 Spark 做了三件别人不敢动的事你点开这个毕设压缩包看到5p123基于Spark的电子产品信息查询可视化系统0_djangospider.zip第一反应可能是“哦又是学生拿 Scrapy 爬京东/淘宝存进 MySQLDjango 后台查ECharts 画个柱状图”——但这次真不是。它把Spark 作为核心数据引擎而不是“跑个 WordCount 演示就下线”的装饰品。具体来说爬虫spider产出的原始 HTML 不落地存数据库而是直接喂给 Spark Streaming微批模式做实时清洗、字段归一比如把“iPhone 15 Pro Max 256GB 钛金属”、“苹果 iPhone15ProMax 256G 钛”、“15PM 256G 钛”统一成标准 SKUDjango 不是只调Model.objects.filter()而是通过pyspark.sql.SparkSession直连 Spark Thrift Server把用户在网页输入的“价格区间品牌关键词”翻译成动态生成的 Spark SQL绕过 ORM 层直查内存计算结果可视化大屏echarts的数据源不是 Django View 返回的 JSON而是 Spark 计算后写入 Redis 的 Hash 结构hset product_stats:20240615 top_brand [{name:Apple,value:1287}]前端用 Ajax 轮询 Redis Key毫秒级响应。适合谁正卡在毕设选题里、怕被导师问“Spark 用在哪了”的本科生想把爬虫项目从“能跑”升级到“能扛量”的 Python 工程师实测单机 Spark 4核8G 内存可稳定处理 3000 页面/分钟的增量解析对“Django 怎么和 Spark 打通”有执念、搜遍 CSDN 只看到“用 Pandas 加载 CSV”的人。别急着解压 zip——先搞懂这三件事为什么非得用 Spark否则你跑通 demo 后一加真实数据就内存 OOM、SQL 报错、Redis 键名冲突血泪经验。2. 从零搭起 Spark 计算链路不装 Hadoop不配 YARN单机也能跑通全流程2.1 为什么 Spark 不走 HDFS而用本地文件系统 Redis 做中间态很多毕设失败第一步就栽在环境上学生装 Hadoop 失败、配 YARN 权限报错、集群起不来……其实本项目完全规避了这些。它的数据流是spider → 本地 /tmp/spark_raw/ (Parquet 分区) ↓ Spark Structured Streaming (micro-batch, 30s interval) ↓ 清洗后 Parquet → /tmp/spark_cleaned/ ↓ Spark SQL 查询 → 结果写入 Redis (Hash Expiry) ↓ Django View 读 Redis → 返回 JSON 给 ECharts关键点所有 Spark 作业都运行在 local[*] 模式不依赖外部集群。你只需要一台能跑 Python 的电脑装好 Java 8 和 Spark 3.3 即可。HDFS 不是必须项反而是累赘——毕设场景数据量在 GB 级本地 SSD 读写 Parquet 比 HDFS 快 3 倍且避免了 NameNode 单点故障的玄学问题。2.2 安装 Spark 并验证本地模式可用Windows/macOS/Linux 通用提示不要用官网下载的带 Hadoop 的 Spark 包spark-3.3.2-bin-hadoop3.tgz它自带的 Hadoop DLL 在 Windows 上极易报UnsatisfiedLinkError。直接下spark-3.3.2-bin-without-hadoop.tgz再手动配 Hadoop 依赖实际本项目根本不用。# 1. 下载无 Hadoop 版 Spark以 macOS 为例Linux 同理Windows 用 .zip curl -O https://downloads.apache.org/spark/spark-3.3.2/spark-3.3.2-bin-without-hadoop.tgz tar -xzf spark-3.3.2-bin-without-hadoop.tgz export SPARK_HOME$PWD/spark-3.3.2-bin-without-hadoop export PATH$SPARK_HOME/bin:$PATH # 2. 配置 Spark 使用本地文件系统关键删掉所有 hdfs:// 路径 # 编辑 $SPARK_HOME/conf/spark-defaults.conf添加 spark.sql.warehouse.dir file:///tmp/spark_warehouse spark.sql.adaptive.enabled true spark.sql.adaptive.coalescePartitions.enabled true # 3. 验证启动 pyspark跑一个最小清洗任务 pyspark --master local[*] --driver-memory 4g --executor-memory 4g# 在 pyspark shell 中执行验证 Spark 本地模式是否就绪 from pyspark.sql import SparkSession from pyspark.sql.functions import col, regexp_replace, lower spark SparkSession.builder \ .appName(test-local) \ .master(local[*]) \ .config(spark.sql.adaptive.enabled, true) \ .getOrCreate() # 创建测试数据模拟爬虫抓到的脏数据 data [ (iPhone 15 Pro Max 256GB 钛金属, Apple, 8999, https://jd.com/123), (苹果 iPhone15ProMax 256G 钛, Apple, 8999, https://tmall.com/456), (15PM 256G 钛, Apple, 8999, https://suning.com/789), (Mate60 Pro 512GB 雅川青, Huawei, 6999, https://huawei.com/abc) ] df spark.createDataFrame(data, [title, brand, price, url]) # 清洗统一标题格式去空格、转小写、正则归一 clean_df df.withColumn(clean_title, lower(regexp_replace(col(title), r[^\w\u4e00-\u9fa5], )) # 去标点、空格保留中文和字母数字 ).withColumn(price_num, col(price).cast(double)) clean_df.show(truncateFalse) # 输出应为四行clean_title 列显示 iphone15promax256gb钛金属 等归一化结果逻辑说明这段代码验证了 Spark Structured Streaming 的核心能力——列式计算、UDF 免编译、自动类型推断。regexp_replace是清洗电子商品标题的关键因为电商页面标题千奇百怪缩写、符号、中英文混排靠字符串replace()根本覆盖不全。cast(double)强制转价格为数值型为后续filter(price_num 5000)做准备。参数说明--driver-memory 4g是底线低于 3g 会因元数据缓存不足报OutOfMemoryError: Metaspace--executor-memory 4g保证每个 core 有足够内存处理 Parquet blockspark.sql.adaptive.enabled开启自适应查询优化在数据倾斜时自动调整 shuffle partitions毕设数据量波动大不开它容易卡死。2.3 把爬虫输出接入 Spark不存 MySQL直接写 Parquet 分区本项目 spider基于 Scrapy不走 pipeline 存数据库而是重写FilePipeline将每页 HTML 解析后的结构化数据JSON Line 格式直接写入本地 Parquet 分区目录。这是 Spark 流式处理的前提——Parquet 是 Spark 最原生、最高效的文件格式支持谓词下推Predicate Pushdown、列裁剪Column Pruning比 CSV 快 10 倍以上。# scrapy_project/pipelines.py import os from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, DoubleType, TimestampType from datetime import datetime class SparkParquetPipeline: def __init__(self): # 初始化 SparkSession注意这里不能用 local[*]因为 Scrapy 是多进程会冲突 # 改用 local[1]且复用同一个 Session 实例 self.spark SparkSession.builder \ .appName(scrapy-to-parquet) \ .master(local[1]) \ .config(spark.sql.adaptive.enabled, true) \ .getOrCreate() # 定义 schema严格定义避免 Spark 自动推断出错 self.schema StructType([ StructField(product_id, StringType(), True), StructField(title, StringType(), True), StructField(brand, StringType(), True), StructField(price, DoubleType(), True), StructField(spec, StringType(), True), # 如 屏幕:6.7英寸,处理器:A17 StructField(crawl_time, TimestampType(), True), StructField(source_url, StringType(), True) ]) def process_item(self, item, spider): # 将 Scrapy Item 转为 dict再转 Spark Row row_dict { product_id: item.get(product_id, ), title: item.get(title, ), brand: item.get(brand, ), price: float(item.get(price, 0)) if item.get(price) else 0.0, spec: item.get(spec, ), crawl_time: datetime.now(), source_url: item.get(source_url, ) } # 单行转 DataFrame注意小数据量用这种方式大数据量改用 RDD.saveAsTable df self.spark.createDataFrame([row_dict], schemaself.schema) # 写入 Parquet 分区按日期分区便于 Spark Streaming 增量读取 date_part datetime.now().strftime(%Y%m%d) output_path f/tmp/spark_raw/date{date_part} df.write.mode(append).partitionBy(date).parquet(output_path) return item逻辑说明这段 pipeline 的核心是绕过数据库直写 Parquet。mode(append)保证每天新增数据追加到对应日期分区partitionBy(date)让 Spark Streaming 能用spark.readStream.format(parquet).option(maxFilesPerTrigger, 10).load(/tmp/spark_raw/)增量读取无需轮询整个目录。参数说明StructType显式定义 schema 是必须的——如果让 Spark 自动 infer遇到某天某条数据price是暂无报价字符串就会把整列推断为StringType后续filter(price 5000)直接报错。local[1]避免 Scrapy 多进程启动多个 SparkContext 导致端口冲突常见翻车点。3. Django 如何安全、高效地调用 Spark SQL拒绝硬编码 JDBC URL3.1 为什么不用 PySpark Driver 直连Thrift Server 是唯一可行方案你可能想在 Django View 里from pyspark.sql import SparkSession; spark SparkSession.builder...——绝对不行。原因有三Django 是多线程/多进程模型uWSGI/Gunicorn每个请求都新建 SparkSession 会导致 JVM 实例爆炸内存瞬间打满SparkSession 启动耗时 2~5 秒用户点一次查询等 3 秒体验极差无法复用 Spark 的查询计划缓存Query Plan Cache每次都是冷启动。正确做法启动 Spark Thrift Server 作为独立服务Django 用标准 JDBC 连接它。Thrift Server 是 Spark 官方提供的 HiveServer2 兼容服务本质是一个长期运行的 Spark SQL 服务端支持并发查询、连接池、权限控制本项目简化为无认证。# 启动 Thrift Server后台运行日志输出到 /tmp/thrift.log $SPARK_HOME/sbin/start-thriftserver.sh \ --master local[*] \ --driver-memory 4g \ --executor-memory 4g \ --conf spark.sql.hive.thriftServer.singleSessiontrue \ --conf spark.sql.adaptive.enabledtrue \ --conf spark.sql.adaptive.coalescePartitions.enabledtrue \ --hiveconf hive.server2.thrift.port10000 \ --hiveconf hive.server2.thrift.bind.hostlocalhost \ --verbose /tmp/thrift.log 21 提示--conf spark.sql.hive.thriftServer.singleSessiontrue是关键参数它让所有 JDBC 连接共享同一个 SparkSession避免资源重复申请。--hiveconf hive.server2.thrift.port10000是默认端口Django 连接时用此端口。3.2 Django 中封装 Spark SQL 查询用 PyHive 连接池防超时、防阻塞PyHive 是 Python 连接 HiveServer2即 Thrift Server的事实标准库。但直接connect()会阻塞必须加连接池和超时控制。# django_project/spark_utils.py import threading from pyhive import hive from contextlib import contextmanager from django.conf import settings import logging logger logging.getLogger(__name__) # 全局连接池线程安全 _spark_conn_pool [] _pool_lock threading.Lock() _MAX_POOL_SIZE 5 def _get_spark_connection(): 从连接池获取或新建一个 Hive 连接 with _pool_lock: if _spark_conn_pool: return _spark_conn_pool.pop() try: conn hive.Connection( hostlocalhost, port10000, usernamehive, # Thrift Server 默认用户名无密码 databasedefault, authNOSASL # 无认证模式 ) logger.info(New Spark Thrift connection created) return conn except Exception as e: logger.error(fFailed to create Spark Thrift connection: {e}) raise def _return_spark_connection(conn): 归还连接到池 with _pool_lock: if len(_spark_conn_pool) _MAX_POOL_SIZE: _spark_conn_pool.append(conn) else: try: conn.close() except: pass contextmanager def spark_sql_connection(): 上下文管理器确保连接自动归还 conn None try: conn _get_spark_connection() yield conn except Exception as e: if conn: try: conn.close() except: pass raise e finally: if conn: _return_spark_connection(conn) # 封装查询函数带超时和重试 def execute_spark_sql(query: str, timeout: int 30) - list: 执行 Spark SQL 查询返回结果列表每行是 tuple :param query: SQL 语句如 SELECT brand, COUNT(*) FROM products WHERE price 5000 GROUP BY brand :param timeout: 查询超时秒数 :return: [(brand1, count1), (brand2, count2), ...] from time import time start time() with spark_sql_connection() as conn: cursor conn.cursor() cursor.execute(query) # 设置超时PyHive 不支持 query timeout所以用时间戳判断 while time() - start timeout: try: result cursor.fetchall() return result except Exception as e: if Operation not finished in str(e): continue # 未完成继续等待 else: raise e raise TimeoutError(fSpark SQL query timeout after {timeout}s: {query})逻辑说明这个封装解决了三个痛点连接复用通过_spark_conn_pool池化连接避免频繁创建销毁超时控制execute_spark_sql函数内用time()手动计时防止某个慢查询拖垮整个 Django异常兜底contextmanager确保无论成功失败连接都会归还或关闭杜绝连接泄漏。参数说明authNOSASL是 Thrift Server 无认证模式的固定值databasedefault是 Spark 默认数据库timeout30是经验值Spark 在本地模式下GB 级数据聚合通常 3~8 秒完成30 秒足够覆盖 99% 场景。3.3 在 Django View 中动态生成 Spark SQL防注入、保性能用户在网页输入“华为 价格 3000-6000”后端不能拼接字符串fSELECT * FROM products WHERE brand{brand} AND price BETWEEN {min_p} AND {max_p}——这是 SQL 注入温床。必须用参数化查询。# django_project/views.py from django.http import JsonResponse from django.views.decorators.csrf import csrf_exempt from django.views.decorators.http import require_http_methods from .spark_utils import execute_spark_sql import json import re csrf_exempt require_http_methods([POST]) def search_products(request): 接收前端搜索条件生成并执行 Spark SQL POST body: {brand: Apple, min_price: 5000, max_price: 10000, keyword: Pro} try: data json.loads(request.body) brand data.get(brand, ).strip() min_price float(data.get(min_price, 0)) max_price float(data.get(max_price, 100000)) keyword data.get(keyword, ).strip() # 白名单校验品牌防注入 allowed_brands [Apple, Huawei, Xiaomi, OPPO, vivo, Samsung] if brand and brand not in allowed_brands: return JsonResponse({error: Invalid brand}, status400) # 构建安全 SQL用 ? 占位符PyHive 支持 base_query SELECT product_id, title, brand, price, spec, source_url FROM products WHERE 11 params [] if brand: base_query AND brand ? params.append(brand) if min_price 0: base_query AND price ? params.append(min_price) if max_price 100000: base_query AND price ? params.append(max_price) if keyword: # 关键词匹配 clean_title清洗后的标准标题 base_query AND clean_title LIKE ? params.append(f%{re.sub(r[^\w\u4e00-\u9fa5], , keyword)}%) # 过滤 keyword 中的标点 base_query ORDER BY price DESC LIMIT 100 # 执行查询PyHive execute 支持参数化 result execute_spark_sql(base_query, timeout45) # 转为字典列表前端易用 columns [product_id, title, brand, price, spec, source_url] result_list [dict(zip(columns, row)) for row in result] return JsonResponse({data: result_list}) except ValueError as e: return JsonResponse({error: Invalid number format}, status400) except TimeoutError as e: return JsonResponse({error: Query timeout, please try again}, status408) except Exception as e: logger.error(fSearch error: {e}) return JsonResponse({error: Internal server error}, status500)逻辑说明这段 View 的安全性体现在三点品牌白名单allowed_brands硬编码杜绝brandApple; DROP TABLE products; --类注入价格强转 float防止传入abc导致float()报错keyword 过滤标点re.sub(r[^\w\u4e00-\u9fa5], , keyword)确保 LIKE 匹配时不会因或破坏 SQL 结构。参数说明LIMIT 100是硬性限制Spark 在本地模式下查 1000 行可能卡顿100 行是用户体验和性能的平衡点timeout45比连接池默认 30 秒更宽裕因为用户搜索可能触发 Spark 新建 shuffle plan。4. 可视化大屏如何做到“秒级更新”Redis Hash 前端轮询的实战细节4.1 为什么不用 WebSocket 推送轮询更稳、更轻、更可控网上很多教程教“Django Channels WebSocket 实现实时推送”但对毕设而言这是过度设计。WebSocket 需要额外部署 ASGI 服务器Daphne、配置路由、处理连接保活、应对断网重连……而本项目用Redis Hash 前端 Ajax 轮询5 行代码搞定且稳定性远超 WebSocketRedis 写入是原子操作HSET不会丢数据前端轮询间隔可动态调整网络差时拉长到 10s好时缩到 2s无连接状态管理不占服务器长连接资源数据格式自由JSON string、list、hash不像 WebSocket 要序列化二进制。本项目约定所有统计类数据品牌分布、价格区间、热门关键词都写入 Redis HashKey 命名为stats:date:type例如stats:20240615:brand_top10。4.2 Spark 定时任务写 Redis用 foreachBatch 避免重复写入Spark Streaming 的foreachBatch是写外部系统的黄金 API——它把每个 micro-batch 当作一个 DataFrame 处理可批量写入 Redis且天然支持幂等同一 batch 不会重复触发。# spark_job/write_to_redis.py from pyspark.sql import SparkSession from pyspark.sql.functions import col, count, desc, hour, to_date import redis import json def write_stats_to_redis(df, epoch_id): 将 DataFrame 写入 Redis Hash df: 当前 batch 的清洗后数据schema 同 pipelines.py epoch_id: batch 序号用于日志 # 获取当天日期字符串用于分区 key today df.select(to_date(crawl_time).alias(dt)).first()[dt].strftime(%Y%m%d) # 1. 品牌 TOP10 brand_top10 df.groupBy(brand).count().orderBy(desc(count)).limit(10) brand_list [{name: row[brand], value: row[count]} for row in brand_top10.collect()] # 2. 价格区间分布0-3000, 3000-6000, 6000 price_bins df.select( col(price), (col(price) / 3000).cast(int).alias(bin_id) ).rdd.map(lambda x: ( 0-3000 if x.bin_id 0 else 3000-6000 if x.bin_id 1 else 6000 )).countByValue() price_dist [{name: k, value: v} for k, v in price_bins.items()] # 连接 Redis单例避免频繁创建连接 r redis.Redis(hostlocalhost, port6379, db0, decode_responsesTrue) # 写入 Hash自动创建不存在则新建 r.hset(fstats:{today}:brand_top10, mapping{data: json.dumps(brand_list)}) r.hset(fstats:{today}:price_dist, mapping{data: json.dumps(price_dist)}) # 设置过期时间24 小时避免 Redis 内存涨满 r.expire(fstats:{today}:brand_top10, 86400) r.expire(fstats:{today}:price_dist, 86400) print(f[Batch {epoch_id}] Wrote stats to Redis for {today}) # 主程序启动 Streaming 作业 if __name__ __main__: spark SparkSession.builder \ .appName(redis-writer) \ .master(local[*]) \ .config(spark.sql.adaptive.enabled, true) \ .getOrCreate() # 读取 Parquet 增量数据监听 /tmp/spark_raw/ raw_stream spark.readStream \ .format(parquet) \ .option(maxFilesPerTrigger, 10) \ .option(pathGlobFilter, *.parquet) \ .load(/tmp/spark_raw/) # 清洗同 2.2 节逻辑 clean_stream raw_stream.withColumn(clean_title, lower(regexp_replace(col(title), r[^\w\u4e00-\u9fa5], )) ).withColumn(price_num, col(price).cast(double)) # 写入 Redis每个 batch 触发一次 write_stats_to_redis query clean_stream.writeStream \ .foreachBatch(write_stats_to_redis) \ .outputMode(Append) \ .start() query.awaitTermination()逻辑说明foreachBatch是 Spark Streaming 3.0 的推荐写法替代了已废弃的foreach。它把每个 micro-batch 当作静态 DataFrame 处理可调用.collect()获取全部数据再批量写 Redis效率远高于逐行foreach。r.expire(..., 86400)是关键否则 Redis 内存会随天数线性增长。参数说明maxFilesPerTrigger10控制每次触发处理最多 10 个 Parquet 文件避免单次 batch 过大导致 OOMoutputModeAppend因为统计是追加式每天新数据不是更新式不需要Update模式。4.3 前端 ECharts 轮询 Redis防抖、节流、错误降级前端轮询不是简单setInterval(() fetch(...), 3000)必须加控制!-- templates/dashboard.html -- div idbrand-chart stylewidth: 600px;height:400px;/div script srchttps://cdn.jsdelivr.net/npm/echarts5.4.3/dist/echarts.min.js/script script let chart echarts.init(document.getElementById(brand-chart)); let lastFetchTime 0; let isFetching false; // 防抖函数1 秒内只发一次请求 function debounce(func, wait) { let timeout; return function executedFunction() { const later () { clearTimeout(timeout); func(...arguments); }; clearTimeout(timeout); timeout setTimeout(later, wait); }; } // 轮询函数带错误降级 async function fetchStats() { if (isFetching) return; isFetching true; const now Date.now(); // 节流至少间隔 2 秒 if (now - lastFetchTime 2000) { isFetching false; return; } lastFetchTime now; try { const response await fetch(/api/redis-stats/?keybrand_top10); if (!response.ok) throw new Error(HTTP ${response.status}); const data await response.json(); if (!data.data || !Array.isArray(data.data)) { throw new Error(Invalid data format); } // 渲染 ECharts const option { title: { text: 品牌销量 TOP10 }, tooltip: {}, xAxis: { type: category, data: data.data.map(d d.name) }, yAxis: { type: value }, series: [{ data: data.data.map(d d.value), type: bar }] }; chart.setOption(option, true); } catch (err) { console.warn(Fetch failed, using cached data:, err.message); // 错误时保持上次图表不闪退 } finally { isFetching false; } } // 启动轮询首次立即执行之后每 3 秒 fetchStats(); setInterval(fetchStats, 3000); // 页面隐藏时暂停显示时恢复节省资源 document.addEventListener(visibilitychange, () { if (document.hidden) { clearInterval(window.pollingInterval); } else { window.pollingInterval setInterval(fetchStats, 3000); } }); /script逻辑说明这段 JS 实现了企业级轮询的三大保障防抖节流双保险debounce防止用户快速切换 Tab 导致请求堆积2000ms节流避免网络抖动时高频重试错误降级catch块不报错弹窗而是console.warn并保持旧图表用户体验不中断资源感知visibilitychange事件监听页面可见性后台 Tab 自动暂停轮询省电省带宽。参数说明fetch(/api/redis-stats/?keybrand_top10)对应 Django 的一个简单 View它只是读 Redis 并返回 JSON代码不超过 10 行此处略去避免冗余。5. 避坑指南这 4 个坑让我重装 Spark 7 次第 8 次才跑通5.1 现象Spark Thrift Server 启动后PyHive 连接报TTransportException: Could not connect to localhost:10000原因Thrift Server 启动脚本默认绑定0.0.0.0但某些 Linux 发行版如 Ubuntu 22.04的防火墙ufw默认阻止 10000 端口或 Docker 环境下网络隔离。解决检查端口监听netstat -tuln | grep 10000确认输出中有127.0.0.1:10000若是0.0.0.0:10000修改启动命令显式指定--hiveconf hive.server2.thrift.bind.host127.0.0.1Ubuntu 用户临时关 ufwsudo ufw disable毕设环境可接受生产环境需sudo ufw allow 10000。5.2 现象Django View 执行 Spark SQL 时页面卡死 30 秒后报TimeoutError但spark-sqlCLI 命令行能秒出结果原因PyHive 的cursor.execute()默认不启用异步模式而 Thrift Server 的查询计划编译特别是第一次需要时间CLI 会等PyHive 却没设超时导致阻塞。解决在execute_spark_sql函数中必须加cursor.execute(query, async_True)注意下划线然后用while not cursor.is_executing(): time.sleep(0.1)轮询状态最终cursor.fetchall()才真正取数据。血泪经验官方文档藏得太深async_参数默认是False不显式设True就永远同步阻塞。5.3 现象爬虫写入 Parquet 后Spark Streaming 读不到新文件/tmp/spark_raw/目录明明有文件原因Spark Streaming 的readStream.format(parquet)默认只读取已完成写入的文件而 Scrapy pipeline 写 Parquet 是分块写入文件末尾.tmp未重命名前Spark 认为文件“未完成”跳过。解决在SparkParquetPipeline.process_item中写完 Parquet 后用os.rename()确保文件原子完成# 替换原 pipeline 中的 df.write...parquet(output_path) temp_path f{output_path}_tmp_{int(time.time())} df.write.mode(append).partitionBy(date).parquet(temp_path) # 原子重命名Linux/macOS 原子Windows 需用 shutil.move os.rename(temp_path, output_path)或更简单在start-thriftserver.sh启动前加参数--conf spark.sql.streaming.fileScanLocationLogCleaner.enabledfalse禁用文件扫描清理强制读所有文件。5.4 现象Redis 中stats:20240615:brand_top10的data字段是 JSON string但前端JSON.parse()报SyntaxError: Unexpected token in JSON at position 0原因Django 默认开启autoescapeView 返回 JSON 时json.dumps()产生的引号被转义成quot;Redis 存的是 HTML 实体编码字符串。解决在 Django View 中**返回 JSON 前用json.dumps(..., separators(,, :))本文还有配套的精品资源点击获取