大家好!我是一名新疆和田大数据技术专业专科生,因为家庭原因回到和田虽然没有找到专业相关工作但是依然在有空的时候探索深耕并学习技术,不多废话了跟我来

🔥 地域特色+实用场景+轻量化技术栈!新手也能上手的大数据实战项目,解决和田农产品“溯源难、滞销多”痛点~

 1. 项目背景

 

和田地区的红枣、核桃、和田玉枣等特色农产品享誉全国,但长期面临两大问题:① 消费者担心“假特产”,缺乏可信溯源渠道;② 农户盲目种植,易出现“丰收不增收”的滞销困境。

 

本项目以轻量级大数据技术为核心,打造“溯源+预测”一体化系统:既通过数据链实现农产品从种植到销售的全程可查,又基于历史销量数据预测未来需求,帮助农户精准规划生产、商家智能备货。

 

2. 核心亮点(创新性+实践性)

 

🌍 地域绑定:聚焦和田特色农产品,数据来源真实可获取(公开电商数据+模拟溯源数据)

🚀 轻量易实现:不依赖复杂集群,单机即可部署,技术栈门槛低

 🎯 实用导向:同时解决消费者溯源需求和农户滞销痛点,落地性强

 📊 全流程覆盖:从数据采集→预处理→分析→预测→可视化,完整大数据链路

 

 3. 技术栈选型(轻量化为主)

 

- 数据存储:HDFS(存储原始数据)+ MySQL(存储结构化溯源/销量数据)

- 数据处理:MapReduce(批量预处理)+ Python(Pandas/Numpy 清洗)

- 预测模型:LightGBM(轻量级梯度提升树,比LSTM更易训练)

- 可视化&溯源查询:Flask(简易Web页面,支持溯源码查询+销量趋势展示)

- 数据采集:Python爬虫(爬取电商平台和田农产品销量数据)+ Excel导入(农户/合作社溯源基础数据)

 

二、项目思维导图(清晰梳理全流程)

 

 三、核心逻辑题(大数据场景必遇)

 

在项目开发中,以下逻辑题对应实际业务场景,先思考再看解析~

 

1. 逻辑题1:溯源码去重(MapReduce场景)

 

问题:采集的和田红枣溯源数据中,存在重复的“溯源码-农户ID”组合(同一批农产品被重复录入),请设计MapReduce逻辑,保留每条溯源码的第一条有效记录。

 

解析思路:

 

- Map阶段:以“溯源码”为key,value为“农户ID+种植日期+加工厂家”完整记录

- Reduce阶段:对同一个key的value列表,按“录入时间”升序排序,取第一个元素输出

 

2. 逻辑题2:销量异常值识别(数据预处理场景)

 

问题:电商爬取的销量数据中,存在“单日销量为0”(商品下架)或“单日销量远超均值10倍”(促销活动)的异常值,如何筛选出“正常销售日”的数据用于预测?

 

解析思路:

 

- 计算近30天销量均值μ和标准差σ

- 异常值判定规则:销量 < 1 或 销量 > μ + 3σ(3σ原则)

- 保留 1 ≤ 销量 ≤ μ + 3σ 的数据,既排除下架状态,又剔除极端促销干扰

 

3. 逻辑题3:溯源码关联查询(MySQL场景)

 

问题:用户输入溯源码后,需展示“种植户信息→种植地块→加工企业→物流信息→销售渠道”完整链路,已知MySQL中有4张表(种植户表、加工表、物流表、销售表),均含“溯源码”字段,如何设计SQL查询?

 

解析思路:

 

sql

SELECT 

  z.农户姓名, z.地块位置, z.种植日期,

  j.加工企业名称, j.加工日期,

  w.物流公司, w.物流单号, w.发货日期,

  x.销售平台, x.销售日期

FROM 种植户表 z

JOIN 加工表 j ON z.溯源码 = j.溯源码

JOIN 物流表 w ON z.溯源码 = w.溯源码

JOIN 销售表 x ON z.溯源码 = x.溯源码

WHERE z.溯源码 = 'HT20250601001'; -- 用户输入的溯源码

 

 

四、核心代码及讲解(可直接复制运行)

 

模块1:数据采集(Python爬虫+Excel导入)

 

1.1 电商销量爬虫(爬取京东和田红枣销量数据)

 

python

import requests

from lxml import etree

import pandas as pd

import time

 

# 目标URL:京东和田红枣搜索结果页(简化版,实际可分页爬取)

url = "https://search.jd.com/Search?keyword=和田红枣&enc=utf-8"

headers = {

    "User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36"

}

 

def crawl_jd_sales():

    sales_data = []

    try:

        response = requests.get(url, headers=headers, timeout=10)

        response.encoding = "utf-8"

        html = etree.HTML(response.text)

        

        # 解析商品名称、价格、销量、店铺名

        names = html.xpath('//div[@class="p-name p-name-type-2"]/a/em/text()')

        prices = html.xpath('//div[@class="p-price"]/strong/i/text()')

        sales = html.xpath('//div[@class="p-commit"]/strong/a/text()')

        shops = html.xpath('//div[@class="p-shop"]/span/a/text()')

        

        # 数据整理(取前10条,避免数据过多)

        for i in range(min(10, len(names))):

            sales_data.append({

                "商品名称": names[i].strip(),

                "价格": float(prices[i]) if prices[i] else 0.0,

                "销量": sales[i].replace("+", "").replace("条评价", "") if sales[i] else "0",

                "店铺名": shops[i] if shops[i] else "未知店铺",

                "爬取日期": time.strftime("%Y-%m-%d", time.localtime())

            })

        

        # 保存为CSV,后续导入HDFS

        df = pd.DataFrame(sales_data)

        df.to_csv("hetian_red_date_sales.csv", index=False, encoding="utf-8-sig")

        print("爬虫成功!已保存销量数据到CSV文件")

        

    except Exception as e:

        print(f"爬虫失败:{str(e)}")

 

if __name__ == "__main__":

    crawl_jd_sales()

 

 

代码讲解:

 

- 用requests+lxml爬取京东公开搜索结果,避开反爬(仅爬取10条,非高频请求)

- 解析核心字段:商品名称、价格、销量、店铺名,保存为CSV方便后续处理

- 异常处理:避免网络超时或页面结构变化导致程序崩溃

 

1.2 溯源基础数据导入(Excel→MySQL)

 

python

import pandas as pd

import pymysql

 

# 连接MySQL(需提前创建数据库hetian_agri)

db = pymysql.connect(

    host="localhost",

    user="root",

    password="123456",

    database="hetian_agri",

    charset="utf8"

)

cursor = db.cursor()

 

# 读取Excel溯源数据(可由农户/合作社提供模板填写)

df = pd.read_excel("hetian_trace_data.xlsx")

 

# 创建溯源表

create_table_sql = """

CREATE TABLE IF NOT EXISTS trace_info (

    trace_code VARCHAR(20) PRIMARY KEY, -- 溯源码(唯一标识)

    farmer_name VARCHAR(50) NOT NULL, -- 种植户姓名

    land_location VARCHAR(100) NOT NULL, -- 种植地块位置

    plant_date DATE NOT NULL, -- 种植日期

    process_company VARCHAR(50) NOT NULL, -- 加工企业

    process_date DATE NOT NULL, -- 加工日期

    logistics_company VARCHAR(50) -- 物流公司

)

"""

cursor.execute(create_table_sql)

db.commit()

 

# 批量插入数据

def insert_trace_data():

    try:

        for index, row in df.iterrows():

            sql = """

            INSERT INTO trace_info (trace_code, farmer_name, land_location, plant_date, process_company, process_date, logistics_company)

            VALUES (%s, %s, %s, %s, %s, %s, %s)

            ON DUPLICATE KEY UPDATE farmer_name = VALUES(farmer_name)

            """

            cursor.execute(sql, (

                row["溯源码"],

                row["种植户姓名"],

                row["地块位置"],

                row["种植日期"],

                row["加工企业"],

                row["加工日期"],

                row["物流公司"]

            ))

        db.commit()

        print("溯源数据成功导入MySQL!")

    except Exception as e:

        db.rollback()

        print(f"数据导入失败:{str(e)}")

    finally:

        db.close()

 

if __name__ == "__main__":

    insert_trace_data()

 

 

代码讲解:

 

- 提前创建Excel模板(含溯源码、种植户、地块等字段),农户可手动填写

- 用pandas读取Excel,pymysql连接MySQL批量插入数据,支持重复溯源码更新

- 适合小批量溯源数据导入,符合和田当地合作社的实际操作场景

 

模块2:数据预处理(MapReduce去重+Python清洗)

 

2.1 MapReduce溯源数据去重(Java代码)

 

java

import org.apache.hadoop.conf.Configuration;

import org.apache.hadoop.fs.Path;

import org.apache.hadoop.io.Text;

import org.apache.hadoop.mapreduce.Job;

import org.apache.hadoop.mapreduce.Mapper;

import org.apache.hadoop.mapreduce.Reducer;

import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;

import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;

 

import java.io.IOException;

import java.util.ArrayList;

import java.util.Collections;

import java.util.Comparator;

import java.util.List;

 

// Mapper:读取原始溯源数据,输出<溯源码, 完整记录>

class TraceDedupMapper extends Mapper<Object, Text, Text, Text> {

    private Text traceCode = new Text();

    private Text fullRecord = new Text();

 

    @Override

    protected void map(Object key, Text value, Context context) throws IOException, InterruptedException {

        // 原始数据格式:溯源码,种植户姓名,地块位置,种植日期,加工企业,录入时间

        String[] fields = value.toString().split(",");

        if (fields.length >= 6) {

            traceCode.set(fields[0]); // 溯源码作为key

            fullRecord.set(value.toString()); // 完整记录作为value

            context.write(traceCode, fullRecord);

        }

    }

}

 

// Reducer:对同一溯源码的记录去重,保留录入时间最早的一条

class TraceDedupReducer extends Reducer<Text, Text, Text, Text> {

    @Override

    protected void reduce(Text key, Iterable<Text> values, Context context) throws IOException, InterruptedException {

         List<String> records = new ArrayList<>();

         for (Text val : values) {

             records.add(val.toString());

         }

         // 按录入时间升序排序,取第一条

         Collections.sort(records, new Comparator<String>() {

             @Override

             public int compare(String s1, String s2) {

                 String time1 = s1.split(",")[5];

                 String time2 = s2.split(",")[5];

                 return time1.compareTo(time2);

             }

         });

         context.write(key, new Text(records.get(0)));

     }

 }

 // 主类:提交MapReduce任务

 public class TraceDataDedup {

     public static void main(String[] args) throws Exception {

         Configuration conf = new Configuration();

         Job job = Job.getInstance(conf, "trace_data_dedup");

         

         job.setJarByClass(TraceDataDedup.class);

         job.setMapperClass(TraceDedupMapper.class);

         job.setReducerClass(TraceDedupReducer.class);

         

         job.setOutputKeyClass(Text.class);

         job.setOutputValueClass(Text.class);

         

         // 输入路径:HDFS上的原始溯源数据,输出路径:去重后的数据(需不存在)

         FileInputFormat.addInputPath(job, new Path(args[0]));

         FileOutputFormat.setOutputPath(job, new Path(args[1]));

         

         System.exit(job.waitForCompletion(true) ? 0 : 1);

     }

 }

代码讲解:

  解决原始溯源数据重复问题,核心是“按溯源码分组,按录入时间排序取第一条”

​可在Windows本地用Hadoop伪分布式运行,无需搭建复杂集群

​编译运行命令: hadoop jar TraceDataDedup.jar /input/trace_raw /output/trace_dedup 

 2.2 Python销量数据清洗(处理异常值/缺失值)

import pandas as pd

import numpy as np

 

def clean_sales_data():

    # 读取爬虫获取的销量数据

    df = pd.read_csv("hetian_red_date_sales.csv", encoding="utf-8-sig")

    

    # 1. 处理销量字段(统一格式,去除非数字字符)

    df["销量"] = df["销量"].str.replace("万", "").str.replace("+", "").str.replace(",", "")

    df["销量"] = pd.to_numeric(df["销量"], errors="coerce").fillna(0)

    

    # 2. 异常值处理(3σ原则)

    mu = df["销量"].mean()

    sigma = df["销量"].std()

    df = df[(df["销量"] >= 1) & (df["销量"] <= mu + 3*sigma)]

    

    # 3. 缺失值处理(店铺名为空填充为“未知店铺”)

    df["店铺名"] = df["店铺名"].fillna("未知店铺")

    

    # 4. 保存清洗后的数据,用于后续预测

    df.to_csv("cleaned_sales_data.csv", index=False, encoding="utf-8-sig")

    print("销量数据清洗完成!清洗后数据量:", len(df))

    print("清洗后数据预览:")

    print(df.head())

 

if __name__ == "__main__":

    clean_sales_data()

模块3:销量预测(LightGBM模型)

import pandas as pd

import numpy as np

from sklearn.model_selection import train_test_split

from sklearn.metrics import mean_squared_error

import lightgbm as lgb

import matplotlib.pyplot as plt

 

# 设置中文显示

plt.rcParams['font.sans-serif'] = ['SimHei']

plt.rcParams['axes.unicode_minus'] = False

 

def sales_prediction():

    # 1. 加载清洗后的销量数据(需补充日期字段,模拟30天销量)

    df = pd.read_csv("cleaned_sales_data.csv", encoding="utf-8-sig")

    # 模拟日期序列(实际可从爬虫数据中提取历史日期)

    df["日期"] = pd.date_range(start="2025-01-01", periods=len(df), freq="D")

    df = df.sort_values("日期")

    

    # 2. 特征工程(简单特征,适合轻量级模型)

    df["星期"] = df["日期"].dt.dayofweek # 星期特征(0-6)

    df["是否周末"] = df["星期"].apply(lambda x: 1 if x >=5 else 0)

    df["价格等级"] = pd.cut(df["价格"], bins=[0, 30, 60, 100], labels=[1, 2, 3]) # 价格分档

    

    # 3. 划分训练集和测试集

    X = df[["价格", "星期", "是否周末", "价格等级"]]

    y = df["销量"]

    X_train, X_test, y_train, y_test = train_test_split(X, y, test_size=0.2, random_state=42)

    

    # 4. 训练LightGBM模型

    lgb_model = lgb.LGBMRegressor(

        n_estimators=100, # 弱学习器数量

        learning_rate=0.1, # 学习率

        max_depth=3, # 树深度(避免过拟合)

        random_state=42

    )

    lgb_model.fit(X_train, y_train, eval_set=[(X_test, y_test)], early_stopping_rounds=10)

    

    # 5. 预测与评估

    y_pred = lgb_model.predict(X_test)

    mse = mean_squared_error(y_test, y_pred)

    rmse = np.sqrt(mse)

    print(f"模型评估:RMSE = {rmse:.2f}")

    

    # 6. 可视化预测结果

    plt.figure(figsize=(12, 6))

    plt.plot(y_test.values, label="实际销量")

    plt.plot(y_pred, label="预测销量", linestyle="--")

    plt.xlabel("测试集样本")

    plt.ylabel("销量")

    plt.title("和田红枣销量预测结果对比")

    plt.legend()

    plt.savefig("sales_prediction_result.png", dpi=300, bbox_inches="tight")

    plt.show()

    

    # 7. 保存模型(用于后续部署)

    import joblib

    joblib.dump(lgb_model, "sales_prediction_model.pkl")

    print("模型已保存为 sales_prediction_model.pkl")

 

if __name__ == "__main__":

    sales_prediction()

代码讲解:

  选用LightGBM轻量级模型,训练速度快、占用资源少,适合单机部署

特征工程简单实用,无需复杂特征构造,新手易理解

 可视化预测结果,直观展示模型效果,支持后续优化

​ 保存模型文件,可直接用于FlaskWeb部署

 

模块4:可视化&溯源查询(FlaskWeb页面)

from flask import Flask, render_template, request

import pymysql

import joblib

import pandas as pd

 

app = Flask(__name__)

 

# 加载预测模型

model = joblib.load("sales_prediction_model.pkl")

 

# 连接MySQL(溯源查询用)

def get_db_connection():

    return pymysql.connect(

        host="localhost",

        user="root",

        password="123456",

        database="hetian_agri",

        charset="utf8"

    )

 

# 首页:包含溯源查询和销量预测入口

@app.route("/")

def index():

    return render_template("index.html")

 

# 溯源查询接口

@app.route("/trace", methods=["POST"])

def trace_query():

    trace_code = request.form.get("trace_code")

    db = get_db_connection()

    cursor = db.cursor(pymysql.cursors.DictCursor)

    

    try:

        cursor.execute("SELECT * FROM trace_info WHERE trace_code = %s", (trace_code,))

        result = cursor.fetchone()

        if result:

            # 格式化日期字段

            result["plant_date"] = result["plant_date"].strftime("%Y-%m-%d")

            result["process_date"] = result["process_date"].strftime("%Y-%m-%d")

            return render_template("trace_result.html", data=result)

        else:

            return "未查询到该溯源码对应的和田农产品信息,请检查输入!"

    except Exception as e:

        return f"查询失败:{str(e)}"

    finally:

        db.close()

 

# 销量预测接口

@app.route("/predict", methods=["POST"])

def sales_predict():

    # 获取用户输入的特征(价格、星期、是否周末、价格等级)

    price = float(request.form.get("price"))

    weekday = int(request.form.get("weekday"))

    is_weekend = int(request.form.get("is_weekend"))

    price_level = int(request.form.get("price_level"))

    

    # 构造特征数据

    features = pd.DataFrame({

        "价格": [price],

        "星期": [weekday],

        "是否周末": [is_weekend],

        "价格等级": [price_level]

    })

    

    # 预测销量

    prediction = model.predict(features)[0]

    return render_template("predict_result.html", prediction=round(prediction, 0))

 

if __name__ == "__main__":

    # 开启调试模式,便于开发

    app.run(debug=True, port=5000)

五、常见问题及针对性解决方案

 这个轻量级项目完美结合了和田地域特色与大数据实用场景,技术栈门槛低、实践性强,新手跟着代码一步步操作就能跑通全流程!

 

✨ 如果你觉得项目有用,别忘了点赞+关注哦~ 后续会更新更多大数据实战项目!

💬 互动提问:你觉得这个项目还能怎么优化?比如如何获取更精准的和田农产品数据?或者有其他地域特色的大数据项目想法?欢迎在评论区留言讨论!

 

更多推荐