【问题标题】:Generic protocol buffers de-serializer in Apache Flink using Java使用 Java 的 Apache Flink 中的通用协议缓冲区反序列化器
【发布时间】: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


    【解决方案1】:
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2011-05-03
    • 2017-05-15
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多