【问题标题】:Parallel method invocation in spark and using of spark session in the passed methodspark中的并行方法调用和传递方法中spark session的使用
【发布时间】:2017-02-23 04:18:02
【问题描述】:

首先让我告诉大家,我是 Spark 的新手。

我需要处理一个表中的大量记录,当它按电子邮件分组时,它大约是 100 万。我需要根据数据集针对个人电子邮件执行多个逻辑计算并根据逻辑计算更新数据库

我的代码结构大概是这样的

//Initial Data Load ...
import sparkSession.implicits._
var tableData = sparkSession.read.jdbc(<JDBC_URL>, <TABLE NAME>, connectionProperties).select("email").where(<CUSTOM CONDITION>)

//Data Frame with Records with grouping on email count greater than one
var recordsGroupedBy =tableData.groupBy("email").count().withColumnRenamed("count", "recordcount").filter("recordcount > 1 ").toDF()

//Now comes the processing after grouping against email using processDataAgainstEmail() method 

recordsGroupedBy.collect().foreach(x=>processDataAgainstEmail(x.getAs("email"),sparkSession))

在这里我看到 foreach 不是并行执行的。我需要并行调用 processDataAgainstEmail(,) 方法。 但是如果我尝试通过做并行化

您好,我可以通过调用获取列表

val emailList =dataFrameWithGroupedByMultipleRecords.select("email").rdd.map(r => r(0).asInstanceOf[String]).collect().toList

var rdd = sc.parallelize(emailList )

rdd.foreach(x => processDataAgainstEmail(x.getAs("email"),sparkSession))

这是不支持的,因为我在使用并行化时无法通过 sparkSession。

任何人都可以帮助我解决这个问题,例如 processDataAgainstEmail(,) 将执行与数据库插入和更新相关的多个操作,并且还需要执行 spark dataframe 和 spark SQL 操作?

总结我需要使用 sparksession 并行调用 processDataAgainstEmail(,)

如果无法通过 spark 会话,该方法将无法在数据库上执行任何操作。我不确定替代方法是什么,因为电子邮件的并行性对于我的方案来说是必须的。

【问题讨论】:

  • 您可能需要foreachPartition,检查此post。在 foreachPartition 中,您可以调用 processDataAgainstEmail 并将数据保存到数据库,尽管您不能使用 SparkSession。 SparkSession 仅适用于驱动程序的代码,因此您需要以不同的方式组织代码。

标签: apache-spark pyspark apache-spark-sql spark-streaming


【解决方案1】:

forEach 是按顺序对列表中的每个元素进行操作的列表方法,因此您一次只对它执行一个操作,并将其传递给processDataAgainstEmail 方法。

获得结果列表后,然后调用sc.parallelize 以从您在上一步中创建/操作的记录列表中并行创建数据框。正如我在 pySpark 中看到的那样,并行化是创建数据帧的属性,而不是任何操作的结果。

【讨论】:

  • 实际上,我想使用 processDataAgainstEmail() 方法并行处理电子邮件列表的结果,该方法将电子邮件和火花会话作为其两个参数。这可能吗?或任何替代方式?
猜你喜欢
  • 1970-01-01
  • 2023-04-07
  • 2016-04-16
  • 1970-01-01
  • 2015-10-10
  • 2017-07-28
  • 1970-01-01
  • 2015-09-25
  • 2015-06-16
相关资源
最近更新 更多