【问题标题】:Send key in Flink Kafka Producer在 Flink Kafka Producer 中发送密钥
【发布时间】:2021-03-22 16:15:03
【问题描述】:

我是 Flink Stream 处理的新手,需要一些关于 Flink Kafka 生产者的帮助,因为经过一段时间的搜索后找不到与之相关的太多东西。我目前正在从 Kafka 主题读取流,然后在执行一些计算后,我想将其写入 Kafka 中的新单独主题。但我面临的问题是我无法将密钥发送到 Kafka 主题。我正在使用 Flink Kafka 连接器,它为我提供了 FlinkKafkaConsumer 和 FlinkKafkaProducer。以下是我的代码的更详细信息,我可以在我的代码中更改它可以工作的内容,目前在 Kafka 上,我正在生成我的消息在 Key 中使用 null ,因为值是我需要的:

Properties consumerProperties = new Properties();
    
    consumerProperties.setProperty("bootstrap.servers", serverURL);
    consumerProperties.setProperty("group.id", groupID);
    StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
    FlinkKafkaConsumer<String> kafkaConsumer = new FlinkKafkaConsumer<>(consumerTopicName,
            new SimpleStringSchema(), consumerProperties);

    kafkaConsumer.setStartFromEarliest();
    DataStream<String> kafkaConsumerStream = env.addSource(kafkaConsumer);
    final int[] tVoteCount = {0};
    
    DataStream<String> kafkaProducerStream = kafkaConsumerStream.map(new MapFunction<String, String>() {
        @Override
        public String map(String value) throws InterruptedException, IOException {
            JsonNode jsonNode = jsonParser.readValue(value, JsonNode.class);
            Tcount = Tcount + jsonNode.get(key1).asInt();
            int nameCandidate = jsonNode.get(key2).asInt();
            System.out.println(Tcount);
            String tCountT = Integer.toString(Tcount);
            //tVoteCount = tVoteCount + voteCount;
             
            //waitForEventTime(timeStamp);
            return tCountT;
        }
    });
    kafkaConsumerStream.print();
    System.out.println("sdjknvksjdnv"+Tcount);
    Properties producerProperties = new Properties();
    producerProperties.setProperty("bootstrap.servers", serverURL);
    FlinkKafkaProducer<String> kafkaProducer = new FlinkKafkaProducer<>(producerTopicName,
            new SimpleStringSchema(), producerProperties);
    kafkaProducerStream.addSink(kafkaProducer);
    env.execute();

谢谢。

【问题讨论】:

    标签: apache-flink kafka-producer-api


    【解决方案1】:

    blog 中,您将找到一个关于如何将键和主题写入主题的示例:

    您需要将创建的 new FlinkKafkaProducer 替换为以下内容:

    FlinkKafkaProducer<KafkaRecord> kafkaProducer = 
      new FlinkKafkaProducer<KafkaRecord>(
        producerTopicName, 
        ((record, timestamp) -> new ProducerRecord<byte[], byte[]>(producerTopicName, record.key.getBytes(), record.value.getBytes())), 
        producerProperties
      );
    

    【讨论】:

    • 你好,迈克。你的信息很有帮助。但我需要更多帮助。您提供的解决方案按基于键的值分组。但是当一条消息到达并且没有其他具有相同键的消息时,则不会执行计算,并且相同的消息会在 Producer Topic 上发布。有没有一种方法可以使用我的自定义密钥发送新消息并执行计算?谢谢
    • 嗨@QasimKhan,很抱歉我并不完全熟悉它。请提出一个新问题,准确说明您要实现的目标以及当前代码的外观。
    【解决方案2】:

    如果您提供自己的KafkaSerializationSchema 而不是使用SimpleStringSchema,那么您将可以完全控制所写的内容。 @mike 在他的回答中提供了一个如何做到这一点的例子。

    【讨论】:

      猜你喜欢
      • 2020-12-14
      • 2018-11-07
      • 1970-01-01
      • 2019-12-10
      • 2021-11-08
      • 2021-11-25
      • 1970-01-01
      • 2019-07-18
      • 1970-01-01
      相关资源
      最近更新 更多