【问题标题】:Apache Flink - kafka producer to sink messages to kafka topics but on different partitionsApache Flink - kafka 生产者将消息接收到 kafka 主题但在不同的分区上
【发布时间】:2020-09-09 03:41:56
【问题描述】:

现在我的 flink 代码正在处理一个文件,并在 1 个分区的 kafka 主题上下沉数据。

现在我有一个有 2 个分区的主题,我希望 flink 代码使用 DefaultPartitioner 在这 2 个分区上接收数据。

你能帮我解决这个问题吗?

这是我当前代码的sn-p代码:

DataStream<String> speStream = inputStream..map(new MapFunction<Row, String>(){....}
Properties props = Producer.getProducerConfig(propertiesFilePath);
speStream.addSink(new FlinkKafkaProducer011(kafkaTopicName, new KeyedSerializationSchemaWrapper<>(new SimpleStringSchema()), props, FlinkKafkaProducer011.Semantic.EXACTLY_ONCE));

【问题讨论】:

    标签: apache-flink


    【解决方案1】:

    通过将 flinkproducer 改为

    解决了这个问题
     speStream.addSink(new FlinkKafkaProducer011(kafkaTopicName,new SimpleStringSchema(), 
     props));
    

    我之前用过

    speStream.addSink(new FlinkKafkaProducer011(kafkaTopicName,
    new KeyedSerializationSchemaWrapper<>(new SimpleStringSchema()), props,
    FlinkKafkaProducer011.Semantic.EXACTLY_ONCE));
    

    【讨论】:

      【解决方案2】:

      Flink 版本1.11(我与Java 一起使用)中,SimpleStringSchema 需要一个包装器(即KeyedSerializationSchemaWrapper),@Ankit 也使用它,但在我建议的解决方案中删除由于相同原因,低于constructor 相关错误。

      FlinkKafkaProducer<String> producer = new FlinkKafkaProducer<String>(
                              topic_name, new KeyedSerializationSchemaWrapper<>(new SimpleStringSchema()), 
                              properties, FlinkKafkaProducer.Semantic.EXACTLY_ONCE);
      

      错误:

      The constructor FlinkKafkaProducer<String>(String, SimpleStringSchema, Properties, FlinkKafkaProducer.Semantic) is undefined
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2019-11-10
        • 2019-02-04
        • 2016-04-18
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多