0%

Spark开发集锦

Spark

Spark 2.2.x中文文档
spark两个重要概念

  • RDD
  • 共享变量

共享变量

广播变量

累加器

dev

依赖

spark core依赖

1
2
3
4
5
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-core_2.11</artifactId>
<version>2.3.1</version>
</dependency>

hdfs client

1
2
3
4
5
6
<dependency>
<groupId>org.apache.hadoop</groupId>
<artifactId>hadoop-hdfs-client</artifactId>
<version>3.2.0</version>
<scope>provided</scope>
</dependency>

导包

1
2
import org.apache.spark.SparkContext
import org.apache.spark.SparkConf

初始化

每个JVM进程中,只能有一个活跃(active)的 SparkContext 对象。如果你非要再新建一个,那首先必须将之前那个活跃的 SparkContext 对象stop()掉。

如何保证SparkContext是单例

1
2
val conf = new SparkConf().setAppName(appName).setMaster(master)
val sc = new SparkContext(conf)
1
spark-shell –help 可以查看完整的选项列表

RDD: 弹性分布式数据集

可容错、可并行操作的分布式元素集合

有两种方法可以创建 RDD 对象:由驱动程序中的集合对象通过并行化操作创建,或者从外部存储系统中数据集加载(如:共享文件系统、HDFS、HBase或者其他Hadoop支持的数据源)。

并行集合

Spark SQL partition个数与宽窄依赖

groupByKey

默认值

  1. 不指定partition大小

默认的 Spark SQL 会使用 spark.sql.shuffle.partitions 的数量来进行 aggregationjoin,默认值为 200。这会导致 partition 膨胀的问题,200个 partition 都需要执行,无论大小,尽管有些 partition 是没有数据的。

  1. 通过repartition设定大小后,再groupByKey

仍然是200

  1. Using repartition Operator With Explicit Number of Partitions

    1
    repartition(numPartitions: Int, partitionExprs: Column*): Dataset[T]

指定partitionExprs,则可以实现指定的分区数

RDD

rdd的几个基本方法:

getPartitions()

compute()

Hadoop split: Hadoop分片

hdfs dfs blockSize

hdfs-site.xml中修改dfs.blockSize, spark 2.7.x 默认值为128M

spark 3.0 读取 orcfile

1
2
3
4
5
6
7
8
9
10
11
val path="/data/sample.orc"
val orcRdd: RDD[(NullWritable, OrcStruct)] =
sc.hadoopFile(
path,
classOf[OrcInputFormat],
classOf[NullWritable],
classOf[OrcStruct],
10)
val result= orcRdd.map((line: (NullWritable, OrcStruct)) => {
line._2.getNumFields
}).collect()

原理?

TODO

toDebugString

参考文献:

  1. Number of Partitions for groupBy Aggregation