【发布时间】:2017-10-25 18:20:18
【问题描述】:
在 RDD 上调用 collect() 会将整个数据集返回给驱动程序,这可能会导致内存不足,我们应该避免这种情况。
collect() 在数据帧上调用时的行为方式是否相同?select() 方法呢?
如果在数据帧上调用,它是否也与 collect() 一样工作?
【问题讨论】:
标签: dataframe apache-spark apache-spark-sql
在 RDD 上调用 collect() 会将整个数据集返回给驱动程序,这可能会导致内存不足,我们应该避免这种情况。
collect() 在数据帧上调用时的行为方式是否相同?select() 方法呢?
如果在数据帧上调用,它是否也与 collect() 一样工作?
【问题讨论】:
标签: dataframe apache-spark apache-spark-sql
- 收集(操作) - 在驱动程序中将数据集的所有元素作为数组返回。这通常在过滤器或 其他返回足够小的数据子集的操作。
select(*cols)(转换)- 投影一组表达式并返回一个新的 DataFrame。
参数:cols – 列名列表(字符串)或表达式 (柱子)。如果其中一个列名是“*”,则扩展该列 包含当前 DataFrame 中的所有列。**
df.select('*').collect() [Row(age=2, name=u'Alice'), Row(age=5, name=u'Bob')] df.select('name', 'age').collect() [Row(name=u'Alice', age=2), Row(name=u'Bob', age=5)] df.select(df.name, (df.age + 10).alias('age')).collect() [Row(name=u'Alice', age=12), Row(name=u'Bob', age=15)]
在数据帧上执行select(column-name1,column-name2,etc) 方法,返回一个新数据帧,其中仅包含在select() 函数中选择的列。
例如假设df 有几列,包括“名称”和“值”等。
df2 = df.select("name","value")
df2 将仅包含 df 的整个列中的两列(“名称”和“值”)
select 的结果 df2 将在执行程序中而不是在驱动程序中(如使用 collect() 的情况)
df.printSchema()
# root
# |-- age: long (nullable = true)
# |-- name: string (nullable = true)
# Select only the "name" column
df.select("name").show()
# +-------+
# | name|
# +-------+
# |Michael|
# | Andy|
# | Justin|
# +-------+
您可以在数据帧 (spark docs) 上运行 collect()
>>> l = [('Alice', 1)]
>>> spark.createDataFrame(l).collect()
[Row(_1=u'Alice', _2=1)]
>>> spark.createDataFrame(l, ['name', 'age']).collect()
[Row(name=u'Alice', age=1)]
要打印驱动程序上的所有元素,可以使用 collect() 方法 首先将 RDD 带到驱动节点,因此: rdd.collect().foreach(println)。 这可能会导致驱动程序耗尽 但是,因为 collect() 将整个 RDD 提取到一个 单机;如果您只需要打印 RDD 的几个元素,则 更安全的方法是使用 take(): rdd.take(100).foreach(println)。
【讨论】:
select 的一个很好的解释,但我仍然觉得我不明白collect 的作用。如果你没有collect(),你所有的例子会返回什么?
collect(),那么我猜全部代码可以正常工作,只有每个工作人员会给出所有结果,所以假设如果我们使用rdd.foreach() 代替rdd.collec().foreach(println),那么它将打印工作人员中的所有行(可能是日志文件),记住我们在笔记本中看到的是驱动程序而不是工人的输出。所以我们将无法在笔记本中看到输出
调用select 将得到lazy 评估:例如:
val df1 = df.select("col1")
val df2 = df1.filter("col1 == 3")
上述两个语句都会创建惰性路径,当您对该 df 调用操作时将执行该路径,例如 show、collect 等。
val df3 = df2.collect()
在转换结束时使用.explain 以遵循其计划
这里有更详细的信息Transformations and Actions
【讨论】:
Select 用于投影dataframe 的部分或全部字段。它不会给你一个value 作为输出,而是一个新的dataframe。它是transformation。
【讨论】:
Select 是一个转换,而不是一个动作,所以它被延迟评估(实际上不会进行计算,只是映射操作)。 Collect 是一个动作。
试试:
df.limit(20).collect()
【讨论】:
直接回答问题:
collect()在数据帧上调用时的行为方式是否相同?
是的,spark.DataFrame.collect 在功能上与spark.RDD.collect 相同。它们在这些不同的对象上具有相同的目的。
select()方法呢?
没有spark.RDD.select这样的东西,所以不能和spark.DataFrame.select一样。
如果在数据帧上调用,它的工作方式是否也与
collect()相同?
select 和 collect 之间唯一的相似之处在于它们都是 DataFrame 上的函数。它们在功能上的重叠绝对为零。
这是我自己的描述:collect 是 sc.parallelize 的反义词。 select 与任何 SQL 语句中的 SELECT 相同。
如果您仍然无法理解 collect 的实际作用(对于 RDD 或 DataFrame),那么您需要查找一些有关 spark 在幕后做什么的文章。例如:
【讨论】:
粗体字简答:
collect 主要用于序列化
(失去并行性,保留数据帧的所有其他数据特征)
例如使用 PrintWriter pw 你可以' t 直接使用df.foreach( r => pw.write(r) ),必须在foreach、df.collect.foreach(etc) 之前使用collect。
PS:“并行性损失”并不是“完全损失”,因为在序列化之后它可以再次分发给执行者。
select主要用于选择列,类似于projection in relational algebra
(仅在框架上下文中类似,因为 Spark select 不会重复数据删除)。
所以,它也是框架上下文中filter 的补充。
其他答案的评论解释:我喜欢 transformations 中的 Jeff's classification of Spark operations(如 select)和 actions(如 collect)。请记住,transforms(包括select)是lazily evaluated。
【讨论】: