RDD 的一些行动操作会以普通集合或者值的形式将 RDD 的部分或全部数据返回驱动器程 序中。
collect() 通常在单元测试中使用,因为此时 RDD 的整个内容不会很大,可以放在内 存中
take(n) 返回 RDD 中的 n 个元素,并且尝试只访问尽量少的分区,因此该操作会得到一个 不均衡的集合。这些操作返回元素的顺序与你预期的可能不一样。
这些操作对于单元测试和快速调试都很有用,但是在处理数据时会遇到瓶颈。
可以用 JSON 格式把数据发送到一个网络服务器上,或者把数 据存到中。都可以使用 foreach() 行动操作来对 RDD 中的每个元 素进行操作,而不需要把 RDD 发回本地。
在不同RDD类型间转换
有些函数只能用于特定类型的 RDD,比如 mean() 和 variance() 只能用在数值 RDD 上, 而 join() 只能用在键值对 RDD 上
Java
要从 T 类型的 RDD 创建出一个 DoubleRDD,我们就应当在映射操作中使用 DoubleFunction<T> 来替代 Function<T, Double>
生成JavaDoubleRDD、计算 RDD 中每个元素的平方值,这样就可以调用 DoubleRDD 独有的函数了,比如 mean() 和 variance()。
JavaDoubleRDD result = rdd.mapToDouble(
new DoubleFunction<Integer>() {
public double call(Integer x) {
return (double) x * x;
} });
System.out.println(result.mean());
持久化(缓存)
Spark RDD 是惰性求值的,而有时我们希望能多次使用同一个 RDD。如果简单 地对 RDD 调用行动操作,Spark 每次都会重算 RDD 以及它的所有依赖
迭代算法中消耗格外大,因为迭代算法常常会多次使用同一组数据
为了避免多次计算同一个 RDD,可以让 Spark 对数据进行持久化。当我们让 Spark 持久化 存储一个 RDD 时,计算出 RDD 的节点会分别保存它们所求出的分区数据。如果一个有持 久化数据的节点发生故障,Spark 会在需要用到缓存的数据时重算丢失的数据分区。如果 希望节点故障的情况不会拖累我们的执行速度,也可以把数据备份到多个节点上。
Java 中,默认情况下 persist() 会把数据以序列化的形式缓存在 JVM 的堆空间中
result.persist(StorageLevel.DISK_ONLY)
persist() 调 用本身不会触发强制求值
如果要缓存的数据太多,内存中放不下,Spark 会自动利用最近最少使用(LRU)的缓存 策略把最老的分区从内存中移除
对于仅把数据存放在内存中的缓存级别,下一次要用到 已经被移除的分区时,这些分区就需要重新计算
但是对于使用内存与磁盘的缓存级别的 分区来说,被移除的分区都会写入磁盘
RDD 还有一个方法叫作 unpersist(),调用该方法可以手动把持久化的 RDD 从缓 存中移除
本文来自电脑杂谈,转载请注明本文网址:
http://www.pc-fly.com/a/sanxing/article-57960-5.html
一个连海洋法公约都没签署的国家却天天叫嚷遵守国际法
我国其实可以在南海举行实弹演习的
当然凭他的数学水平也确实管不了