How to handle Input Schema Change dynamically in Flink Streaming Job

Viewed 51

Goal is to build a flink job which enriches the input json data from Kafka Topic with API call and send enriched information to downstream kafka in json format.

Schema of Input json data from Kafka topic changes frequently. How to handle this scenario without restarting the flink streaming job ?

Options considered.

  1. deserialize Kafka value from bytes to Jackson ObjectNode and use get methods to get the field of interest and enrich it. This approach will not require restart of flink job. But I will miss the functionalities of Table API.
  2. Build Scala case class representing the input json data , this approach will allow me to use TableAPI , however any change in schema, requires changes in Scala Case class and restarting of job.

Is there any other approach that I am missing ? and how do people handle schema changes in Flink production jobs

Thanks.

0 Answers
Related