Apache spark custom log unfiltered data (LazyLogging)

Viewed 57

I'm filtering a column to comply with some validations and I can filter using Spark built-in functions, but I need to log the invalid data with a proper message (I am using LazyLogging), is there any way I can do it without using a custom UDF, so I can keep Spark optimization?

for example filtering names that are shorter then 20 characters:

df.filter(length($"name") <= lit(20))

in this scenario how can I log the names that are more than 20 characters without custom UDF?

1 Answers

In case the result of the filter operation is not too large that it does not fit into your driver, you can collect the result and print it out to your default Logger.

val logCollection = df.filter(length($"name") > lit(20)).collectAsList
logCollection.foreach(logger.info(_))

As an alternative you can create a separate stream by applying another writeStream format to write the names into a database, console etc. Just keep in mind that when you do this, you will actually create multiple streaming queries within your SparkSession which are consuming the data independently:

val originalDf = df.[...]
val logDf = df.filter(length($"name") > lit(20))

val originalQuery = originalDf.writeStream.[...].start() // keep logic as is
val logQuery = logDf.writeStream.format("console").[...].start()

spark.streams.awaitAnyTermination()
Related