【发布时间】:2021-06-29 18:22:35
【问题描述】:
场景:Apache Flink、Kafka、Protocol buffers 数据消费者。
数据源是协议缓冲区格式的Kafka主题(多个主题:主题#1、主题#3、主题#3)。 消费者是 Apache Flink 消费者。每个主题都有一个唯一的 protobuf 定义。
List<String> topicList = Arrays.asList("topic#1,topic#2,topic#3".split(","));
inputStream = env.addSource(new FlinkKafkaConsumer[ProtobufDeserializationSchema](topicList, new ProtobufDeserializationSchema(), properties));
我正在尝试在 Apache Flink 中开发通用数据摄取作业,以将 Kafka 中的数据摄取到数据库中。
如何为 Apache Flink 实现一个通用的 protobuf 反序列化器?我正在寻找将 Kafka 主题链接到 protobuf 定义以进行反序列化的实现。
最初的做法是将字节数组带入 Flink 数据流中,然后根据 Kafka 主题名称确定 protobuf 定义,对 map 函数中的消息进行反序列化。我怎样才能以通用方式做到这一点?
【问题讨论】:
标签: java apache-kafka protocol-buffers apache-flink