【问题标题】:Apache Flink dynamic number of SinksApache Flink 动态接收器数量
【发布时间】:2018-01-05 14:36:21
【问题描述】:

我正在使用 Apache Flink 和 KafkaConsumer 从 Kafka 主题中读取一些值。 我还有一个通过读取文件获得的流。

根据收到的值,我想在不同的 Kafka 主题上写这个流。

基本上,我有一个与许多孩子相关联的领导者网络。对于每个孩子,Leader 需要将读取的流写入一个孩子特定的 Kafka Topic 中,以便孩子可以阅读。 当 child 启动时,它会在从 Leader 读取的 Kafka 主题中注册自己。 问题是我不知道我有多少个孩子。

例如,我从 Kafka Topic 中读取 1,我想将流写入一个名为 Topic1 的 Kafka Topic 中。 我读了1-2,我想写两个Kafka主题(Topic1Topic2)。

我不知道这是否可能,因为为了在主题上写作,我使用 Kafka Producer 以及 addSink 方法,据我了解(以及根据我的尝试)似乎 Flink 需要先验地知道汇的数量。

那么,有没有办法获得这样的行为?

【问题讨论】:

    标签: apache-kafka apache-flink


    【解决方案1】:

    如果我很好地理解了您的问题,我认为您可以使用单个接收器来解决它,因为您可以根据正在处理的记录选择 Kafka 主题。似乎源中的一个元素可能被写入多个主题,在这种情况下,您需要 FlatMapFunction 来复制每个源记录 N 次(每个输出主题一个)。我建议与(主题,记录)成对输出(又名Tuple2)。

    DataStream<Tuple2<String, MyValue>> stream = input.flatMap(new FlatMapFunction<>() {
        public void flatMap(MyValue value, Collector<Tupple2<String, MyValue>> out) {
            for (String topic : topics) {
                out.collect(Tuple2.of(topic, value));
            }
        }
    });
    

    然后,您可以使用先前通过创建带有 KeyedSerializationSchemaFlinkKafkaProducer 计算的主题,在其中您实现 getTargetTopic 以返回该对的第一个元素。

    stream.addSink(new FlinkKafkaProducer10<>(
            "default-topic",
            new KeyedSerializationSchema<>() {
                public String getTargetTopic(Tuple2<String, MyValue> element) {
                    return element.f0;
                }
                ...
            },
            kafkaProperties)
    );
    

    【讨论】:

      【解决方案2】:

      KeyedSerializationSchema 现在已弃用。相反,您必须使用“KafkaSerializationSchema”

      同样可以通过重写序列化方法来实现。

          public ProducerRecord<byte[], byte[]> serialize(
      String inputString, @Nullable Long aLong){ 
              return new ProducerRecord<>(customTopicName,
       key.getBytes(StandardCharsets.UTF_8), inputString.getBytes(StandardCharsets.UTF_8));
      }
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2019-07-25
        • 2018-12-07
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2021-09-09
        相关资源
        最近更新 更多