三亩地 三亩地SAN MU DI · CODE DIARY
ARTICLE DETAIL

日记详情

真实记录编程学习的某一天,欢迎挑你感兴趣的翻一翻。

大数据处理:Spark与分布式计算

大数据处理:Spark与分布式计算

大数据处理:Spark与分布式计算

大家好,我是欧阳瑞(Rich Own)。今天想和大家聊聊大数据处理这个重要话题。作为一个全栈开发者,处理大规模数据是现代应用的常见需求。今天就来分享一下Spark和分布式计算的实战经验。

大数据概述

大数据特点

特点说明
Volume数据量大
Velocity数据产生速度快
Variety数据类型多样
Veracity数据质量不一

处理框架对比

框架说明适用场景
Spark内存计算批处理、流处理
Hadoop磁盘计算大规模批处理
Flink流批一体实时处理
PrestoSQL查询交互式查询

Spark基础

安装Spark

# 下载Spark wget https://downloads.apache.org/spark/spark-3.5.0/spark-3.5.0-bin-hadoop3.tgz tar -xzf spark-3.5.0-bin-hadoop3.tgz cd spark-3.5.0-bin-hadoop3 # 启动Spark Shell ./bin/spark-shell

基本操作

from pyspark.sql import SparkSession # 创建SparkSession spark = SparkSession.builder \ .appName('MyApp') \ .master('local[*]') \ .getOrCreate() # 读取数据 df = spark.read.csv('data.csv', header=True, inferSchema=True) # 查看数据 df.show() df.printSchema() # 数据操作 result = df.filter(df['age'] > 30) \ .groupBy('department') \ .count() \ .orderBy('count', ascending=False) # 保存结果 result.write.csv('output.csv', header=True)

Spark SQL

创建表

# 创建临时视图 df.createOrReplaceTempView('users') # 执行SQL查询 result = spark.sql(''' SELECT department, COUNT(*) as count FROM users WHERE age > 30 GROUP BY department ORDER BY count DESC ''') result.show()

窗口函数

from pyspark.sql import Window from pyspark.sql.functions import row_number window = Window.partitionBy('department').orderBy('salary', ascending=False) df.withColumn('rank', row_number().over(window)) \ .filter('rank <= 3') \ .show()

分布式计算

RDD操作

# 创建RDD rdd = spark.sparkContext.parallelize([1, 2, 3, 4, 5]) # 转换操作 result = rdd \ .map(lambda x: x * 2) \ .filter(lambda x: x > 5) \ .reduce(lambda a, b: a + b) print(result)

分区操作

# 设置分区数 df = df.repartition(10) # 查看分区数 print(df.rdd.getNumPartitions()) # 自定义分区 def custom_partitioner(key): return hash(key) % 10 rdd = rdd.partitionBy(10, custom_partitioner)

实战案例:数据分析

class DataAnalyzer: def __init__(self, spark): self.spark = spark def analyze_sales(self, input_path): # 读取数据 df = self.spark.read.parquet(input_path) # 数据清洗 df_clean = df.filter(df['amount'].isNotNull()) \ .filter(df['amount'] > 0) # 计算指标 daily_sales = df_clean.groupBy('date') \ .sum('amount') \ .orderBy('date') category_sales = df_clean.groupBy('category') \ .sum('amount') \ .orderBy('sum(amount)', ascending=False) return daily_sales, category_sales

最佳实践

1. 性能优化

# 使用广播变量 broadcast_var = spark.sparkContext.broadcast(lookup_table) # 使用累加器 accumulator = spark.sparkContext.accumulator(0) # 持久化数据 df.cache() df.persist()

2. 资源配置

# 提交作业 spark-submit \ --master yarn \ --deploy-mode cluster \ --num-executors 10 \ --executor-memory 8g \ --driver-memory 4g \ my_script.py

总结

Spark是处理大数据的强大工具。通过Spark SQL、RDD和分布式计算,可以高效处理大规模数据。

我的鬃狮蜥Hash对大数据处理也有自己的理解——它总是能从环境中筛选出有用的信息,这也许就是自然界的"大数据分析"吧!

如果你对大数据处理有任何问题,欢迎留言交流!我是欧阳瑞,极客之路,永无止境!


技术栈:Spark · 大数据 · 分布式计算

← 返回列表