什么时候需要大数据?

在前面的文章中,我们一直用 Pandas 处理数据。对于百万级的数据集,Pandas 表现良好。但当数据量达到千万、亿级时,你会遇到以下问题:

  • 内存不足:Pandas 需要把所有数据加载到内存中。10GB 的 CSV 文件需要至少 10GB 内存来进行处理。
  • 计算时间过长:groupby 或 join 操作可能需要几分钟甚至几小时。
  • 单机资源瓶颈:一台机器的 CPU 核心数有限,无法并行处理大量数据。

这就是大数据技术介入的时刻。

大数据 4V 特征

特征 英文 说明
体量 Volume 数据量巨大,TB 甚至 PB 级别
速度 Velocity 数据生产和处理速度快,实时流数据
多样 Variety 数据类型多样:结构化、半结构化、非结构化
真实 Veracity 数据质量和真实性不一致,需要验证

分布式计算基础

大数据的核心理念是分布式计算:把大任务拆成小块,在多台机器上并行执行。

MapReduce 思想

MapReduce 是分布式计算的经典范式,包含两个阶段:

  1. Map(映射):把数据分片,每个分片独立处理,产生键值对
  2. Reduce(归约):将相同键的结果聚合起来
输入: [A, B, C, D, E, F]
        |   |   |   |   |   |
      Map  Map Map Map Map Map  (并行处理)
        |   |   |   |   |   |
        \   /   \   /   \   /
        Reduce Reduce Reduce  (聚合结果)
           |       |       |
        [结果1]  [结果2]  [结果3]

Apache Spark 概述

Spark 是目前最流行的大数据处理框架。它比传统的 Hadoop MapReduce 快 10-100 倍,因为它利用内存计算,避免了频繁的磁盘读写。

安装 PySpark

pip install pyspark

SparkSession

所有 Spark 程序的入口点。

from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .appName("BigDataAnalysis") \
    .config("spark.sql.adaptive.enabled", "true") \
    .config("spark.sql.adaptive.coalescePartitions.enabled", "true") \
    .config("spark.driver.memory", "4g") \
    .getOrCreate()

print(f'Spark 版本: {spark.version}')

RDD 与 DataFrame

Spark 提供两种核心抽象。

RDD(Resilient Distributed Dataset)

RDD 是最底层的抽象,代表一个不可变的分布式数据集合。你可以显式控制每个操作。

# 从列表创建 RDD
data = [1, 2, 3, 4, 5, 6, 7, 8, 9, 10]
rdd = spark.sparkContext.parallelize(data, numSlices=4)

# RDD 操作
rdd_map = rdd.map(lambda x: x * 2)
rdd_filter = rdd.filter(lambda x: x > 5)
rdd_flat = rdd.flatMap(lambda x: [x, x * 10])

print(f'原始 RDD: {rdd.collect()}')
print(f'Map 结果: {rdd_map.collect()}')
print(f'Filter 结果: {rdd_filter.collect()}')
print(f'FlatMap 结果: {rdd_flat.collect()}')

# 聚合操作
sum_result = rdd.reduce(lambda a, b: a + b)
count = rdd.count()
print(f'求和: {sum_result}, 计数: {count}')

DataFrame(推荐)

DataFrame 建立在 RDD 之上,提供了类似 Pandas 的 API,但底层是分布式的。

# 从 Pandas DataFrame 创建
import pandas as pd

pdf = pd.DataFrame({
    'name': ['Alice', 'Bob', 'Charlie', 'David'],
    'age': [25, 30, 35, 28],
    'city': ['北京', '上海', '广州', '深圳'],
    'salary': [15000, 20000, 25000, 18000]
})

df = spark.createDataFrame(pdf)
df.show()

# 从 CSV 文件读取
sales_df = spark.read.csv(
    'huge_sales_data.csv',
    header=True,
    inferSchema=True
)

# 从 Parquet 读取(推荐,速度更快)
sales_df = spark.read.parquet('sales_data.parquet')

DataFrame 核心操作

数据查看

# 基本查看
sales_df.printSchema()
sales_df.show(10, truncate=False)
sales_df.describe().show()

# 统计信息
total_rows = sales_df.count()
print(f'总行数: {total_rows}')

# 列选择
sales_df.select('order_id', 'customer_id', 'amount').show(5)

过滤

