How to save streaming aggregation in Complete output mode to Parquet?

Viewed 1223

I have applied aggregation on streaming dataframe using complete mode. To save dataframe in local, I have implemented foreach sink. I am able to save dataframe in text form. But I need to save it in Parquet format.

val writerForText = new ForeachWriter[Row] {
    var fileWriter: FileWriter = _

    override def process(value: Row): Unit = {
      fileWriter.append(value.toSeq.mkString(","))
    }

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

    override def open(partitionId: Long, version: Long): Boolean = {
      FileUtils.forceMkdir(new File(s"src/test/resources/${partitionId}"))
      fileWriter = new FileWriter(new File(s"src/test/resources/${partitionId}/temp"))
      true

    }
  }

val columnName = "col1"
frame.select(count(columnName),count(columnName),min(columnName),mean(columnName),max(columnName),first(columnName), last(columnName), sum(columnName))
              .writeStream.outputMode(OutputMode.Complete()).foreach(writerForText).start()

How can I achieve that? Thanks in advance!

1 Answers

To save dataframe in local, i have implemented foreach sink. I am able to save dataframe in text form. But i need to save it in parquet format.

The default format when saving a streaming Dataset is...parquet. With that said, you don't have to use a fairly advanced foreach sink, but merely parquet.

The query could be as follows:

scala> :type in
org.apache.spark.sql.DataFrame

scala> in.isStreaming
res0: Boolean = true

in.writeStream.
  option("checkpointLocation", "/tmp/checkpoint-so").
  start("/tmp/parquets")
Related