开发者社区> 问答> 正文

对于kafka partition 设置时间戳及watermark

我按官网操作,重写了序列化方式

val kafkaSource = new FlinkKafkaConsumer09[MyType]("myTopic", schema,

props)kafkaSource.assignTimestampsAndWatermarks(new

AscendingTimestampExtractor[MyType] {

def extractAscendingTimestamp(element: MyType): Long =

element.eventTimestamp})

val stream: DataStream[MyType] = env.addSource(kafkaSource)

有个疑问,这样写完之后是不是不用设置setAutoWatermarkInterval 呢?*来自志愿者整理的flink邮件归档

展开
收起
又出bug了-- 2021-12-02 11:44:27 685 0
1 条回答
写回答
取消 提交回答
  • setAutoWatermarkInterval这个只是设置interval。决定你那个抽取ts的函数的执行频率的。*来自志愿者整理的FLINK邮件归档

    2021-12-02 14:26:06
    赞同 展开评论 打赏
问答排行榜
最热
最新

相关电子书

更多
Java Spring Boot开发实战系列课程【第16讲】:Spring Boot 2.0 实战Apache Kafka百万级高并发消息中间件与原理解析 立即下载
MaxCompute技术公开课第四季 之 如何将Kafka数据同步至MaxCompute 立即下载
消息队列kafka介绍 立即下载