Spark 实战:Spark SQL 处理大数据
·
Spark SQL 实战:高效处理大数据
1. Spark SQL 核心优势
- 统一数据处理:整合结构化/半结构化数据(JSON, Parquet, CSV)
- Catalyst优化器:自动优化查询计划,提升计算效率
- 内存计算:基于RDD的分布式内存模型,比Hive快$10\sim100$倍
- 多语言支持:Python/Scala/Java/R API
2. 数据处理流程
graph LR
A[数据源] --> B(Spark DataFrame)
B --> C{操作}
C --> D[转换:filter/groupBy]
C --> E[行动:count/show]
E --> F[结果输出]
3. 关键操作示例
(1) 创建SparkSession(入口)
from pyspark.sql import SparkSession
spark = SparkSession.builder \
.appName("BigDataProcessing") \
.config("spark.sql.shuffle.partitions", "200") \ # 优化并行度
.getOrCreate()
(2) 数据加载与转换
# 加载1TB Parquet数据
df = spark.read.parquet("hdfs://data/transactions.parquet")
# 数据转换(惰性执行)
transformed = df.filter(df.amount > 1000) \
.groupBy("user_id") \
.agg({"amount": "sum", "transaction_id": "count"}) \
.withColumnRenamed("sum(amount)", "total_spent")
(3) SQL查询
df.createOrReplaceTempView("transactions")
result = spark.sql("""
SELECT user_id,
PERCENTILE(amount, 0.5) AS median_amount,
COUNT(*) FILTER (WHERE status='FAILED') AS failed_count
FROM transactions
GROUP BY user_id
HAVING COUNT(*) > 100
""")
4. 性能优化技巧
-
分区策略
- 按时间分区:
df.repartition(52, "week_number") - Bucket分桶:
.bucketBy(100, "user_id").saveAsTable("bucketed_data")
- 按时间分区:
-
缓存机制
df.persist(StorageLevel.MEMORY_AND_DISK) # 复用中间结果 -
Join优化
SELECT /*+ BROADCAST(small_table) */ * FROM large_table JOIN small_table ON key
5. 实战案例:用户行为分析
# 计算用户留存率(日级)
spark.sql("""
WITH daily_users AS (
SELECT user_id, DATE_TRUNC('day', event_time) AS day
FROM events
GROUP BY 1,2
)
SELECT
first_day,
COUNT(DISTINCT d1.user_id) AS d0_users,
COUNT(DISTINCT d2.user_id) / COUNT(DISTINCT d1.user_id) AS d1_retention
FROM daily_users d1
LEFT JOIN daily_users d2
ON d1.user_id = d2.user_id
AND d2.day = DATE_ADD(d1.day, 1)
GROUP BY first_day
""").show()
6. 输出结果
result.write \
.format("parquet") \
.mode("overwrite") \
.save("hdfs://results/user_analysis")
# 或导出到外部系统
result.write.jdbc(url="jdbc:mysql://dbserver", table="results", mode="append")
关键指标:在100节点集群处理1PB数据时,Spark SQL比MapReduce减少$70%$的计算时间,主要因:
$$T_{total} = T_{io} + \frac{T_{compute}}{N_{cores}} + T_{network}$$
其中$T_{network}$通过数据本地化优化显著降低
更多推荐
所有评论(0)