【问题标题】:How to create dataframe inside ForeachWriter[Row]如何在 ForeachWriter[Row] 中创建数据框
【发布时间】:2021-09-07 23:46:45
【问题描述】:

我有一个从 Kafka 作为源读取的流式查询。我想对从流中接收到的每个批次执行一些逻辑。到目前为止,这是我的做法

val streamDF = spark
               .readStream
               ...
               .load()

//val bc = spark.sparkContext.broadcast(spark)

streamDF
     .writeStream
     .foreach( new ForeachWriter[Row] {
             def open(partitionId: Long, version: Long): Boolean = {true}
            
             def process(record: String) = {
                         val aRDD = spark.sparkContext.parallelize(Seq('a','b','C'))
                         val aDF = spark.createDataframe(aRDD)
                         //val aDF = bc.vlaue.createDataframe(aRDD)
                 
                             // do something with aDF
               }

             def close(errorOrNull: Throwable): Unit = {}
  }
).start()

我使用的是 Spark 2.3.2,所以我坚持使用 ForeachWriter(我不能使用 foreachBatch,这会让我的生活更简单)。我也知道 foreach() 对执行程序执行。 所以,记住这一点,我向所有执行者广播了 sparkSession。但这也无济于事。这是代码sn-p的注释部分。

我正在寻找一种解决方案,在 Spark 2.3.2 中将数据处理为 foreach 中的数据帧(我必须使用数据帧/数据集,因为操作非常繁重......它们也包括操作)

我发现了一个类似的问题,但没有任何回应 --> similar q

【问题讨论】:

  • 不也是一个答案。

标签: apache-spark spark-structured-streaming spark-kafka-integration


【解决方案1】:

抱歉,不是真的,但不可能在 Executor 上创建数据框。

数据框是 Spark 中的分布式集合。它们只能在 Driver 节点上或通过 Spark 应用程序中的转换(通过操作)创建。

【讨论】:

  • 是的.. 没错。是否有任何替代方法可以将每个批次作为常规数据帧处理? (不是流数据帧)
  • 指南举例说明如何做
  • 为什么不能使用foreachbatch?
  • 如问题中所述,我使用的是没有 foreachbatch 的 Spark 2.3.2(这是在 2.4 中引入的)。我正在研究一种解决方法。如果可行,我将在此处发布解决方案:) #fingerscrossed
  • 是的,但是一个奇怪的理由是要求创建一个违背 Spark 精神的 DF。可能是说明您想要实现的目标的想法。
猜你喜欢
  • 2020-06-24
  • 1970-01-01
  • 1970-01-01
  • 2018-05-31
  • 1970-01-01
  • 1970-01-01
  • 2021-10-09
  • 1970-01-01
  • 2015-06-10
相关资源
最近更新 更多