How to parse confluent avro messages in Spark

Viewed 93

Currently I am using Abris library to de-serialize Confluent Avro messages getting from KAFKA and it works well when topic has only messages with one version of schema as soon as topic has data with different versions it start giving me malformed data found error which is obvious because while creating the config I am passing the SchemaManager.PARAM_VALUE_SCHEMA_ID=-> "latest"

But my questions is how to know the schema Id at run time basically for each record and then pass it to the Abris config here is the sample code:

Spark version: Spark 2.4.0 Scala :2.11.12 Abris:5.0.0

def getTopicSchemaMap(topicNm: String): Map[String, String] = {
  Map(
    SchemaManager.PARAM_SCHEMA_REGISTRY_TOPIC -> topicNm,
    SchemaManager.PARAM_SCHEMA_REGISTRY_URL -> schemaRegUrl,
    SchemaManager.PARAM_VALUE_SCHEMA_NAMING_STRATEGY -> "topic.name",
    SchemaManager.PARAM_VALUE_SCHEMA_ID -> "latest")

}



val kafkaDataFrameRaw = spark
    .read
    .format("kafka")
    .option("kafka.bootstrap.servers", kafkaUrl)
    .option("subscribe", topics)
    .option("maxOffsetsPerTrigger", maxOffsetsPerTrigger)
    .option("startingOffsets", "earliest")
    .option("failOnDataLoss", false)        
    .load()

val df= kafkaDataFrameRaw.select(
                  from_confluent_avro(col("value"), getTopicSchemaMap(topicNm)) as 'value, col("offset").as("offsets"))
0 Answers
Related