# 过滤条件
high_value_orders = sales_df.filter(sales_df['amount'] > 1000)
high_value_orders.show(5)

# 多条件过滤
filtered = sales_df.filter(
    (sales_df['amount'] > 500) & 
    (sales_df['status'] == 'completed')
)
print(f'符合条件的订单数: {filtered.count()}')

# ISIN 操作
cities = ['北京', '上海', '深圳']
city_orders = sales_df.filter(sales_df['city'].isin(cities))

分组聚合

from pyspark.sql import functions as F

# 按类别聚合
category_stats = sales_df.groupBy('category').agg(
    F.count('order_id').alias('order_count'),
    F.sum('amount').alias('total_revenue'),
    F.avg('amount').alias('avg_order_value'),
    F.max('amount').alias('max_order'),
    F.min('amount').alias('min_order'),
    F.stddev('amount').alias('std_amount')
).orderBy(F.desc('total_revenue'))

category_stats.show(10)

# 多维度聚合
daily_stats = sales_df.groupBy('order_date', 'category').agg(
    F.sum('amount').alias('daily_revenue'),
    F.countDistinct('customer_id').alias('unique_customers')
)

# 窗口函数
from pyspark.sql.window import Window

window_spec = Window.partitionBy('category').orderBy(F.desc('amount'))
sales_df.withColumn('rank_in_category', 
    F.rank().over(window_spec)
).filter(F.col('rank_in_category') <= 3).show()

关联操作

# 读取客户表和订单表
customers = spark.read.parquet('customers.parquet')
orders = spark.read.parquet('orders.parquet')

# Inner Join
customer_orders = customers.join(
    orders,
    customers['customer_id'] == orders['customer_id'],
    how='inner'
)

# Left Join
all_customers = customers.join(
    orders,
    customers['customer_id'] == orders['customer_id'],
    how='left'
)

# 关联后聚合
customer_ltv = customer_orders.groupBy(
    customers['customer_id'], 
    customers['name'], 
    customers['city']
).agg(
    F.sum(orders['amount']).alias('total_spent'),
    F.count(orders['order_id']).alias('order_count'),
    F.max(orders['order_date']).alias('last_order_date')
).orderBy(F.desc('total_spent'))

customer_ltv.show(10)

UDF(用户自定义函数)

from pyspark.sql.types import StringType

# 注册 UDF
def price_level(amount):
    if amount > 1000:
        return '高'
    elif amount > 500:
        return '中'
    else:
        return '低'

price_level_udf = F.udf(price_level, StringType())

# 使用 UDF
sales_with_level = sales_df.withColumn(
    'price_level', 
    price_level_udf(sales_df['amount'])
)
sales_with_level.groupBy('price_level').count().show()

使用 Pandas UDF(向量化 UDF)

Pandas UDF 比普通 UDF 快很多,因为它利用 Apache Arrow 进行零拷贝数据传输。

import pandas as pd
from pyspark.sql.types import DoubleType
from pyspark.sql.functions import pandas_udf

@pandas_udf(DoubleType())
def discount_price(price, quantity):
    """批量计算折扣价格"""
    total = price * quantity
    # 满 1000 打 9 折
    discount = np.where(total > 1000, total * 0.9, total)
    return pd.Series(discount)

sales_with_discount = sales_df.withColumn(
    'discounted_amount',
    discount_price(sales_df['price'], sales_df['quantity'])
)

Spark SQL

如果你熟悉 SQL,完全可以用 SQL 来操作 DataFrame。

# 注册临时视图
sales_df.createOrReplaceTempView('sales')
customers.createOrReplaceTempView('customers')

# 执行 SQL 查询
result = spark.sql("""
    SELECT 
        c.city,
        c.segment,
        COUNT(DISTINCT s.order_id) AS total_orders,
        SUM(s.amount) AS total_revenue,
        AVG(s.amount) AS avg_order_value,
        COUNT(DISTINCT c.customer_id) AS customer_count,
        SUM(s.amount) / COUNT(DISTINCT c.customer_id) AS arpu
    FROM sales s
    JOIN customers c ON s.customer_id = c.customer_id
    WHERE s.order_date >= '2025-01-01'
    GROUP BY c.city, c.segment
    HAVING total_orders > 100
    ORDER BY total_revenue DESC
""")

result.show(20)

