【问题标题】:Flink route stream dynamically to different kafka topicsFlink 路由流动态到不同的 kafka 主题
【发布时间】:2019-09-17 18:16:45
【问题描述】:

Flink 动态路由 Kafka 主题的解决方案是实现 KeyedSerializationSchema 并覆盖 getTargetTopic,但不推荐使用 KeyedSerializationSchema,而应该使用 KafkaSerializationSchema。该接口不提供 getTargetTopic 或类似的东西。

那么,在 Flink 中,在 getTargetTopic 不存在的情况下,Kafka 动态路由应该如何工作?

【问题讨论】:

    标签: dynamic apache-kafka routing apache-flink


    【解决方案1】:

    KafkaSerializationSchema. serialize 返回ProducerRecord<byte[], byte[]> 这个ProducerRecord 包含一个主题。 您可以使用像 https://kafka.apache.org/10/javadoc/org/apache/kafka/clients/producer/ProducerRecord.html#ProducerRecord-java.lang.String-K-V- 这样的构造函数 注入主题。

    考虑到这一点,您只需要创建一个类似的方法

    String dynamicTopic(T element, @Nullable Long timestamp)
    

    你的KafkaSerializationSchema 实现只需要使用它

    ProducerRecord<byte[], byte[]> serialize(T element, @Nullable Long timestamp){
        ...
        return new ProducerRecord(dynamicTopic(element, timestamp), aKey, aValue);
    }
    

    【讨论】:

    • 绝对正确 100% 并且非常简单。昨天肯定有捷径,不记得要创建ProducerRecord,反正你要提供topic!
    猜你喜欢
    • 2019-06-28
    • 2022-06-30
    • 1970-01-01
    • 2019-01-24
    • 1970-01-01
    • 1970-01-01
    • 2020-03-10
    • 2019-12-15
    • 2019-12-07
    相关资源
    最近更新 更多