【问题标题】:Flink: assign watermark to FlinkKafkaConsumerFlink:为 FlinkKafkaConsumer 分配水印
【发布时间】:2020-06-12 00:51:08
【问题描述】:

我有一个 FlinkKafkaConsumer 定义如下 FlinkKafkaConsumer[String]("topic", new SimpleStringSchema(), properties),我正在使用 setStreamTimeCharacteristic(TimeCharacteristic.EventTime) 处理事件时间。

现在我想用函数assignTimestampsAndWatermarks 分配一个周期性水印,但我不知道我应该传递给该函数什么,因为在文档中,此函数的示例接收一个类型为MyType 的元素getCreationTime() 而我的消费者是字符串类型。

在这种情况下是否可以分配事件时间?

编辑:我想用作事件时间的时间是每个寄存器存储在 Kafka 中的时间。

【问题讨论】:

    标签: apache-kafka apache-flink flink-streaming


    【解决方案1】:

    EventTime 的概念至少在定义中与事件的创建时间而非接收时间严格相关。因此,如果您从 Kafka 消费的事件具有某种时间戳(例如,如果您将 JSON 作为字符串消费然后对其进行解析),那么您可以在 assignTimestampsAndWatermarks 函数中使用此时间戳。

    如果您正在解析普通的 String 对象,那么您可以做的最好的事情是使用自定义 KafkaDeserializationSchema 为每个事件提取 Kafka 时间戳并使用它。

    从技术上讲,您甚至可以使用为每条记录人为增加时间戳的计数器(例如通过将其增加 1),但这在 EventTime 处理方面似乎没有意义。

    【讨论】:

    • 我使用的字符串是 JSON,但没有任何带有时间戳的字段。这种情况我该怎么办?
    • 你为什么要使用EventTime呢?:)
    • 原来是使用IntervalJoin。我认为可能有一种方法可以从 Kafka 主题中获取事件时间,但只有在接收到的数据具有隐式时间戳(即带有它的 json 中的字段)时才能使用它。我理解正确吗?
    • 另一个问题:是否可以将存储时间(数据存储在kafka中时)用作水印?
    • 如果您使用KafkaDeserializationSchema,那么在反序列化时您可以访问整个记录,这意味着您还可以访问 Kafka 在将记录写入主题时分配的时间戳。你可以使用它:) 它不需要嵌入到 JSON 本身中。
    猜你喜欢
    • 2021-07-09
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多