【问题标题】:How to write protobuf byte array from Flink to Kafka如何将 protobuf 字节数组从 Flink 写入 Kafka
【发布时间】:2021-01-05 02:55:24
【问题描述】:

我是 Flink 的新手。我要做的就是将我的 protobuf POJO 作为字节数组放入 kafka。所以我的FlinkKafkaProducer 看起来像这样:

FlinkKafkaProducer<String> flinkKafkaProducer = createStringProducer(outputTopic, address);
        stringInputStream
                .map(//here returns byte[])
                .addSink(flinkKafkaProducer);

public static FlinkKafkaProducer<String> createStringProducer(String topic, String kafkaAddress) {
        return new FlinkKafkaProducer<>(kafkaAddress, topic, new SimpleStringSchema());
    }

现在它工作正常,但我的输出是字符串。我尝试添加TypeInformationSerializationSchema() 而不是new SimpleStringSchema() 来更改输出,但我不知道如何正确调整它。找不到教程。有人可以帮忙吗?

【问题讨论】:

    标签: java apache-kafka apache-flink


    【解决方案1】:

    所以,我终于弄清楚了如何将 protobuf 作为字节数组写入 kafka 生产者。问题在于序列化。在 POJO 的情况下,flink 使用 libery Kryo 进行自定义反序列化。编写 protobuf 的最佳方式是使用 ProtobufSerializer.class。在这个例子中,我将从 kafka String 消息中读取并写入字节数组。 Gradle 依赖项:

     compile (group: 'com.twitter', name: 'chill-protobuf', version: '0.7.6'){
            exclude group: 'com.esotericsoftware.kryo', module: 'kryo'
        }
        implementation 'com.google.protobuf:protobuf-java:3.11.0'
    

    注册:

    StreamExecutionEnvironment environment = StreamExecutionEnvironment.getExecutionEnvironment();
    environment.getConfig().registerTypeWithKryoSerializer(MyProtobuf.class, ProtobufSerializer.class);
    

    KafkaSerializerClass

    @Data
    @RequiredArgsConstructor
    public class MyProtoKafkaSerializer implements KafkaSerializationSchema<MyProto> {
        private final String topic;
        private final byte[] key;
    
        @Override
        public ProducerRecord<byte[], byte[]> serialize(MyProto element, Long timestamp) {
                    
            return new ProducerRecord<>(topic, key, element.toByteArray());
        }
    }
    

    工作

      public static FlinkKafkaProducer<MyProto> createProtoProducer(String topic, String kafkaAddress) {
            MyProtoKafkaSerializer myProtoKafkaSerializer = new MyProtoKafkaSerializer(topic);
            Properties props = new Properties();
            props.setProperty("bootstrap.servers", kafkaAddress);
            props.setProperty("group.id", consumerGroup);
            return new FlinkKafkaProducer<>(topic, myProtoKafkaSerializer, props, FlinkKafkaProducer.Semantic.AT_LEAST_ONCE);
        }
    
     public static FlinkKafkaConsumer<String> createProtoConsumerForTopic(String topic, String kafkaAddress, String kafkaGroup) {
            Properties props = new Properties();
            props.setProperty("bootstrap.servers", kafkaAddress);
            props.setProperty("group.id", kafkaGroup);
            return new FlinkKafkaConsumer<>(topic, new SimpleStringSchema(), props);
        }
    
    DataStream<String> stringInputStream = environment.addSource(flinkKafkaConsumer);
            FlinkKafkaProducer<MyProto> flinkKafkaProducer = createProtoProducer(outputTopic, address);
            stringInputStream
                    .map(hashtagMapFunction)
                    .addSink(flinkKafkaProducer);
    
            environment.execute("My test job");
    

    来源:

    1. https://ci.apache.org/projects/flink/flink-docs-release-1.10/dev/custom_serializers.html#register-a-custom-serializer-for-your-flink-program
    2. https://flink.apache.org/news/2020/04/15/flink-serialization-tuning-vol-1.html#protobuf-via-kryo

    【讨论】:

      【解决方案2】:

      在这件事上找到文档确实似乎很棘手。我假设你使用 Flink >= 1.9。在这种情况下,以下应该有效:

      private static class PojoKafkaSerializationSchema implements KafkaSerializationSchema<YourPojo> {
          @Override
          public void open(SerializationSchema.InitializationContext context) throws Exception {}
      
          @Override
          public ProducerRecord<byte[], byte[]> serialize(YourPojo element,@Nullable Long timestamp) {
              // serialize your POJO here and return a Kafka `ProducerRecord`
              return null;
          }
      }
      
      // Elsewhere: 
      PojoKafkaSerializationSchema schema = new PojoKafkaSerializationSchema();
      FlinkKafkaProducer<Integer> kafkaProducer = new FlinkKafkaProducer<>(
          "test-topic",
          schema,
          properties,
          FlinkKafkaProducer.Semantic.AT_LEAST_ONCE
      );
      

      这段代码的灵感主要来自this test case,但我没有时间实际运行它。

      【讨论】:

        猜你喜欢
        • 2020-12-03
        • 2019-08-19
        • 1970-01-01
        • 2018-05-09
        • 1970-01-01
        • 1970-01-01
        • 2021-10-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多