【问题标题】:Spark - How to use SparkContext within classes?Spark - 如何在类中使用 SparkContext?
【发布时间】:2015-10-18 07:56:25
【问题描述】:

我正在 Spark 中构建一个应用程序,并希望在我的类的方法中使用 SparkContext 和/或 SQLContext,主要是从文件或 SQL 查询中提取/生成数据集。

例如,我想创建一个 T2P 对象,其中包含收集数据的方法(在这种情况下需要访问 SparkContext):

class T2P (mid: Int, sc: SparkContext, sqlContext: SQLContext) extends Serializable {

  def getImps(): DataFrame = {
      val imps = sc.textFile("file.txt").map(line => line.split("\t")).map(d => Data(d(0).toInt, d(1), d(2), d(3))).toDF()
      return imps
   }

  def getX(): DataFrame = {
      val x = sqlContext.sql("SELECT a,b,c FROM table")
      return x
  }
}

//creating the T2P object
class App {
    val conf = new SparkConf().setAppName("T2P App").setMaster("local[2]")
    val sc = new SparkContext(conf)
    val sqlContext = new SQLContext(sc)
    val t2p = new T2P(0, sc, sqlContext);
}

将 SparkContext 作为参数传递给 T2P 类不起作用,因为 SparkContext 不可序列化(创建 T2P 对象时出现task not serializable 错误)。在我的类中使用 SparkContext/SQLContext 的最佳方式是什么?或者这可能是在 Spark 中设计数据拉取类型流程的错误方法?

更新 从这篇文章中的 cmets 意识到 SparkContext 不是问题,而是我在“map”函数中使用了一个方法,导致 Spark 尝试序列化整个类。这将导致错误,因为 SparkContext 不可序列化。

def startMetricTo(userData: ((Int, String), List[(Int, String)]), startMetric: String) : T2PUser = {
  //do something 
}

def buildUserRollup() = {
  this.userRollup = this.userSorted.map(line=>startMetricTo(line, this.startMetric))
}

这会导致“任务不可序列化”异常。

【问题讨论】:

  • 为什么它不起作用?你能提供例子吗?除了 rdd 转换闭包中的那些部分之外,您的所有代码都在主节点上运行。在主节点上,您不需要序列化。在内部 rdd 转换中,您不需要 sparkContext。你能阐明你的目标吗?
  • 已更新。这有帮助吗?
  • 好吧,为什么需要序列化?重新考虑一下:所有分布式工作都是在像map\filter 这样的rdd 转换中完成的。您只是不需要对存储库帮助程序进行序列化。
  • 你能显示你创建它的地方吗?看起来您在 rdd 转换中执行此操作。另外,能不能把class换成object
  • 我正在使用同样的方法,它工作正常。将 SparkContext 和 SqlContext 作为参数传递不应在类实例化时生成此错误(至少在 1.4 中)。我可以看到有几件事会导致问题。一个是,如果您的类“数据”由于某种原因不可序列化,您可能会遇到该异常。另一种可能是,如果您的示例中未列出的方法正在执行更复杂的操作(例如,闭包可能会破坏任务序列化),这也可能导致问题。

标签: java scala apache-spark


【解决方案1】:

我通过创建一个单独的 MetricCalc 对象来存储我的 startMetricTo() 方法解决了这个问题(在评论者和其他 StackOverflow 用户的帮助下)。然后我更改了 buildUserRollup() 方法以使用这个新的 startMetricTo()。这允许整个MetricCalc 对象被序列化而不会出现问题。

//newly created object
object MetricCalc {
    def startMetricTo(userData: ((Int, String), List[(Int, String)]), startMetric: String) : T2PUser = {
    //do something
  }
}

//using function in T2P
def buildUserRollup(startMetric: String) = {
   this.userRollup = this.userSorted.map(line=>MetricCalc.startMetricTo(line, startMetric))
}

【讨论】:

    【解决方案2】:

    我尝试了几个选项,这最终对我有用..

       object SomeName extends App {
    
       val conf = new SparkConf()...
       val sc = new SparkContext(conf)
    
       implicit val sqlC = SQLContext.getOrCreate(sc)
       getDF1(sqlC)
    
       def getDF1(sqlCo: SQLContext): Unit = {
        val query1 =  SomeQuery here  
        val df1 = sqlCo.read.format("jdbc").options(Map("url" -> dbUrl,"dbtable" -> query1)).load.cache()
    
         //iterate through df1 and retrieve the 2nd DataFrame based on some values in the Row of the first DataFrame
    
        df1.foreach(x => {
         getDF2(x.getString(0), x.getDecimal(1).toString, x.getDecimal(3).doubleValue) (sqlCo)
       })     
      }
    
      def getDF2(a: String, b: String, c: Double)(implicit sqlCont: SQLContext) :  Unit = {
         val query2 = Somequery
    
         val sqlcc = SQLContext.getOrCreate(sc)
        //val sqlcc = sqlCont //Did not work for me. Also, omitting (implicit sqlCont: SQLContext) altogether did not work
         val df2 = sqlcc.read.format("jdbc").options(Map("url" -> dbURL, "dbtable" -> query2)).load().cache()
           .
           .
           .
         }
      }
    

    注意:在上面的代码中,如果我在 getDF2 方法签名中省略了 (implicit sqlCont: SQLContext) 参数,它将不起作用。我尝试了其他几种将 sqlContext 从一种方法传递到另一种方法的选项,它总是给我 NullPointerException 或 Task not serializable Excpetion。好消息是它最终以这种方式工作,我可以从 DataFrame1 的一行中检索参数并在加载 DataFrame 2 时使用这些值。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2015-10-02
      • 2018-08-12
      • 2017-02-13
      • 2019-02-23
      相关资源
      最近更新 更多