# 将结果保存
result.write.mode('overwrite').parquet('output/city_segment_analysis.parquet')

Parquet 格式

Parquet 是 Hadoop/Spark 生态中最常用的列式存储格式。

Parquet 的优势

特性 CSV Parquet
存储方式 行式 列式
压缩率 一般 极高(通常 5-10 倍)
读取速度 需要读整行 只读取需要的列
Schema 无(需要推断) 自带 Schema
类型安全 弱(全部是字符串) 强(原生类型)
# 写入 Parquet
sales_df.write \
    .mode('overwrite') \
    .partitionBy('order_year', 'order_month') \
    .parquet('data_lake/sales/')

# 读取 Parquet
spark.read.parquet('data_lake/sales/').show()

# Parquet 的谓词下推(Predicate Pushdown)
# Spark 会自动利用 Parquet 的统计信息来跳过不相关的数据块
fast_query = spark.read.parquet('data_lake/sales/') \
    .filter(F.col('amount') > 10000)  # 只读取包含大额订单的数据块
fast_query.explain()

PySpark vs Pandas:何时用谁

# 对比:同样的操作在 Pandas 和 PySpark 中的写法

# Pandas 版本
import pandas as pd
pdf = pd.read_csv('sales.csv')
result_pd = pdf.groupby('category').agg(
    revenue=('amount', 'sum'),
    orders=('order_id', 'count'),
    avg_value=('amount', 'mean')
).reset_index()

# PySpark 版本
from pyspark.sql import functions as F
sdf = spark.read.csv('sales.csv', header=True, inferSchema=True)
result_spark = sdf.groupBy('category').agg(
    F.sum('amount').alias('revenue'),
    F.count('order_id').alias('orders'),
    F.avg('amount').alias('avg_value')
)

选择指南

场景 推荐 理由
数据量 < 1GB Pandas 简单、生态丰富
数据量 1-100GB Pandas + 分块 无需搭建集群
数据量 > 100GB PySpark 分布式处理
复杂统计分析 Pandas statsmodels, scikit-learn 支持更好
简单的 ETL 和大规模聚合 PySpark 更快、更省内存
ML 特征工程 结合使用 用 Spark 做大规模预聚合,Pandas 做精细处理

简单的分布式分析示例

下面用 PySpark 做一个完整的数据分析。

from pyspark.sql import SparkSession, functions as F
from pyspark.sql.window import Window
import matplotlib.pyplot as plt

spark = SparkSession.builder \
    .appName('EcommerceAnalysis') \
    .config('spark.sql.adaptive.enabled', 'true') \
    .getOrCreate()

# 1. 读取数据
orders = spark.read.csv('orders_2024.csv', header=True, inferSchema=True)
orders.cache()  # 缓存到内存,加速后续操作

print(f'订单总数: {orders.count():,}')
orders.printSchema()

# 2. 数据清洗
clean_orders = orders.filter(
    (F.col('amount') > 0) & 
    (F.col('status') != 'cancelled')
)

# 3. 每日销售额
daily_sales = clean_orders.groupBy('order_date').agg(
    F.sum('amount').alias('revenue'),
    F.count('order_id').alias('orders'),
    F.countDistinct('customer_id').alias('customers')
).orderBy('order_date')

# 4. 产品品类分析
category_analysis = clean_orders.groupBy('category').agg(
    F.sum('amount').alias('revenue'),
    F.count('order_id').alias('order_count'),
    F.countDistinct('customer_id').alias('unique_customers'),
    F.avg('amount').alias('avg_order_value')
).withColumn(
    'revenue_rank',
    F.row_number().over(Window.orderBy(F.desc('revenue')))
).orderBy(F.desc('revenue'))

category_analysis.show()

# 5. 客户分组
customer_value = clean_orders.groupBy('customer_id').agg(
    F.sum('amount').alias('total_spent'),
    F.count('order_id').alias('order_count'),
    F.datediff(F.lit('2025-01-01'), F.max('order_date')).alias('recency')
)

# RFM 打分
customer_value = customer_value.withColumn(
    'r_score',
    F.ntile(4).over(Window.orderBy(F.col('recency')))
).withColumn(
    'f_score',
    F.ntile(4).over(Window.orderBy(F.desc('order_count')))
).withColumn(
    'm_score',
    F.ntile(4).over(Window.orderBy(F.desc('total_spent')))
)

