[FLINK-8354] Add KafkaDeserializationSchema that uses ConsumerRecord
We now directly use the ConsumerRecord from the Kafka API instead of trying to forward what we need to the deserialization schema ourselves. This makes it more future-proof, if Kafka adds new fields to the ConsumerRecord. The previously used KeyedDeserializationSchema now extends KafkaDeserializationSchema and has a default method to bridge the interface. This way existing uses of KeyedDeserializationSchema still work.
Showing
想要评论请 注册 或 登录