【免费】基于Spark实时社交媒体舆情分析与趋势预测(Python版本+pyspark+可视化大屏+Kafka+FastAPI+Vue3) 锋哥原创出品,必属精品
大家好我是Java1234_小锋老师分享一套锋哥原创的基于Spark实时社交媒体舆情分析与趋势预测(Python版本pyspark可视化大屏KafkaFastAPIVue3)项目介绍随着微博、抖音、知乎、小红书等社交媒体的快速发展网络舆情呈现出数据规模大、传播速度快、情感变化剧烈等特点。传统基于离线批处理的舆情分析方法难以满足“秒级感知、分钟级研判”的业务需求。针对上述问题本文设计并实现了一套基于 Spark 的实时社交媒体舆情分析与趋势预测系统。系统采用前后端分离架构前端基于 Vue3、Element Plus 与 ECharts 构建管理后台和可视化大屏后端基于 Python 与 FastAPI 提供统一 REST API消息层引入 Kafka 承接高并发舆情事件流计算层使用 Spark StreamingStructured Streaming完成按小时窗口的帖文量、独立用户数、正中负情感分布与热度指数聚合预测层基于 Spark ML 线性回归结合滞后热度特征与小时特征对舆情热度进行趋势预测并以 RMSE、MAE、MAPE 评估模型误差。数据持久化采用 MySQL数据库名为 db_social_opinion。系统还设计了 Kafka/Spark 不可用时的 pandas 与 scikit-learn 降级方案保证演示与实验环境的可用性。测试结果表明系统能够稳定完成舆情数据采集、实时统计、趋势预测与可视化展示功能完整、结构清晰达到本科毕业设计要求。源码下载链接: https://pan.baidu.com/s/15-uuTt3lRlIFc0AzRH8ANw?pwd1234提取码: 1234系统展示核心代码 Spark Streaming 流式计算 - 消费 Kafka 数据并实时统计舆情 from config import settings def run_spark_streaming(events: list None) - list: 运行 Spark Structured Streaming 处理社交媒体数据流 若传入 events 列表则直接处理用于降级模式复用逻辑 返回窗口统计结果列表 try: from pyspark.sql import SparkSession from pyspark.sql.functions import ( col, count, sum as spark_sum, countDistinct, window, from_json, to_timestamp ) from pyspark.sql.types import ( StructType, StructField, StringType, IntegerType, DoubleType ) spark SparkSession.builder \ .appName(settings.SPARK_APP_NAME) \ .master(settings.SPARK_MASTER) \ .config(spark.sql.shuffle.partitions, 4) \ .config(spark.driver.memory, 2g) \ .getOrCreate() spark.sparkContext.setLogLevel(WARN) schema StructType([ StructField(user_id, IntegerType()), StructField(platform_id, IntegerType()), StructField(topic_id, IntegerType()), StructField(content, StringType()), StructField(sentiment, StringType()), StructField(heat, DoubleType()), StructField(event_time, StringType()), ]) if events: df spark.createDataFrame(events) else: raw_df spark.readStream \ .format(kafka) \ .option(kafka.bootstrap.servers, settings.KAFKA_BOOTSTRAP_SERVERS) \ .option(subscribe, settings.KAFKA_TOPIC) \ .option(startingOffsets, earliest) \ .load() df raw_df.select( from_json(col(value).cast(string), schema).alias(data) ).select(data.*) df df.withColumn(ts, to_timestamp(col(event_time), yyyy-MM-dd HH:mm:ss)) df df.filter(col(ts).isNotNull()) windowed df.groupBy(window(col(ts), 1 hour)).agg( count(*).alias(post_count), countDistinct(user_id).alias(uv), spark_sum((col(sentiment) positive).cast(int)).alias(positive), spark_sum((col(sentiment) neutral).cast(int)).alias(neutral), spark_sum((col(sentiment) negative).cast(int)).alias(negative), spark_sum(heat).alias(heat_index), ) if events: rows windowed.collect() results [] for row in rows: start row[window].start results.append({ window_time: start.strftime(%Y-%m-%d %H:00:00), post_count: int(row[post_count]), uv: int(row[uv]), positive: int(row[positive]), neutral: int(row[neutral]), negative: int(row[negative]), heat_index: float(row[heat_index]), }) spark.stop() return results spark.stop() return [] except Exception as e: print(f[Spark Streaming] 运行失败: {e}) return None def save_stats_to_db(stats: list): 将统计结果写入数据库 from database import SessionLocal from models.realtime_stat import RealtimeStat db SessionLocal() try: for s in stats: existing db.query(RealtimeStat).filter( RealtimeStat.window_time s[window_time] ).first() if existing: for k, v in s.items(): setattr(existing, k, v) else: db.add(RealtimeStat(**s)) db.commit() print(f[Spark Streaming] 已写入 {len(stats)} 条统计结果) finally: db.close()template div classpage-container div classpage-card div classpage-title舆情热度预测分析/div div classerror-cards div classerror-card div classmetric-labelRMSE (均方根误差)/div div classmetric-value{{ errorMetric.rmse }}/div /div div classerror-card div classmetric-labelMAE (平均绝对误差)/div div classmetric-value{{ errorMetric.mae }}/div /div div classerror-card div classmetric-labelMAPE (平均绝对百分比误差 %)/div div classmetric-value{{ errorMetric.mape }}%/div /div /div div refcompareRef classpred-chart pred-chart-compare/div div refresidualRef classpred-chart pred-chart-residual/div el-table :datatableData stripe border stylewidth:100% el-table-column propwindow_time label时间窗口 min-width170 template #default{ row }{{ formatWindowTime(row.window_time) }}/template /el-table-column el-table-column proptrue_heat label真实热度 min-width120 template #default{ row }span stylecolor:#409eff;font-weight:600{{ row.true_heat }}/span/template /el-table-column el-table-column proppred_heat label预测热度 min-width120 template #default{ row }span stylecolor:#67c23a;font-weight:600{{ row.pred_heat }}/span/template /el-table-column el-table-column label误差 min-width100 template #default{ row } span :style{ color: Math.abs(row.true_heat - row.pred_heat) 50 ? #f56c6c : #909399 } {{ (row.true_heat - row.pred_heat).toFixed(2) }} /span /template /el-table-column el-table-column propcreate_time label生成时间 min-width170 template #default{ row }{{ formatDateTime(row.create_time) }}/template /el-table-column /el-table el-pagination stylemargin-top:16px;justify-content:flex-end v-model:current-pagepage v-model:page-sizesize :totaltotal layouttotal, prev, pager, next changeloadTable / /div /div /template script setup /** * 预测分析页面真实 vs 预测对比图 误差分析 */ import { ref, onMounted, onUnmounted } from vue import * as echarts from echarts import request from /utils/request import { formatDateTime, formatWindowTime } from /utils/format const errorMetric ref({ rmse: 0, mae: 0, mape: 0 }) const tableData ref([]) const page ref(1) const size ref(10) const total ref(0) const compareRef ref(null) const residualRef ref(null) let charts [] function buildAxisLabel() { return { rotate: 30, interval: auto, fontSize: 11, margin: 16, formatter(val) { const text formatWindowTime(val); return text.length 16 ? ${text.slice(0,10)}\n${text.slice(11)} : text }, } } function initCompareChart(data) { const chart echarts.init(compareRef.value) const labels data.map(d formatWindowTime(d.window_time)) chart.setOption({ title: { text: 真实热度 vs 预测热度 对比, left: center, textStyle: { fontSize: 15 } }, tooltip: { trigger: axis }, legend: { data: [真实热度, 预测热度], top: 32 }, xAxis: { type: category, data: labels, axisLabel: buildAxisLabel() }, yAxis: { type: value, name: 热度指数 }, series: [ { name: 真实热度, type: line, smooth: true, data: data.map(d Number(d.true_heat)), itemStyle: { color: #409eff }, lineStyle: { width: 3 } }, { name: 预测热度, type: line, smooth: true, data: data.map(d Number(d.pred_heat)), itemStyle: { color: #67c23a }, lineStyle: { width: 3, type: dashed } }, ], grid: { left: 20, right: 24, bottom: 28, top: 72, containLabel: true }, }) charts.push(chart) } function initResidualChart(data) { const chart echarts.init(residualRef.value) const labels data.map(d formatWindowTime(d.window_time)) chart.setOption({ title: { text: 预测残差分析 (真实值 - 预测值), left: center, textStyle: { fontSize: 15 } }, tooltip: { trigger: axis }, xAxis: { type: category, data: labels, axisLabel: buildAxisLabel() }, yAxis: { type: value, name: 残差 }, series: [{ type: bar, data: data.map(d ({ value: d.residual, itemStyle: { color: d.residual 0 ? #409eff : #f56c6c } })), barWidth: 20 }], grid: { left: 20, right: 24, bottom: 28, top: 56, containLabel: true }, }) charts.push(chart) } async function loadData() { const [errorRes, compareRes, residualRes] await Promise.all([ request.get(/prediction/error), request.get(/prediction/compare), request.get(/prediction/residual), ]) errorMetric.value errorRes.data charts.forEach(c c.dispose()) charts [] initCompareChart(compareRes.data) initResidualChart(residualRes.data) } async function loadTable() { const res await request.get(/prediction/list, { params: { page: page.value, size: size.value } }) tableData.value res.data.items total.value res.data.total } onMounted(() { loadData(); loadTable() }) onUnmounted(() charts.forEach(c c.dispose())) /script style scoped .pred-chart { width: 100%; margin-bottom: 24px; } .pred-chart-compare { height: 480px; } .pred-chart-residual { height: 420px; } /style