customer_value.select('customer_id', 'r_score', 'f_score', 'm_score').show(10)

# 6. 结果保存为 Parquet
daily_sales.write.mode('overwrite').parquet('output/daily_sales_agg/')
category_analysis.write.mode('overwrite').parquet('output/category_analysis/')

# 7. 收集少量结果到驱动节点做可视化
daily_sales_local = daily_sales.limit(365).toPandas()
daily_sales_local['order_date'] = pd.to_datetime(daily_sales_local['order_date'])
daily_sales_local.set_index('order_date')['revenue'].plot(
    figsize=(12, 5), title='2024 每日销售额'
)
plt.show()

# 释放缓存
orders.unpersist()

本地运行 Spark

你不需要搭建集群就可以开始学习 Spark。

安装和配置

# 1. 安装 Java(Spark 依赖 JVM)
# Windows: 下载并安装 JDK 11+

# 2. 安装 PySpark
pip install pyspark

# 3. 设置环境变量(Windows)
$env:JAVA_HOME = "C:\Program Files\Java\jdk-11"
$env:SPARK_HOME = "C:\Program Files\spark"

本地模式

# local[*] 表示使用所有可用的 CPU 核心
spark = SparkSession.builder \
    .master('local[*]') \
    .appName('LocalSpark') \
    .getOrCreate()

理解分区

# 查看分区数
df = spark.read.csv('large_file.csv', header=True)
print(f'分区数: {df.rdd.getNumPartitions()}')

# 调整分区
df_repartitioned = df.repartition(8)  # 增加分区
df_coalesced = df.coalesce(2)  # 减少分区(避免 shuffle)

扩展阅读:云上的 Spark

在实际生产环境中,你很少直接在单机上跑 Spark。常见的部署方式有:

平台 服务名称 特点
AWS EMR (Elastic MapReduce) 与 S3 深度集成
GCP Dataproc 按秒计费,与 BigQuery 集成
Azure HDInsight / Synapse 与 Azure Data Lake 集成
Databricks 统一分析平台 最友好的 Spark 体验
阿里云 E-MapReduce 国内主流选择

注意事项

惰性求值(Lazy Evaluation)

Spark 的 transformations 是惰性的,只有遇到 action 才会真正执行。

# 下面这些不会立即执行
df1 = spark.read.csv('data.csv', header=True)
df2 = df1.filter(F.col('amount') > 100)
df3 = df2.groupBy('category').count()

# 直到调用 action 才会执行
df3.show()  # 触发真正的计算

尽量避免的陷阱

# ❌ 不要在 UDF 中引用非序列化的外部对象
class MyCalculator:
    def compute(self, x):
        return x * 2

calc = MyCalculator()

# 这会序列化失败
bad_udf = F.udf(lambda x: calc.compute(x))

# ✅ 使用静态函数
@F.udf
def good_udf(x):
    return x * 2

# ❌ 不要在大数据集上使用 .toPandas()
# 这会把所有数据拉到驱动节点,导致 OOM
# bad = large_df.toPandas()  # 危险!

# ✅ 先聚合再收集
small_result = large_df.groupBy('category').count().toPandas()

# ❌ 不要在循环中使用 DataFrame
# 每次迭代都会重新执行整个 DAG
for i in range(10):
    df.filter(F.col('amount') > i).count()  # 重复执行

# ✅ 使用 cache() 或 checkpoint()
df_cached = df.filter(F.col('status') == 'completed').cache()
df_cached.count()  # 第一次触发计算
df_cached.filter(F.col('amount') > 100).count()  # 第二次直接读缓存

小结

大数据技术让分析海量数据集成为可能。本篇文章覆盖了:

  • 大数据 4V 特征和分布式计算的基本思想
  • Spark 的核心概念:SparkSession、RDD、DataFrame
  • PySpark DataFrame 的常用操作:过滤、聚合、关联、UDF
  • Spark SQL 在数据分析中的应用
  • Parquet 列式存储及其优势
  • Pandas 与 PySpark 的选择策略
  • 本地运行 Spark 和云上部署方案

重要提醒:Spark 虽然强大,但并非所有问题都需要它。先问自己数据量是否真的需要分布式计算,避免过度设计。

在下一篇文章中,我们将探讨如何把分析结果产品化,构建数据应用和 API 服务。