【发布时间】:2017-07-03 14:19:55
【问题描述】:
这可能是一个愚蠢的问题。但是,我想知道我是否有这样的东西 - rdd.mapPartitions(func)。 func 中的逻辑应该是线程安全的吗?
谢谢
【问题讨论】:
-
由于 RDD 操作总是不可变的,线程安全问题对于转换 RDD 的底层函数并不重要。
标签: java apache-spark thread-safety rdd
这可能是一个愚蠢的问题。但是,我想知道我是否有这样的东西 - rdd.mapPartitions(func)。 func 中的逻辑应该是线程安全的吗?
谢谢
【问题讨论】:
标签: java apache-spark thread-safety rdd
传递给mapPartitions 的任何函数(或任何其他操作或转换)都必须是线程安全的。 JVM 上的 Spark(对于来宾语言不一定如此)使用执行线程并且不保证各个任务之间的任何隔离。
当您使用未在函数中初始化但通过闭包传递的资源时,这一点尤其重要,例如在主函数中初始化但在函数中引用的对象。
不用说,除非明确允许,否则不应修改任何参数。
【讨论】:
简短的回答是否定的,它不必是线程安全的。
这样做的原因是 spark 在分区之间划分数据。然后它为每个分区创建一个任务,您编写的函数将在该特定分区内作为单线程操作运行(即没有其他线程会访问相同的数据)。
也就是说,您必须确保不会通过访问不是 RDD 数据的资源来手动创建线程“不安全”。例如,如果您创建一个静态对象并访问它,它可能会导致问题,因为多个任务可能在同一个执行程序 (JVM) 中运行并访问它。也就是说,除非您确切知道自己在做什么,否则您不应该一开始就做那样的事情......
【讨论】:
spark.task.cpus 大于1 怎么办?每个任务可以有多个线程,是否存在竞争条件问题?
当你执行“rdd.mapPartitions(func)”时,func 可能实际上在不同的 jvm 中执行!!!线程在 JVM 中没有意义。
如果您在本地模式下运行,并使用全局状态或线程不安全函数,作业可能会按预期工作,但行为未定义或不支持。
【讨论】: