【发布时间】: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