6 minutes
大数据入门
什么时候需要大数据?
在前面的文章中,我们一直用 Pandas 处理数据。对于百万级的数据集,Pandas 表现良好。但当数据量达到千万、亿级时,你会遇到以下问题:
- 内存不足:Pandas 需要把所有数据加载到内存中。10GB 的 CSV 文件需要至少 10GB 内存来进行处理。
- 计算时间过长:groupby 或 join 操作可能需要几分钟甚至几小时。
- 单机资源瓶颈:一台机器的 CPU 核心数有限,无法并行处理大量数据。
这就是大数据技术介入的时刻。
大数据 4V 特征
| 特征 | 英文 | 说明 |
|---|---|---|
| 体量 | Volume | 数据量巨大,TB 甚至 PB 级别 |
| 速度 | Velocity | 数据生产和处理速度快,实时流数据 |
| 多样 | Variety | 数据类型多样:结构化、半结构化、非结构化 |
| 真实 | Veracity | 数据质量和真实性不一致,需要验证 |
分布式计算基础
大数据的核心理念是分布式计算:把大任务拆成小块,在多台机器上并行执行。
MapReduce 思想
MapReduce 是分布式计算的经典范式,包含两个阶段:
- Map(映射):把数据分片,每个分片独立处理,产生键值对
- 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 服务。