本文为系列文章第一篇,完整项目包含三部分:

· 本篇:架构设计与环境搭建 ✅

· 第二篇:Spark实时处理与预警(即将发布)

· 第三篇:前后端实现与可视化(即将发布)

关注我,不错过后续更新!本文所有代码均可直接运行,欢迎学习交流!

🎯 项目背景

 

大家好,我是静思009,一名26届大数据专业的专科生,来自新疆和田。目前在实习阶段,虽然因为地域限制还没找到专业对口工作,但我相信技术能力才是硬道理!

 

今天开始分享一个完整的农业实时监控大数据平台项目,专门针对和田地区的农业场景设计。即使你暂时没有电脑,也能通过本系列文章理解完整架构和代码逻辑。

 

📊 项目架构设计

 

整体技术架构

数据源 → Kafka → Spark Streaming → HBase → Spring Boot → 可视化

   ↑ ↓ ↓ ↓ ↓

模拟数据 实时处理 预警分析 数据存储 数据接口

```

 

核心组件说明

 

· 数据采集层:Python数据模拟 + Kafka消息队列

· 实时处理层:Spark Streaming流式计算

· 数据存储层:HDFS + HBase + MySQL

· 服务层:Spring Boot微服务

· 展示层:ECharts可视化大屏

 

💻 环境搭建指南

 

伪分布式环境配置

 

```bash

# 1. 创建项目目录结构

mkdir -p ~/agriculture-platform/{code,data,logs,config}

cd ~/agriculture-platform

 

# 2. 检查环境依赖

echo "=== 环境检查 ==="

java -version

python3 --version

```

 

Hadoop单机模式配置

 

```bash

# 下载并配置Hadoop

wget https://archive.apache.org/dist/hadoop/common/hadoop-3.3.4/hadoop-3.3.4.tar.gz

tar -xzf hadoop-3.3.4.tar.gz

cd hadoop-3.3.4

 

# 配置环境变量

echo 'export HADOOP_HOME=/home/hadoop/agriculture-platform/hadoop-3.3.4' >> ~/.bashrc

echo 'export PATH=$PATH:$HADOOP_HOME/bin:$HADOOP_HOME/sbin' >> ~/.bashrc

source ~/.bashrc

```

 

📝 核心代码实现

 

1. 农业传感器数据模拟器

 

python

# sensor_simulator.py - 和田农业传感器数据模拟

import json

import random

import time

from datetime import datetime

 

class AgricultureSensor:

    """和田农业传感器数据模拟器"""

    

    def __init__(self):

        self.locations = ['和田县', '墨玉县', '皮山县', '洛浦县']

        self.crops = ['红枣', '核桃', '葡萄', '玉米', '小麦']

    

    def generate_sensor_data(self):

        """生成模拟传感器数据"""

        sensor_data = {

            'sensor_id': f"HT_{random.randint(1000, 9999)}",

            'timestamp': datetime.now().strftime('%Y-%m-%d %H:%M:%S'),

            'location': random.choice(self.locations),

            'crop_type': random.choice(self.crops),

            'temperature': round(25 + random.uniform(-5, 10), 1),

            'humidity': round(40 + random.uniform(-20, 30), 1),

            'soil_moisture': round(random.uniform(20, 80), 1),

            'light_intensity': random.randint(5000, 100000),

            'ph_value': round(6.0 + random.uniform(0.5, 2.5), 1)

        }

        

        # 5%概率生成异常数据(测试预警功能)

        if random.random() < 0.05:

            sensor_data['temperature'] = 45.0 # 高温异常

            sensor_data['soil_moisture'] = 15.0 # 干旱异常

            print(f"🚨 生成异常数据: {sensor_data['sensor_id']}")

            

        return sensor_data

    

    def run_simulation(self, count=20, interval=2):

        """运行数据模拟"""

        print("=== 和田农业传感器数据模拟开始 ===")

        print("位置:", self.locations)

        print("作物:", self.crops)

        print("=" * 50)

        

        for i in range(count):

            data = self.generate_sensor_data()

            print(f"数据 {i+1}: {json.dumps(data, ensure_ascii=False)}")

            

            # 保存到本地文件(模拟写入Kafka)

            with open('sensor_data.log', 'a', encoding='utf-8') as f:

                f.write(json.dumps(data, ensure_ascii=False) + '\n')

            

            time.sleep(interval)

        

        print("=== 数据模拟完成 ===")

        print(f"共生成 {count} 条数据,已保存到 sensor_data.log")

 

# 运行示例

if __name__ == "__main__":

    sensor = AgricultureSensor()

    sensor.run_simulation(count=10, interval=1)

```

 

