有时候需要访问同一组值,不做持久化,会重复生成,计算机代价和开销很大。持久化作用:
仅仅是标记该方法的作用是将一个RDD标记为持久化,并不是真正的持久化操作,行动操作才是真正的持久化,主要的参数是:
memory_only
将反序列化的对象存在JVM中,如果内存不足将会按照先进先出的原则,替换内容。只存入内存中。
RDD.cache() 等价于RDD.persist(memory_only),表示缓存在内存中
Memory_and_disk
先将结果存入内存中,如果内存不够,再存入磁盘中
手动将持久化的RDD对象从缓存中进行清除。
list = ["hadoop", "spark", "hive"]
rdd = sc.parallelize(list) # 生成RDD
rdd.cache() # 标记为持久化
print(rdd.count()) # 第一个行动化操作。触发从头到尾的计算,将结果存入缓存中
print(','.join(rdd.collect())) # 使用上面缓存的结果,不必再次从头到尾的进行计算,使用缓存的RDDRDD分区被保存在不同的节点上,在多个节点上同时进行计算userData和events两个表中的所有数据,都要对中间表joined表进行操作。events中的所有数据和userData中的部分数据进行操作原则是尽量使得:分区个数 = 集群中CPU核心数目。spark的部署模式
local模式(本地模式):默认为本地机器的CPU数目Standalone 模式:集群中所有的CPU数目和2之间比较取较大值yarn模式:集群中所有的CPU数目和2之间比较取较大值mesos模式:Apache,默认是8创建RDD时候指定分区个数
list = [1,2,3,4]
rdd = sc.parallelize(list,4) # 设置4个分区修改分区数目用repartition方法
data = sc.parallelize([1,2,3,4], 4) # 指定4个分区
len(data.glom().collect()) # 显示分区数目
rdd = data.repartition(2) # 重新设置分区数目为2spark自带的分区方式
# demo.py
from pyspark import SparkConf, SparkContext
def myPartitioner(key):
print("mypartitioner is running")
print("the key is %d" %key)
return key%10.
def main():
conf = SparkConf().setMaster("local").setAppName("myapp")
sc = SparkContext(conf=conf) # 生成对象,就是指挥官
data = sc.parallelize(range(10), 5) # 分成5个分区
data.map(lambda x: (x,1)) \ # 生成键值对,下图1
.partitionBy(10, myPartitioner) \ # 函数只接受键值对作为参数,将上面的data变成键值对形式传进来
.map(lambda x:x[0]) \ # 取出键值对的第一个元素,下图2
.saveAsTextFile("file:///usr/local/spark/mycode/rdd/partitioner") \ # 写入目录地址,生成10个文件
if __name__ == "__main__":
main()首先进入文件所在的目录,运行方式有两种:
python3 demo.py/usr/local/spark/bin/spark-submit demo.py