告别Java门槛:用Python+MapReduce轻松清洗海量网站日志

如果你所在的技术团队正面临这样的困境:业务数据量日益增长,每天产生的网站日志动辄几十GB,传统的单机脚本处理已经力不从心,但团队里又缺乏专业的Java大数据开发人员,那么这篇文章正是为你准备的。

很多中小型技术团队在初次接触大数据处理时,往往被Hadoop生态的Java开发门槛吓退。传统的MapReduce开发需要编写冗长的Java代码,配置复杂的依赖环境,调试过程更是令人头疼。但实际上,Hadoop提供了灵活的Hadoop Streaming接口,允许你使用任何能够处理标准输入输出的语言来编写MapReduce程序,而Python正是其中最受欢迎的选择之一。

今天,我将带你绕过Java的复杂性,直接使用Python脚本配合Hadoop Streaming,构建一套完整的网站日志清洗与分析流水线。这套方案不仅降低了技术门槛,还能让你在几天内就搭建起可用的日志分析系统,快速验证业务需求。

1. 环境准备与数据理解

在开始编写代码之前,我们需要先搭建好基础环境。对于中小团队,我建议从云服务商提供的托管Hadoop服务开始,比如阿里云的E-MapReduce(EMR)或AWS的EMR。这些服务简化了集群的部署和维护,让你能专注于业务逻辑开发。

如果你选择自建环境,需要确保以下组件就绪:

  • Hadoop集群(建议版本2.7+或3.x)
  • Python 3.6+(确保所有节点都安装)
  • Hadoop Streaming JAR包(通常位于$HADOOP_HOME/share/hadoop/tools/lib/hadoop-streaming-*.jar

1.1 日志数据格式解析

典型的Nginx或Apache网站日志格式如下所示:

112.97.63.118 - - [10/May/2023:15:32:01 +0800] "GET /article/12345 HTTP/1.1" 200 4321 "https://www.example.com/" "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36"

这条日志包含了多个关键字段,我们需要提取的有价值信息包括:

字段名示例值说明
客户端IP112.97.63.118用户访问来源IP
访问时间10/May/2023:15:32:01 +0800请求发生的时间戳
请求方法GETHTTP请求方法
请求URL/article/12345用户访问的页面路径
状态码200HTTP响应状态
响应大小4321服务器返回的数据字节数
Refererhttps://www.example.com/用户来源页面
User-AgentMozilla/5.0...客户端浏览器信息

注意:实际生产环境中的日志格式可能因配置而异,在编写清洗脚本前,务必先分析你的日志样本,确认字段分隔符和顺序。

1.2 数据上传到HDFS

在开始处理前,我们需要将原始日志文件上传到HDFS分布式文件系统中。假设你的日志文件存储在本地/data/logs/目录下:

# 创建HDFS目录结构
hadoop fs -mkdir -p /user/hadoop/weblog/raw

# 上传日志文件到HDFS
hadoop fs -put /data/logs/access_2023_05_10.log /user/hadoop/weblog/raw/

# 验证文件是否上传成功
hadoop fs -ls /user/hadoop/weblog/raw/

对于生产环境,我建议设置定时任务自动上传日志。可以编写一个简单的Shell脚本,每天凌晨执行:

#!/bin/bash
# upload_logs.sh

# 获取昨天的日期
YESTERDAY=$(date -d "yesterday" +%Y_%m_%d)

# 本地日志路径
LOCAL_LOG_PATH="/data/logs/access_${YESTERDAY}.log"

# HDFS目标路径
HDFS_LOG_PATH="/user/hadoop/weblog/raw/access_${YESTERDAY}.log"

# 检查文件是否存在
if [ -f "$LOCAL_LOG_PATH" ]; then
    # 上传到HDFS
    hadoop fs -put "$LOCAL_LOG_PATH" "$HDFS_LOG_PATH"
    
    # 上传成功后可选择性删除本地文件(或移动到备份目录)
    # mv "$LOCAL_LOG_PATH" "/data/logs/backup/"
    
    echo "$(date): 成功上传 ${YESTERDAY} 日志到HDFS"
else
    echo "$(date): 警告:${LOCAL_LOG_PATH} 不存在"
fi

将上述脚本加入crontab,每天凌晨1点执行:

0 1 * * * /path/to/upload_logs.sh >> /var/log/log_upload.log 2>&1

2. Python编写MapReduce清洗脚本

现在进入核心部分:用Python编写MapReduce程序。与Java相比,Python版本的MapReduce代码更加简洁直观。

2.1 Mapper:日志解析与过滤

Mapper的任务是读取原始日志的每一行,解析出我们需要的关键字段,并进行初步过滤。创建一个名为log_mapper.py的文件:

#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""
网站日志清洗Mapper脚本
功能:解析原始日志,提取关键字段,过滤无效记录
"""

import sys
import re
from datetime import datetime

def parse_log_line(line):
    """
    解析单行日志,支持Nginx默认格式
    格式:$remote_addr - $remote_user [$time_local] "$request" $status $body_bytes_sent "$http_referer" "$http_user_agent"
    """
    # 使用正则表达式匹配日志格式
    # 这个正则匹配标准的Nginx组合日志格式
    pattern = r'(\d+\.\d+\.\d+\.\d+)\s+\S+\s+\S+\s+\[([^\]]+)\]\s+"([^"]*)"\s+(\d+)\s+(\d+)\s+"([^"]*)"\s+"([^"]*)"'
    
    match = re.match(pattern, line.strip())
    if not match:
        # 如果正则匹配失败,尝试更宽松的匹配
        return None
    
    try:
        ip, time_str, request, status, size, referer, user_agent = match.groups()
        
        # 解析HTTP请求,提取方法和URL
        # 请求格式示例:GET /article/12345 HTTP/1.1
        request_parts = request.split()
        if len(request_parts) >= 2:
            method = request_parts[0]
            url = request_parts[1]
        else:
            method = "UNKNOWN"
            url = request
        
        # 转换时间格式
        # 原始格式:10/May/2023:15:32:01 +0800
        try:
            # 移除时区信息,只保留日期时间部分
            time_str_clean = time_str.split()[0]
            dt = datetime.strptime(time_str_clean, "%d/%b/%Y:%H:%M:%S")
            formatted_time = dt.strftime("%Y-%m-%d %H:%M:%S")
        except ValueError:
            formatted_time = time_str
        
        # 过滤条件:只保留状态码为200、301、302的请求
        # 排除静态资源请求(如图片、CSS、JS等)
        if status not in ["200", "301", "302"]:
            return None
        
        # 排除常见的静态资源
        static_extensions = ['.jpg', '.jpeg', '.png', '.gif', '.css', '.js', 
                           '.ico', '.svg', '.woff', '.woff2', '.ttf', '.eot']
        if any(url.lower().endswith(ext) for ext in static_extensions):
            return None
        
        # 构建清洗后的记录
        # 格式:IP,时间,URL,状态码,响应大小,Referer,User-Agent
        cleaned_record = f"{ip}\t{formatted_time}\t{url}\t{status}\t{size}\t{referer}\t{user_agent}"
        
        return cleaned_record
        
    except Exception as e:
        # 解析出错,跳过这一行
        sys.stderr.write(f"解析错误: {str(e)}, 行: {line[:50]}...\n")
        return None

def main():
    """主函数:处理标准输入"""
    for line in sys.stdin:
        # 跳过空行
        if not line.strip():
            continue
            
        cleaned = parse_log_line(line)
        if cleaned:
            # 输出清洗后的记录
            # 使用制表符分隔,便于后续处理
            print(cleaned)

if __name__ == "__main__":
    main()

这个Mapper脚本做了几件重要的事情:

  1. 格式解析:使用正则表达式匹配日志格式,提取关键字段
  2. 数据清洗:转换时间格式,提取URL中的路径部分
  3. 数据过滤:排除错误请求和静态资源请求
  4. 格式标准化:输出统一的制表符分隔格式

2.2 Reducer:数据聚合与统计

Reducer的任务是对Mapper的输出进行汇总统计。对于日志清洗场景,Reducer可以相对简单,主要是对数据进行分组或去重。创建一个名为log_reducer.py的文件:

#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""
网站日志清洗Reducer脚本
功能:对清洗后的日志进行简单聚合,可按需求定制
"""

import sys
from collections import defaultdict

def main():
    """
    基础Reducer:按IP分组统计访问次数
    可根据实际需求修改为其他聚合逻辑
    """
    ip_count = defaultdict(int)
    
    # 读取Mapper的输出
    for line in sys.stdin:
        if not line.strip():
            continue
            
        try:
            # 解析清洗后的记录
            parts = line.strip().split('\t')
            if len(parts) < 7:
                continue
                
            ip = parts[0]
            ip_count[ip] += 1
            
        except Exception as e:
            sys.stderr.write(f"Reducer处理错误: {str(e)}\n")
            continue
    
    # 输出统计结果:IP和访问次数
    for ip, count in ip_count.items():
        print(f"{ip}\t{count}")

if __name__ == "__main__":
    main()

在实际应用中,Reducer的逻辑可以根据具体需求定制。比如,你可能需要:

  • 按小时统计访问量
  • 统计最热门的URL
  • 识别异常访问模式
  • 计算用户会话信息

2.3 本地测试脚本

在提交到Hadoop集群前,强烈建议先在本地进行测试。创建一个测试脚本test_mr_locally.py

#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""
本地测试MapReduce脚本
"""

import subprocess
import os

def test_mapreduce():
    """在本地小样本数据上测试MapReduce流程"""
    
    # 准备测试数据
    test_data = """112.97.63.118 - - [10/May/2023:15:32:01 +0800] "GET /article/12345 HTTP/1.1" 200 4321 "https://www.example.com/" "Mozilla/5.0"
114.100.100.100 - - [10/May/2023:15:33:22 +0800] "GET /static/image.jpg HTTP/1.1" 200 12345 "https://www.example.com/article/12345" "Mozilla/5.0"
115.200.200.200 - - [10/May/2023:15:34:33 +0800] "POST /login HTTP/1.1" 302 0 "https://www.example.com/" "Mozilla/5.0"
112.97.63.118 - - [10/May/2023:15:35:44 +0800] "GET /article/67890 HTTP/1.1" 200 5678 "https://www.example.com/" "Mozilla/5.0"
116.88.99.100 - - [10/May/2023:15:36:55 +0800] "GET /api/data HTTP/1.1" 404 123 "https://www.example.com/" "Mozilla/5.0"
"""
    
    print("=" * 60)
    print("开始本地MapReduce测试")
    print("=" * 60)
    
    # 步骤1:测试Mapper
    print("\n1. 测试Mapper输出:")
    print("-" * 40)
    
    # 将测试数据写入临时文件
    with open('test_input.txt', 'w', encoding='utf-8') as f:
        f.write(test_data)
    
    # 运行Mapper
    with open('test_input.txt', 'r', encoding='utf-8') as infile:
        with open('mapper_output.txt', 'w', encoding='utf-8') as outfile:
            mapper_cmd = ['python3', 'log_mapper.py']
            result = subprocess.run(mapper_cmd, stdin=infile, stdout=outfile, 
                                  stderr=subprocess.PIPE, text=True)
            
            if result.stderr:
                print("Mapper错误输出:")
                print(result.stderr)
    
    # 显示Mapper输出
    with open('mapper_output.txt', 'r', encoding='utf-8') as f:
        mapper_output = f.read()
        print(mapper_output)
    
    # 步骤2:测试Reducer
    print("\n2. 测试Reducer输出:")
    print("-" * 40)
    
    with open('mapper_output.txt', 'r', encoding='utf-8') as infile:
        with open('reducer_output.txt', 'w', encoding='utf-8') as outfile:
            reducer_cmd = ['python3', 'log_reducer.py']
            result = subprocess.run(reducer_cmd, stdin=infile, stdout=outfile,
                                  stderr=subprocess.PIPE, text=True)
            
            if result.stderr:
                print("Reducer错误输出:")
                print(result.stderr)
    
    # 显示Reducer输出
    with open('reducer_output.txt', 'r', encoding='utf-8') as f:
        reducer_output = f.read()
        print(reducer_output)
    
    # 清理临时文件
    for f in ['test_input.txt', 'mapper_output.txt', 'reducer_output.txt']:
        if os.path.exists(f):
            os.remove(f)
    
    print("\n" + "=" * 60)
    print("本地测试完成")
    print("=" * 60)

if __name__ == "__main__":
    test_mapreduce()

运行这个测试脚本,你可以看到Mapper和Reducer的处理结果,确保逻辑正确后再提交到Hadoop集群。

3. 使用Hadoop Streaming运行Python脚本

Hadoop Streaming是Hadoop提供的一个工具,允许使用任何可执行文件或脚本作为Mapper和Reducer。它的工作原理是通过标准输入输出传递数据。

3.1 准备运行环境

首先,确保你的Python脚本在所有Hadoop节点上都可执行:

# 给Python脚本添加执行权限
chmod +x log_mapper.py
chmod +x log_reducer.py

# 将脚本上传到HDFS(可选,如果使用-file参数则不需要)
hadoop fs -put log_mapper.py /user/hadoop/scripts/
hadoop fs -put log_reducer.py /user/hadoop/scripts/

3.2 运行MapReduce作业

使用以下命令提交MapReduce作业:

#!/bin/bash
# run_log_cleaning.sh

# Hadoop Streaming JAR路径(根据你的Hadoop版本调整)
STREAMING_JAR="/usr/local/hadoop/share/hadoop/tools/lib/hadoop-streaming-3.3.4.jar"

# 输入输出路径
INPUT_PATH="/user/hadoop/weblog/raw/access_2023_05_10.log"
OUTPUT_PATH="/user/hadoop/weblog/cleaned/2023_05_10"

# 如果输出目录已存在,先删除(Hadoop要求输出目录不存在)
hadoop fs -test -e $OUTPUT_PATH
if [ $? -eq 0 ]; then
    echo "输出目录已存在,删除中..."
    hadoop fs -rm -r $OUTPUT_PATH
fi

# 运行Hadoop Streaming作业
hadoop jar $STREAMING_JAR \
    -files log_mapper.py,log_reducer.py \
    -mapper "python3 log_mapper.py" \
    -reducer "python3 log_reducer.py" \
    -input $INPUT_PATH \
    -output $OUTPUT_PATH \
    -numReduceTasks 2 \
    -inputformat org.apache.hadoop.mapred.TextInputFormat \
    -outputformat org.apache.hadoop.mapred.TextOutputFormat

# 检查作业是否成功
if [ $? -eq 0 ]; then
    echo "MapReduce作业执行成功!"
    
    # 查看输出结果
    echo "清洗后的数据样例:"
    hadoop fs -cat $OUTPUT_PATH/part-* | head -20
    
    # 统计清洗后的记录数
    echo -e "\n清洗后记录统计:"
    hadoop fs -cat $OUTPUT_PATH/part-* | wc -l
else
    echo "MapReduce作业执行失败!"
    exit 1
fi

3.3 关键参数详解

Hadoop Streaming提供了丰富的配置选项,下面是一些常用参数的解释:

参数说明示例值
-files将本地文件分发到所有计算节点log_mapper.py,log_reducer.py
-mapperMapper可执行命令python3 log_mapper.py
-reducerReducer可执行命令python3 log_reducer.py
-input输入数据在HDFS中的路径/user/hadoop/weblog/raw/
-output输出数据在HDFS中的路径/user/hadoop/weblog/cleaned/
-numReduceTasksReducer任务数量2(根据数据量调整)
-inputformat输入格式类org.apache.hadoop.mapred.TextInputFormat
-outputformat输出格式类org.apache.hadoop.mapred.TextOutputFormat
-D设置Hadoop配置参数-D mapreduce.job.reduces=2

3.4 性能优化技巧

在处理大规模日志数据时,性能优化至关重要。以下是一些实用技巧:

1. 合理设置Reducer数量

# 根据数据量动态计算Reducer数量
# 经验法则:每个Reducer处理1-2GB数据
TOTAL_SIZE=$(hadoop fs -du -s $INPUT_PATH | awk '{print $1}')
REDUCER_COUNT=$((TOTAL_SIZE / 1073741824 + 1))  # 1GB = 1073741824 bytes

# 限制最大Reducer数量
MAX_REDUCERS=50
if [ $REDUCER_COUNT -gt $MAX_REDUCERS ]; then
    REDUCER_COUNT=$MAX_REDUCERS
fi

# 在Streaming命令中使用
-D mapreduce.job.reduces=$REDUCER_COUNT

2. 使用Combiner减少网络传输 Combiner是在Mapper端进行的本地Reduce操作,可以减少Mapper到Reducer的数据传输量:

#!/usr/bin/env python3
# log_combiner.py
"""Combiner脚本,在Mapper端进行本地聚合"""

import sys
from collections import defaultdict

def main():
    current_key = None
    current_count = 0
    
    for line in sys.stdin:
        key, count = line.strip().split('\t', 1)
        count = int(count)
        
        if current_key == key:
            current_count += count
        else:
            if current_key:
                print(f"{current_key}\t{current_count}")
            current_key = key
            current_count = count
    
    # 输出最后一组
    if current_key:
        print(f"{current_key}\t{current_count}")

if __name__ == "__main__":
    main()

在Streaming命令中添加Combiner:

-combiner "python3 log_combiner.py"

3. 内存和CPU调优

# 调整Mapper和Reducer的内存设置
-D mapreduce.map.memory.mb=2048 \
-D mapreduce.reduce.memory.mb=4096 \
-D mapreduce.map.java.opts=-Xmx1638m \
-D mapreduce.reduce.java.opts=-Xmx3276m \
-D mapreduce.map.cpu.vcores=2 \
-D mapreduce.reduce.cpu.vcores=2

4. 使用Hive进行高级分析

清洗后的数据已经具备了结构化特征,我们可以使用Hive进行更复杂的分析。Hive提供了类SQL的查询语言(HiveQL),让数据分析变得更加简单。

4.1 创建外部表映射清洗数据

首先,在Hive中创建一个外部表,映射到我们清洗后的数据:

-- 创建数据库(如果不存在)
CREATE DATABASE IF NOT EXISTS weblog_analysis;
USE weblog_analysis;

-- 创建外部表,映射到清洗后的数据
CREATE EXTERNAL TABLE IF NOT EXISTS cleaned_logs (
    ip STRING COMMENT '客户端IP地址',
    access_time TIMESTAMP COMMENT '访问时间',
    url STRING COMMENT '请求URL',
    status_code INT COMMENT 'HTTP状态码',
    response_size INT COMMENT '响应大小(字节)',
    referer STRING COMMENT '来源页面',
    user_agent STRING COMMENT '用户代理'
)
COMMENT '清洗后的网站日志数据'
ROW FORMAT DELIMITED
FIELDS TERMINATED BY '\t'
STORED AS TEXTFILE
LOCATION '/user/hadoop/weblog/cleaned/';

-- 查看表结构
DESCRIBE FORMATTED cleaned_logs;

-- 验证数据加载
SELECT COUNT(*) AS total_records FROM cleaned_logs;
SELECT * FROM cleaned_logs LIMIT 10;

4.2 创建分区表提升查询性能

对于按日期存储的日志数据,使用分区表可以大幅提升查询性能:

-- 创建分区表(按日期分区)
CREATE EXTERNAL TABLE IF NOT EXISTS cleaned_logs_partitioned (
    ip STRING,
    access_time TIMESTAMP,
    url STRING,
    status_code INT,
    response_size INT,
    referer STRING,
    user_agent STRING
)
PARTITIONED BY (log_date STRING)
ROW FORMAT DELIMITED
FIELDS TERMINATED BY '\t'
STORED AS TEXTFILE;

-- 添加分区(假设数据在2023-05-10目录下)
ALTER TABLE cleaned_logs_partitioned 
ADD PARTITION (log_date='2023-05-10') 
LOCATION '/user/hadoop/weblog/cleaned/2023_05_10/';

-- 查看分区信息
SHOW PARTITIONS cleaned_logs_partitioned;

4.3 关键业务指标分析

现在我们可以使用HiveQL进行各种业务分析。以下是一些常见的分析场景:

1. 基础流量统计

-- 每日PV(页面浏览量)
SELECT 
    DATE(access_time) as visit_date,
    COUNT(*) as pv,
    COUNT(DISTINCT ip) as uv,
    COUNT(DISTINCT ip) / COUNT(*) * 100 as uv_pv_ratio
FROM cleaned_logs_partitioned
WHERE log_date = '2023-05-10'
GROUP BY DATE(access_time)
ORDER BY visit_date;

-- 每小时访问趋势
SELECT 
    HOUR(access_time) as hour_of_day,
    COUNT(*) as pv,
    COUNT(DISTINCT ip) as uv
FROM cleaned_logs_partitioned
WHERE log_date = '2023-05-10'
GROUP BY HOUR(access_time)
ORDER BY hour_of_day;

2. 页面热度分析

-- 最热门的页面TOP 20
SELECT 
    url,
    COUNT(*) as visit_count,
    COUNT(DISTINCT ip) as unique_visitors,
    AVG(response_size) as avg_response_size
FROM cleaned_logs_partitioned
WHERE log_date = '2023-05-10'
  AND status_code = 200
  AND url NOT LIKE '%/static/%'
  AND url NOT LIKE '%.css'
  AND url NOT LIKE '%.js'
  AND url NOT LIKE '%.png'
  AND url NOT LIKE '%.jpg'
  AND url NOT LIKE '%.gif'
GROUP BY url
ORDER BY visit_count DESC
LIMIT 20;

3. 用户行为分析

-- 用户访问深度分析
WITH user_session AS (
    SELECT 
        ip,
        COUNT(DISTINCT url) as pages_visited,
        MIN(access_time) as first_visit,
        MAX(access_time) as last_visit,
        (UNIX_TIMESTAMP(MAX(access_time)) - UNIX_TIMESTAMP(MIN(access_time))) / 60 as session_duration_minutes
    FROM cleaned_logs_partitioned
    WHERE log_date = '2023-05-10'
      AND status_code = 200
    GROUP BY ip
)
SELECT 
    CASE 
        WHEN pages_visited = 1 THEN '跳出用户'
        WHEN pages_visited BETWEEN 2 AND 5 THEN '轻度浏览'
        WHEN pages_visited BETWEEN 6 AND 10 THEN '中度浏览'
        ELSE '深度浏览'
    END as user_type,
    COUNT(*) as user_count,
    AVG(pages_visited) as avg_pages,
    AVG(session_duration_minutes) as avg_duration_minutes
FROM user_session
GROUP BY 
    CASE 
        WHEN pages_visited = 1 THEN '跳出用户'
        WHEN pages_visited BETWEEN 2 AND 5 THEN '轻度浏览'
        WHEN pages_visited BETWEEN 6 AND 10 THEN '中度浏览'
        ELSE '深度浏览'
    END
ORDER BY user_count DESC;

4. 流量来源分析

-- Referer来源分析
SELECT 
    CASE 
        WHEN referer = '' OR referer = '-' THEN '直接访问'
        WHEN referer LIKE '%google.%' THEN 'Google'
        WHEN referer LIKE '%baidu.%' THEN '百度'
        WHEN referer LIKE '%bing.%' THEN 'Bing'
        WHEN referer LIKE '%yahoo.%' THEN 'Yahoo'
        WHEN referer LIKE '%example.com%' THEN '本站内链'
        ELSE '其他来源'
    END as traffic_source,
    COUNT(*) as visits,
    COUNT(DISTINCT ip) as unique_visitors,
    ROUND(COUNT(*) * 100.0 / SUM(COUNT(*)) OVER(), 2) as percentage
FROM cleaned_logs_partitioned
WHERE log_date = '2023-05-10'
  AND status_code = 200
GROUP BY 
    CASE 
        WHEN referer = '' OR referer = '-' THEN '直接访问'
        WHEN referer LIKE '%google.%' THEN 'Google'
        WHEN referer LIKE '%baidu.%' THEN '百度'
        WHEN referer LIKE '%bing.%' THEN 'Bing'
        WHEN referer LIKE '%yahoo.%' THEN 'Yahoo'
        WHEN referer LIKE '%example.com%' THEN '本站内链'
        ELSE '其他来源'
    END
ORDER BY visits DESC;

4.4 创建物化视图加速查询

对于频繁查询的指标,可以创建物化视图(Materialized View)来提升查询性能:

-- 创建每日汇总物化视图
CREATE MATERIALIZED VIEW IF NOT EXISTS daily_summary_mv
STORED AS ORC
TBLPROPERTIES ("transactional"="true")
AS
SELECT 
    DATE(access_time) as visit_date,
    COUNT(*) as total_pv,
    COUNT(DISTINCT ip) as total_uv,
    SUM(CASE WHEN status_code = 200 THEN 1 ELSE 0 END) as success_pv,
    SUM(CASE WHEN status_code >= 400 THEN 1 ELSE 0 END) as error_pv,
    AVG(response_size) as avg_response_size,
    COUNT(DISTINCT CASE WHEN referer NOT IN ('', '-') THEN referer END) as referer_count
FROM cleaned_logs_partitioned
WHERE access_time >= DATE_SUB(CURRENT_DATE, 30)
GROUP BY DATE(access_time);

-- 查询物化视图(性能远快于原始表查询)
SELECT * FROM daily_summary_mv ORDER BY visit_date DESC LIMIT 7;

5. 自动化流水线与监控

将整个流程自动化是生产环境的关键。下面是一个完整的自动化脚本示例:

#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""
网站日志分析自动化流水线
每天自动执行:数据上传 -> MapReduce清洗 -> Hive分析 -> 结果导出
"""

import subprocess
import sys
import os
from datetime import datetime, timedelta
import logging
from typing import Dict, Any, Optional

class LogAnalysisPipeline:
    """日志分析流水线管理类"""
    
    def __init__(self, config: Dict[str, Any]):
        """
        初始化流水线配置
        
        Args:
            config: 配置字典,包含Hadoop、Hive等连接信息
        """
        self.config = config
        self.setup_logging()
        
    def setup_logging(self):
        """配置日志"""
        log_dir = self.config.get('log_dir', '/var/log/weblog_analysis')
        os.makedirs(log_dir, exist_ok=True)
        
        log_file = os.path.join(log_dir, f"pipeline_{datetime.now().strftime('%Y%m%d')}.log")
        
        logging.basicConfig(
            level=logging.INFO,
            format='%(asctime)s - %(name)s - %(levelname)s - %(message)s',
            handlers=[
                logging.FileHandler(log_file),
                logging.StreamHandler(sys.stdout)
            ]
        )
        self.logger = logging.getLogger(__name__)
    
    def run_command(self, cmd: str, description: str) -> bool:
        """
        执行Shell命令并记录日志
        
        Args:
            cmd: 要执行的命令
            description: 命令描述
            
        Returns:
            执行是否成功
        """
        self.logger.info(f"开始执行: {description}")
        self.logger.debug(f"命令: {cmd}")
        
        try:
            result = subprocess.run(
                cmd,
                shell=True,
                check=True,
                capture_output=True,
                text=True,
                timeout=3600  # 1小时超时
            )
            
            self.logger.info(f"{description} 执行成功")
            if result.stdout:
                self.logger.debug(f"输出: {result.stdout[:500]}...")
            
            return True
            
        except subprocess.CalledProcessError as e:
            self.logger.error(f"{description} 执行失败: {e}")
            self.logger.error(f"错误输出: {e.stderr}")
            return False
        except subprocess.TimeoutExpired:
            self.logger.error(f"{description} 执行超时")
            return False
    
    def upload_logs_to_hdfs(self, date_str: str) -> bool:
        """
        上传指定日期的日志到HDFS
        
        Args:
            date_str: 日期字符串,格式 YYYY_MM_DD
            
        Returns:
            上传是否成功
        """
        local_path = os.path.join(
            self.config['local_log_dir'], 
            f"access_{date_str}.log"
        )
        hdfs_path = os.path.join(
            self.config['hdfs_raw_dir'], 
            f"access_{date_str}.log"
        )
        
        # 检查本地文件是否存在
        if not os.path.exists(local_path):
            self.logger.warning(f"本地日志文件不存在: {local_path}")
            return False
        
        cmd = f"hadoop fs -put -f {local_path} {hdfs_path}"
        return self.run_command(cmd, f"上传 {date_str} 日志到HDFS")
    
    def run_mapreduce_cleaning(self, date_str: str) -> bool:
        """
        运行MapReduce清洗作业
        
        Args:
            date_str: 日期字符串,格式 YYYY_MM_DD
            
        Returns:
            作业是否成功
        """
        input_path = os.path.join(
            self.config['hdfs_raw_dir'], 
            f"access_{date_str}.log"
        )
        output_path = os.path.join(
            self.config['hdfs_cleaned_dir'], 
            date_str.replace('_', '-')
        )
        
        # 如果输出目录存在,先删除
        check_cmd = f"hadoop fs -test -e {output_path}"
        if subprocess.run(check_cmd, shell=True).returncode == 0:
            self.run_command(
                f"hadoop fs -rm -r {output_path}",
                f"删除已存在的输出目录 {output_path}"
            )
        
        # 构建Streaming命令
        streaming_cmd = f"""
        hadoop jar {self.config['streaming_jar']} \
            -files {self.config['mapper_script']},{self.config['reducer_script']} \
            -mapper "python3 {os.path.basename(self.config['mapper_script'])}" \
            -reducer "python3 {os.path.basename(self.config['reducer_script'])}" \
            -input {input_path} \
            -output {output_path} \
            -numReduceTasks {self.config.get('num_reduce_tasks', 2)} \
            -D mapreduce.map.memory.mb={self.config.get('map_memory_mb', 2048)} \
            -D mapreduce.reduce.memory.mb={self.config.get('reduce_memory_mb', 4096)}
        """
        
        return self.run_command(streaming_cmd, f"MapReduce清洗 {date_str} 日志")
    
    def create_hive_partition(self, date_str: str) -> bool:
        """
        在Hive中创建分区
        
        Args:
            date_str: 日期字符串,格式 YYYY_MM_DD
            
        Returns:
            分区创建是否成功
        """
        hive_date = date_str.replace('_', '-')
        hdfs_path = os.path.join(
            self.config['hdfs_cleaned_dir'], 
            hive_date
        )
        
        hive_cmd = f"""
        hive -e "
        USE {self.config['hive_database']};
        ALTER TABLE {self.config['hive_table']} 
        ADD IF NOT EXISTS PARTITION (log_date='{hive_date}') 
        LOCATION '{hdfs_path}';
        "
        """
        
        return self.run_command(hive_cmd, f"创建Hive分区 {hive_date}")
    
    def run_daily_analysis(self, date_str: str) -> bool:
        """
        运行每日分析查询
        
        Args:
            date_str: 日期字符串,格式 YYYY_MM_DD
            
        Returns:
            分析是否成功
        """
        hive_date = date_str.replace('_', '-')
        
        # 创建分析结果表(如果不存在)
        create_table_sql = f"""
        CREATE TABLE IF NOT EXISTS {self.config['hive_database']}.daily_stats (
            log_date STRING,
            pv BIGINT,
            uv BIGINT,
            avg_session_duration DOUBLE,
            bounce_rate DOUBLE,
            top_pages STRING
        )
        PARTITIONED BY (stats_date STRING)
        STORED AS ORC;
        """
        
        # 插入分析结果
        analysis_sql = f"""
        INSERT INTO TABLE {self.config['hive_database']}.daily_stats
        PARTITION (stats_date='{hive_date}')
        SELECT 
            '{hive_date}' as log_date,
            COUNT(*) as pv,
            COUNT(DISTINCT ip) as uv,
            AVG(session_duration) as avg_session_duration,
            SUM(CASE WHEN page_count = 1 THEN 1 ELSE 0 END) * 100.0 / COUNT(DISTINCT ip) as bounce_rate,
            CONCAT_WS(',', COLLECT_LIST(top_page)) as top_pages
        FROM (
            SELECT 
                ip,
                COUNT(*) as page_count,
                (MAX(UNIX_TIMESTAMP(access_time)) - MIN(UNIX_TIMESTAMP(access_time))) / 60 as session_duration,
                FIRST_VALUE(url) OVER (PARTITION BY ip ORDER BY COUNT(*) DESC) as top_page
            FROM {self.config['hive_database']}.{self.config['hive_table']}
            WHERE log_date = '{hive_date}'
            GROUP BY ip
        ) t
        GROUP BY ip
        LIMIT 1;
        """
        
        hive_cmd = f'hive -e "{create_table_sql}{analysis_sql}"'
        return self.run_command(hive_cmd, f"运行 {hive_date} 分析查询")
    
    def export_to_mysql(self, date_str: str) -> bool:
        """
        将分析结果导出到MySQL
        
        Args:
            date_str: 日期字符串,格式 YYYY_MM_DD
            
        Returns:
            导出是否成功
        """
        hive_date = date_str.replace('_', '-')
        
        sqoop_cmd = f"""
        sqoop export \
            --connect jdbc:mysql://{self.config['mysql_host']}:{self.config['mysql_port']}/{self.config['mysql_database']} \
            --username {self.config['mysql_user']} \
            --password {self.config['mysql_password']} \
            --table daily_stats \
            --export-dir /user/hive/warehouse/{self.config['hive_database']}.db/daily_stats/stats_date={hive_date} \
            --input-fields-terminated-by '\\001' \
            --input-lines-terminated-by '\\n' \
            --update-mode allowinsert \
            --update-key log_date
        """
        
        return self.run_command(sqoop_cmd, f"导出 {hive_date} 数据到MySQL")
    
    def run_pipeline(self, target_date: Optional[str] = None):
        """
        运行完整的分析流水线
        
        Args:
            target_date: 目标日期,默认为昨天
        """
        if target_date is None:
            yesterday = datetime.now() - timedelta(days=1)
            target_date = yesterday.strftime('%Y_%m_%d')
        
        self.logger.info(f"开始处理 {target_date} 的日志数据")
        
        # 步骤1: 上传日志到HDFS
        if not self.upload_logs_to_hdfs(target_date):
            self.logger.error("日志上传失败,终止流水线")
            return False
        
        # 步骤2: MapReduce清洗
        if not self.run_mapreduce_cleaning(target_date):
            self.logger.error("MapReduce清洗失败,终止流水线")
            return False
        
        # 步骤3: 创建Hive分区
        if not self.create_hive_partition(target_date):
            self.logger.error("创建Hive分区失败")
            # 继续执行,分区可能已存在
        
        # 步骤4: 运行分析查询
        if not self.run_daily_analysis(target_date):
            self.logger.error("分析查询失败")
            return False
        
        # 步骤5: 导出到MySQL(可选)
        if self.config.get('export_to_mysql', False):
            if not self.export_to_mysql(target_date):
                self.logger.warning("MySQL导出失败,但分析结果已保存在Hive中")
        
        self.logger.info(f"{target_date} 日志分析流水线执行完成")
        return True

def main():
    """主函数:配置并运行流水线"""
    
    # 配置参数(实际使用时可以从配置文件或环境变量读取)
    config = {
        # 本地日志目录
        'local_log_dir': '/data/logs',
        
        # HDFS路径
        'hdfs_raw_dir': '/user/hadoop/weblog/raw',
        'hdfs_cleaned_dir': '/user/hadoop/weblog/cleaned',
        
        # Hadoop配置
        'streaming_jar': '/usr/local/hadoop/share/hadoop/tools/lib/hadoop-streaming-3.3.4.jar',
        'num_reduce_tasks': 4,
        'map_memory_mb': 2048,
        'reduce_memory_mb': 4096,
        
        # Python脚本路径
        'mapper_script': '/opt/scripts/log_mapper.py',
        'reducer_script': '/opt/scripts/log_reducer.py',
        
        # Hive配置
        'hive_database': 'weblog_analysis',
        'hive_table': 'cleaned_logs_partitioned',
        
        # MySQL配置(可选)
        'export_to_mysql': True,
        'mysql_host': 'localhost',
        'mysql_port': '3306',
        'mysql_database': 'weblog_stats',
        'mysql_user': 'analyst',
        'mysql_password': 'your_password',
        
        # 日志目录
        'log_dir': '/var/log/weblog_analysis'
    }
    
    # 创建流水线实例
    pipeline = LogAnalysisPipeline(config)
    
    # 运行流水线(处理昨天的数据)
    success = pipeline.run_pipeline()
    
    if success:
        print("流水线执行成功!")
        sys.exit(0)
    else:
        print("流水线执行失败,请检查日志")
        sys.exit(1)

if __name__ == "__main__":
    main()

5.1 监控与告警

在生产环境中,监控流水线的运行状态至关重要。我们可以添加简单的监控功能:

# monitoring.py
"""
流水线监控模块
"""

import smtplib
from email.mime.text import MIMEText
from email.mime.multipart import MIMEMultipart
import requests
import json
from datetime import datetime

class PipelineMonitor:
    """流水线监控类"""
    
    def __init__(self, config: Dict[str, Any]):
        self.config = config
    
    def send_email_alert(self, subject: str, body: str):
        """发送邮件告警"""
        if not self.config.get('email_enabled', False):
            return
        
        msg = MIMEMultipart()
        msg['From'] = self.config['email_sender']
        msg['To'] = ', '.join(self.config['email_recipients'])
        msg['Subject'] = f"[日志分析流水线] {subject}"
        
        msg.attach(MIMEText(body, 'plain'))
        
        try:
            with smtplib.SMTP(self.config['smtp_server'], self.config['smtp_port']) as server:
                server.starttls()
                server.login(self.config['email_user'], self.config['email_password'])
                server.send_message(msg)
            print(f"告警邮件已发送: {subject}")
        except Exception as e:
            print(f"发送告警邮件失败: {e}")
    
    def send_slack_alert(self, message: str):
        """发送Slack通知"""
        if not self.config.get('slack_enabled', False):
            return
        
        webhook_url = self.config['slack_webhook']
        payload = {
            "text": message,
            "username": "日志分析流水线",
            "icon_emoji": ":bar_chart:"
        }
        
        try:
            response = requests.post(
                webhook_url,
                data=json.dumps(payload),
                headers={'Content-Type': 'application/json'}
            )
            if response.status_code == 200:
                print("Slack通知已发送")
            else:
                print(f"Slack通知发送失败: {response.status_code}")
        except Exception as e:
            print(f"发送Slack通知失败: {e}")
    
    def check_hdfs_space(self) -> bool:
        """检查HDFS空间使用情况"""
        try:
            # 使用hdfs dfsadmin命令检查空间
            cmd = "hdfs dfsadmin -report | grep 'DFS Used%'"
            result = subprocess.run(cmd, shell=True, capture_output=True, text=True)
            
            if result.returncode == 0:
                used_percentage = float(result.stdout.split(':')[1].strip().replace('%', ''))
                
                if used_percentage > 90:
                    alert_msg = f"HDFS空间使用率过高: {used_percentage}%"
                    self.send_email_alert("存储空间告警", alert_msg)
                    self.send_slack_alert(alert_msg)
                    return False
            return True
            
        except Exception as e:
            print(f"检查HDFS空间失败: {e}")
            return True
    
    def check_yarn_queue(self) -> bool:
        """检查YARN队列状态"""
        try:
            # 这里可以调用YARN REST API检查队列状态
            # 简化示例:检查是否有失败的作业
            cmd = "yarn application -list -appStates FAILED | wc -l"
            result = subprocess.run(cmd, shell=True, capture_output=True, text=True)
            
            if result.returncode == 0:
                failed_jobs = int(result.stdout.strip())
                if failed_jobs > 5:
                    alert_msg = f"YARN队列中有 {failed_jobs} 个失败作业"
                    self.send_email_alert("作业失败告警", alert_msg)
                    return False
            return True
            
        except Exception as e:
            print(f"检查YARN队列失败: {e}")
            return True
    
    def generate_daily_report(self, date_str: str) -> Dict[str, Any]:
        """生成每日执行报告"""
        report = {
            "date": date_str,
            "timestamp": datetime.now().isoformat(),
            "steps": {},
            "summary": {}
        }
        
        # 这里可以添加具体的报告生成逻辑
        # 例如:检查每个步骤的日志文件,提取关键指标
        
        return report

5.2 错误处理与重试机制

在分布式环境中,网络波动、资源竞争等问题可能导致作业失败。实现重试机制可以提高系统的健壮性:

# retry_utils.py
"""
重试工具模块
"""

import time
import random
from functools import wraps
from typing import Callable, Any, Optional

class RetryError(Exception):
    """重试失败异常"""
    pass

def retry_with_backoff(
    max_retries: int = 3,
    initial_delay: float = 1.0,
    max_delay: float = 60.0,
    exponential_base: float = 2.0,
    jitter: bool = True
):
    """
    带指数退避的重试装饰器
    
    Args:
        max_retries: 最大重试次数
        initial_delay: 初始延迟(秒)
        max_delay: 最大延迟(秒)
        exponential_base: 指数基数
        jitter: 是否添加随机抖动
    """
    def decorator(func: Callable) -> Callable:
        @wraps(func)
        def wrapper(*args, **kwargs) -> Any:
            delay = initial_delay
            last_exception = None
            
            for attempt in range(max_retries + 1):
                try:
                    return func(*args, **kwargs)
                except Exception as e:
                    last_exception = e
                    
                    if attempt == max_retries:
                        break
                    
                    # 计算下一次延迟
                    if jitter:
                        # 添加随机抖动,避免惊群效应
                        delay *= exponential_base * (0.5 + random.random())
                    else:
                        delay *= exponential_base
                    
                    delay = min(delay, max_delay)
                    
                    print(f"尝试 {func.__name__} 失败 (第{attempt + 1}次): {e}")
                    print(f"等待 {delay:.2f} 秒后重试...")
                    time.sleep(delay)
            
            # 所有重试都失败
            raise RetryError(
                f"函数 {func.__name__} 在 {max_retries} 次重试后仍然失败"
            ) from last_exception
        
        return wrapper
    return decorator

# 使用示例
@retry_with_backoff(max_retries=3, initial_delay=2.0)
def run_mapreduce_job(job_config: Dict[str, Any]) -> bool:
    """运行MapReduce作业(带重试)"""
    # 实际的作业运行逻辑
    pass

@retry_with_backoff(max_retries=5, initial_delay=1.0)
def upload_to_hdfs(local_path: str, hdfs_path: str) -> bool:
    """上传文件到HDFS(带重试)"""
    # 实际上传逻辑
    pass

6. 性能调优与最佳实践

在实际生产环境中,性能调优是确保系统稳定运行的关键。以下是一些经过验证的最佳实践:

6.1 MapReduce性能优化

1. 合理设置任务数量

def calculate_optimal_tasks(input_size_gb: float) -> Dict[str, int]:
    """
    根据输入数据大小计算最优的Map和Reduce任务数量
    
    经验法则:
    - 每个Map任务处理128-256MB数据
    - 每个Reduce任务处理1-2GB数据
    """
    # 计算Map任务数量
    map_tasks = max(1, int(input_size_gb * 1024 / 256))
    map_tasks = min(map_tasks, 1000)  # 限制最大数量
    
    # 计算Reduce任务数量
    reduce_tasks = max(1, int(input_size_gb / 1))
    reduce_tasks = min(reduce_tasks, 500)  # 限制最大数量
    
    return {
        "map_tasks": map_tasks,
        "reduce_tasks": reduce_tasks
    }

2. 使用压缩减少IO

# 在Streaming命令中添加压缩选项
-D mapreduce.map.output.compress=true \
-D mapreduce.map.output.compress.codec=org.apache.hadoop.io.compress.SnappyCodec \
-D mapreduce.output.fileoutputformat.compress=true \
-D mapreduce.output.fileoutputformat.compress.codec=org.apache.hadoop.io.compress.GzipCodec \
-D mapreduce.output.fileoutputformat.compress.type=BLOCK

3. 优化Python脚本性能

# performance_optimized_mapper.py
"""
性能优化的Mapper脚本
"""

import sys
import re
from datetime import datetime

# 预编译正则表达式,避免重复编译
LOG_PATTERN = re.compile(
    r'(\d+\.\d+\.\d+\.\d+)\s+\S+\s+\S+\s+\[([^\]]+)\]\s+"([^"]*)"\s+(\d+)\s+(\d+)\s+"([^"]*)"\s+"([^"]*)"'
)

# 预定义静态资源扩展名集合
STATIC_EXTENSIONS = {
    '.jpg', '.jpeg', '.png', '.gif', '.ico', '.svg',
    '.css', '.js', '.woff', '.woff2', '.ttf', '.eot',
    '.mp4', '.mp3', '.pdf', '.zip', '.gz'
}

# 预定义有效状态码集合
VALID_STATUS_CODES = {"200", "301", "302", "304"}

def parse_log_line_fast(line: str) -> Optional[str]:
    """
    优化版本的日志解析函数
    
    使用预编译的正则表达式和集合查找
    避免重复的对象创建和函数调用
    """
    line = line.strip()
    if not line:
        return None
    
    match = LOG_PATTERN.match(line)
    if not match:
        return None
    
    try:
        ip, time_str, request, status, size, referer, user_agent = match.groups()
        
        # 快速状态码检查
        if status not in VALID_STATUS_CODES:
            return None
        
        # 快速URL检查
        # 查找URL中的文件扩展名
        question_pos = request.find('?')
        if question_pos != -1:
            url = request[:question_pos]
        else:
            url = request
        
        # 检查是否为静态资源
        dot_pos = url.rfind('.')
        if dot_pos != -1:
            extension = url[dot_pos:].lower()
            if extension in STATIC_EXTENSIONS:
                return None
        
        # 快速时间解析(只解析日期部分)
        try:
            # 假设时间格式固定,直接切片提取
            # [10/May/2023:15:32:01 +0800] -> 2023-05-10 15:32:01
            day = time_str[1:3]
            month_str = time_str[4:7]
            year = time_str[8:12]
            hour = time_str[13:15]
            minute = time_str[16:18]
            second = time_str[19:21]
            
            # 月份映射
            month_map = {
                'Jan': '01', 'Feb': '02', 'Mar': '03', 'Apr': '04',
                'May': '05', 'Jun': '06', 'Jul': '07', 'Aug': '08',
                'Sep': '09', 'Oct': '10', 'Nov': '11', 'Dec': '12'
            }
            
            month = month_map.get(month_str, '01')
            formatted_time = f"{year}-{month}-{day} {hour}:{minute}:{second}"
            
        except (IndexError, KeyError):
            formatted_time = time_str
        
        # 使用join而不是f-string,在某些Python版本中更快
        return '\t'.join([ip, formatted_time, url, status, size, referer, user_agent])
        
    except Exception:
        # 快速失败,不记录详细错误信息以提升性能
        return None

def main():
    """优化的主函数"""
    # 使用sys.stdin.buffer读取二进制数据,然后解码
    # 比直接使用sys.stdin.readline()更快
    for line_bytes in sys.stdin.buffer:
        try:
            line = line_bytes.decode('utf-8', errors='ignore')
            cleaned = parse_log_line_fast(line)
            if cleaned:
                # 直接写入二进制输出,避免编码开销
                sys.stdout.buffer.write(cleaned.encode('utf-8'))
                sys.stdout.buffer.write(b'\n')
        except UnicodeDecodeError:
            # 跳过编码错误的行
            continue

if __name__ == "__main__":
    main()

6.2 数据分区策略优化

对于时间序列数据,合理的分区策略可以大幅提升查询性能:

-- 创建按小时分区的表
CREATE EXTERNAL TABLE IF NOT EXISTS weblog_hourly_partitioned (
    ip STRING,
    url STRING,
    status_code INT,
    response_size INT,
    referer STRING,
    user_agent STRING
)
PARTITIONED BY (
    year INT,
    month INT,
    day INT,
    hour INT
)
STORED AS ORC
LOCATION '/user/hadoop/weblog/hourly_partitioned/'
TBLPROPERTIES (
    'orc.compress'='SNAPPY',
    'orc.create.index'='true',
    'transactional'='true'
);

-- 动态添加分区
SET hive.exec.dynamic.partition=true;
SET hive.exec.dynamic.partition.mode=nonstrict;

INSERT INTO TABLE weblog_hourly_partitioned PARTITION(year, month, day, hour)
SELECT 
    ip,
    url,
    status_code,
    response_size,
    referer,
    user_agent,
    YEAR(access_time) as year,
    MONTH(access_time) as month,
    DAY(access_time) as day,
    HOUR(access_time) as hour
FROM cleaned_logs
WHERE access_time >= '2023-05-01';

6.3 内存管理与调优

# memory_optimization.py
"""
内存优化工具
"""

import resource
import psutil
import gc

class MemoryMonitor:
    """内存使用监控器"""
    
    def __init__(self, process_name: str):
        self.process_name = process_name
        self.peak_memory = 0
    
    def get_memory_usage(self) -> float:
        """获取当前进程内存使用(MB)"""
        process = psutil.Process()
        memory_info = process.memory_info()
        return memory_info.rss / 1024 / 1024  # 转换为MB
    
    def check_memory_limit(self, limit_mb: float = 2048) -> bool:
        """检查内存使用是否超过限制"""
        current = self.get_memory_usage()
        self.peak_memory = max(self.peak_memory, current)
        
        if current > limit_mb:
            print(f"警告:内存使用超过限制 ({current:.2f}MB > {limit_mb}MB)")
            return False
        return True
    
    def optimize_memory(self):
        """执行内存优化"""
        # 强制垃圾回收
        collected = gc.collect()
        print(f"垃圾回收释放了 {collected} 个对象")
        
        # 获取当前内存使用
        before = self.get_memory_usage()
        
        # 尝试释放不必要的内存
        if hasattr(gc, 'get_referrers'):
            # 查找可能的内存泄漏
            pass
        
        after = self.get_memory_usage()
        print(f"内存优化:{before:.2f}MB -> {after:.2f}MB")

# 在Mapper/Reducer中使用
def process_with_memory_monitoring():
    """带内存监控的数据处理函数"""
    monitor = MemoryMonitor("log_processor")
    
    # 处理数据时定期检查内存
    batch_size = 1000
    for i, record in enumerate(data_stream):
        if i % batch_size == 0:
            if not monitor.check_memory_limit():
                # 内存不足,尝试优化
                monitor.optimize_memory()
        
        # 处理记录
        process_record(record)
    
    print(f"峰值内存使用:{monitor.peak_memory:.2f}MB")

6.4 错误处理与数据质量检查

# data_quality.py
"""
数据质量检查工具
"""

from typing import Dict, List, Tuple
import hashlib

class DataQualityChecker:
    """数据质量检查器"""
    
    def __init__(self):
        self.metrics = {
            "total_records": 0,
            "valid_records": 0,
            "invalid_records": 0,
            "field_errors": {},
            "sample_errors": []
        }
    
    def check_record(self, record: str, fields: List[str]) -> Tuple[bool, List[str]]:
        """
        检查单条记录的数据质量
        
        Returns:
            (是否有效, 错误列表)
        """
        self.metrics["total_records"] += 1
        
        errors = []
        parts = record.split('\t')
        
        # 检查字段数量
        if len(parts) != len(fields):
            errors.append(f"字段数量错误: 期望{len(fields)},实际{len(parts)}")
            self.metrics["invalid_records"] += 1
            return False, errors
        
        # 检查每个字段
        for i, (field_name, field_value) in enumerate(zip(fields, parts)):
            if not field_value.strip():
                errors.append(f"字段{i}({field_name})为空")
                self.metrics["field_errors"][field_name] = \
                    self.metrics["field_errors"].get(field_name, 0) + 1
        
        if errors:
            self.metrics["invalid_records"] += 1
            # 记录样本错误(最多记录10条)
            if len(self.metrics["sample_errors"]) < 10:
                self.metrics["sample_errors"].append({
                    "record": record[:100] + "..." if len(record) > 100 else record,
                    "errors": errors
                })
            return False, errors
        else:
            self.metrics["valid_records"] += 1
            return True, []
    
    def calculate_quality_score(self) -> float:
        """计算数据质量分数"""
        if self.metrics["total_records"] == 0:
            return 0.0
        
        valid_ratio = self.metrics["valid_records"] / self.metrics["total_records"]
        
        # 考虑字段错误率
        field_error_score = 1.0
        for field, error_count in self.metrics["field_errors"].items():
            error_ratio = error_count / self.metrics["total_records"]
            field_error_score *= (1.0 - error_ratio)
        
        return valid_ratio * field_error_score * 100
    
    def generate_report(self) -> Dict:
        """生成数据质量报告"""
        return {
            "quality_score": self.calculate_quality_score(),
            "metrics": self.metrics,
            "summary": {
                "total_processed": self.metrics["total_records"],
                "valid_percentage": (
                    self.metrics["valid_records"] / self.metrics["total_records"] * 100
                    if self.metrics["total_records"] > 0 else 0
                ),
                "invalid_percentage": (
                    self.metrics["invalid_records"] / self.metrics["total_records"] * 100
                    if self.metrics["total_records"] > 0 else 0
                )
            },
            "recommendations": self._generate_recommendations()
        }
    
    def _generate_recommendations(self) -> List[str]:
        """根据检查结果生成改进建议"""
        recommendations = []
        
        if self.metrics["valid_records"] / self.metrics["total_records"] < 0.95:
            recommendations.append("数据有效记录率低于95%,建议检查日志格式或清洗逻辑")
        
        for field, error_count in self.metrics["field_errors"].items():
            error_rate = error_count / self.metrics["total_records"]
            if error_rate > 0.1:  # 字段错误率超过10%
                recommendations.append(
                    f"字段'{field}'的错误率较高({error_rate:.1%}),建议加强数据验证"
                )
        
        return recommendations

# 在Reducer中使用数据质量检查
def process_with_quality_check():
    """带数据质量检查的数据处理"""
    checker = DataQualityChecker()
    fields = ["ip", "access_time", "url", "status_code", "response_size", "referer", "user_agent"]
    
    for line in sys.stdin:
        is_valid, errors = checker.check_record(line.strip(), fields)
        
        if is_valid:
            # 处理有效记录
            process_valid_record(line)
        else:
            # 记录无效记录
            log_invalid_record(line, errors)
    
    # 生成质量报告
    report = checker.generate_report()
    
    # 输出报告到标准错误,不影响正常输出
    sys.stderr.write(f"数据质量报告:\n")
    sys.stderr.write(f"质量分数: {report['quality_score']:.2f}\n")
    sys.stderr.write(f"有效记录: {report['summary']['valid_percentage']:.2f}%\n")
    
    if report['recommendations']:
        sys.stderr.write("改进建议:\n")
        for rec in report['recommendations']:
            sys.stderr.write(f"  - {rec}\n")

这套基于Python和Hadoop Streaming的日志分析方案,在实际项目中已经帮助多个中小型团队快速搭建起了大数据处理能力。相比传统的Java开发方式,Python方案的学习曲线更平缓,开发效率更高,特别适合资源有限但需要快速验证分析需求的团队。

我在最近的一个电商项目中实施这套方案时,原本预计需要2-3周的工作量,实际只用5天就完成了从数据清洗到可视化展示的完整流程。团队中的Python开发人员几乎不需要额外的Hadoop培训就能上手,这大大缩短了项目周期。

当然,这套方案也有其局限性。对于超大规模的数据处理(PB级别),或者需要复杂迭代计算的场景,原生的Java MapReduce或Spark可能是更好的选择。但对于日处理量在TB级别以下的网站日志分析,Python + Hadoop Streaming的组合已经足够高效和稳定。

如果你在实施过程中遇到任何问题,或者有特定的业务需求需要调整,欢迎随时交流讨论。大数据处理并不一定需要庞大的团队和复杂的技术栈,选择合适的工具和方法,小团队也能做出大成果。

更多推荐