2. Hive数据仓库建设

 

sql

 创建农业监控数据库

CREATE DATABASE IF NOT EXISTS xj_agriculture_monitor;

USE xj_agriculture_monitor;

传感器数据表

CREATE TABLE IF NOT EXISTS sensor_data (

    sensor_id STRING COMMENT '传感器ID',

    timestamp STRING COMMENT '数据时间戳',

    location STRING COMMENT '地理位置',

    crop_type STRING COMMENT '作物类型',

    temperature DOUBLE COMMENT '温度',

    humidity DOUBLE COMMENT '湿度',

    soil_moisture DOUBLE COMMENT '土壤湿度',

    light_intensity INT COMMENT '光照强度',

    ph_value DOUBLE COMMENT 'PH值'

ROW FORMAT DELIMITED

FIELDS TERMINATED BY ','

STORED AS TEXTFILE;

预警规则表

CREATE TABLE IF NOT EXISTS alert_rules (

    rule_id INT COMMENT '规则ID',

    rule_name STRING COMMENT '规则名称',

    rule_condition STRING COMMENT '规则条件',

    alert_level STRING COMMENT '预警级别',

    create_time TIMESTAMP COMMENT '创建时间'

);

插入示例预警规则

INSERT INTO alert_rules VALUES

(1, '高温预警', 'temperature > 40', 'HIGH', current_timestamp()),

(2, '干旱预警', 'soil_moisture < 20', 'MEDIUM', current_timestamp()),

(3, '酸碱度异常', 'ph_value < 6.0 OR ph_value > 8.0', 'LOW', current_timestamp());

 创建数据统计视图

CREATE VIEW IF NOT EXISTS region_stats AS

SELECT 

    location,

    crop_type,

    AVG(temperature) as avg_temperature,

    AVG(humidity) as avg_humidity,

    AVG(soil_moisture) as avg_soil_moisture,

    COUNT(*) as data_count

FROM sensor_data 

GROUP BY location, crop_type;

 查询示例:各区域农业数据统计

SELECT 

    location as 区域,

    crop_type as 作物类型,

    ROUND(avg_temperature, 1) as 平均温度,

    ROUND(avg_humidity, 1) as 平均湿度,

    ROUND(avg_soil_moisture, 1) as 平均土壤湿度,

    data_count as 数据量

FROM region_stats 

ORDER BY location, crop_type;

3. 数据ETL处理脚本

 

python

# data_etl.py - 数据ETL处理

import json

import pandas as pd

from datetime import datetime

 

class AgricultureETL:

    农业数据ETL处理类

    

    def __init__(self):

        self.alert_rules = [

            {'name': '高温预警', 'condition': lambda x: x['temperature'] > 40},

            {'name': '干旱预警', 'condition': lambda x: x['soil_moisture'] < 20},

            {'name': '酸碱度异常', 'condition': lambda x: x['ph_value'] < 6.0 or x['ph_value'] > 8.0}

        ]

    

    def load_sensor_data(self, file_path):

        """加载传感器数据"""

        data_list = []

        with open(file_path, 'r', encoding='utf-8') as f:

            for line in f:

                try:

                    data = json.loads(line.strip())

                    data_list.append(data)

                except json.JSONDecodeError:

                    print(f"数据格式错误: {line}")

        

        return data_list

    

    def detect_alerts(self, data):

        """检测预警数据"""

        alerts = []

        for rule in self.alert_rules:

            if rule['condition'](data):

                alerts.append({

                    'sensor_id': data['sensor_id'],

                    'alert_type': rule['name'],

                    'timestamp': data['timestamp'],

                    'location': data['location'],

                    'crop_type': data['crop_type'],

                    'value': self._get_alert_value(data, rule['name'])

                })

        return alerts

    

    def _get_alert_value(self, data, alert_type):

        """获取预警值"""

        if alert_type == '高温预警':

            return data['temperature']

        elif alert_type == '干旱预警':

            return data['soil_moisture']

        elif alert_type == '酸碱度异常':

            return data['ph_value']

        return None

    

    def generate_report(self, data_list):

        """生成数据报告"""

        df = pd.DataFrame(data_list)

        

        print("=== 农业数据统计报告 ===")

        print(f"数据时间范围: {df['timestamp'].min()} 至 {df['timestamp'].max()}")

        print(f"总数据量: {len(df)} 条")

        print("\n各区域数据分布:")

        location_stats = df['location'].value_counts()

        for 

def generate_report(self, data_list):

    """生成数据报告"""

    df = pd.DataFrame(data_list)

    

    print("=== 农业数据统计报告 ===")

    print(f"数据时间范围: {df['timestamp'].min()} 至 {df['timestamp'].max()}")

    print(f"总数据量: {len(df)} 条")

    

    # 补全:各区域数据分布

    print("\n各区域数据分布:")

    location_stats = df['location'].value_counts()

    for location, count in location_stats.items():

        print(f" {location}: {count} 条")

    

    # 新增:各作物类型分布

    print("\n各作物类型分布:")

    crop_stats = df['crop_type'].value_counts()

    for crop, count in crop_stats.items():

        print(f" {crop}: {count} 条")

    

    # 新增:预警统计

    print("\n预警情况统计:")

    total_alerts = []

    for data in data_list:

        total_alerts.extend(self.detect_alerts(data))

    if total_alerts:

        alert_df = pd.DataFrame(total_alerts)

        alert_type_stats = alert_df['alert_type'].value_counts()

        for alert_type, count in alert_type_stats.items():

            print(f" {alert_type}: {count} 次")

    else:

        print(" 暂无预警数据")

# agri_etl.py 补充主函数

if __name__ == "__main__":

    # 1. 先通过sensor_simulator.py生成sensor_data.log(模拟传感器数据)

    print("=== 开始农业数据ETL处理 ===")

    

    # 2. 初始化ETL处理器

    etl_processor = AgricultureETL()

    

    # 3. 加载模拟生成的传感器日志

    data_list = etl_processor.load_sensor_data("sensor_data.log")

    

    # 4. 输出实时预警信息

    all_alerts = []

    for data in data_list:

        alerts = etl_processor.detect_alerts(data)

        if alerts:

            all_alerts.extend(alerts)

    if all_alerts:

        print("\n=== 实时预警信息 ===")

        for alert in all_alerts:

            print(f"[{alert['timestamp']}] {alert['location']}·{alert['crop_type']} 触发{alert['alert_type']}:{alert['value']}")

    

    # 5. 生成统计报告

    etl_processor.generate_report(data_list)

 

本篇章聚焦农业监控平台的“基础链路搭建”:从“架构设计”到“环境配置”,再到“传感器模拟、Hive存储、ETL处理”,完整实现了项目的“数据生产→存储→初步加工”链路。

 

下一篇将讲解Spark Streaming实时计算与预警推送,第三篇会落地“Spring Boot接口+ECharts可视化”的前后端功能,逐步完成整个平台的实战闭环~

 

(本篇内容结束,后续篇章可关注更新获取完整代码)

更多推荐