Iterate a column values in a Stream dataframe and assign each value to a common list using Scala and Spark

Viewed 194

I have the following Stream Dataframe

+------------------------------------+
|______sentence______________________|
| Representative is a scientist      |
| Norman did a good job in the exam  |
| you want to go on shopping?        |
--------------------------------------

I have list as follows

val myList

as the final output i need myList contain above three sentences in the stream dataframe

output

myList = [Representative is a scientist, Norman did a good job in the exam, you want to go on shopping? ]

I tried the following which gives stream error

val myList =   sentenceDataframe.select("sentence").rdd.map(r => r(0)).collect.toList

Error thrown with above method

org.apache.spark.sql.AnalysisException: Queries with streaming sources must be executed with writeStream.start()

Please note that above method work with normal datframe but not with stream dataframe.

Is there a way to iterate through each row of the stream dataframe and assign the row value into the common list using scala and spark ?

1 Answers

That sounds like a very weird Use-Case as the stream could theoretically never end. Are you sure you are not just looking for common spark DataFrames?

if that is not the case what you can do is use Accumulators and sparks streaming foreachBatch sink. I used a simple socket connection to demonstrate this. You can start a simple socket server under e.g. ubuntu with nc -lp 3030 and just past message there to the stream, the resulting DataFrame will have a schema of [value: String]

val acc = spark.sparkContext.collectionAccumulator[String]

val stream = spark.readStream.format("socket").option("host", "localhost").option("port", "3030").load()

val query = stream.writeStream.foreachBatch((df: DataFrame, l: Long) => {
     df.collect.foreach(v => acc.add(v(0).asInstanceOf[String]))
  }).start()

...

// For some reason you are stopping the stream here
query.stop()
val myList = acc.value

Now one question you might have is why are we using Accumulators and not just an ArrayBuffer. ArrayBuffers would work locally but on a cluster the code in foreachBatch might be executed on a total different node. That means it would not have any effect and thats also the reason Accumulators exist in the first place (see https://spark.apache.org/docs/latest/rdd-programming-guide.html#accumulators)

Related