首页
学习
活动
专区
圈层
工具
发布
社区首页 >问答首页 >Flink AssignTimeStampsAndWaterMark

Flink AssignTimeStampsAndWaterMark
EN

Stack Overflow用户
提问于 2022-05-25 17:42:35
回答 1查看 70关注 0票数 0

我想通过request_time计算健康检查数据的状态代码,从当前时间起有一个1分钟的窗口。

据了解,健康检查每分钟发送大约60个请求。所以结果应该是{environment: "XX", host: "XX", 200: 60, 300: 0, 400: 0, 500: 0}

但是实际的结果是{environment: "XX", host: "XX", 200: 50000, 300: 0, 400: 0, 500: 0},它计算了许多以前的数据。

我的代码就像

代码语言:javascript
复制
  env.fromSource(kafkasource).filter().flatMap()
    .assignTimeStampAndWaterMarks(
    WaterMarkStrategy.<OnjectNode>forboundedOutOfOrderness(DurationOfSeconds(60).withTimeStampAssigner(assigner))
    .keyBy("env","host")
    .window(1m)
    .reduce()

有人知道缺了什么或者我在逻辑上错了吗?

EN

回答 1

Stack Overflow用户

发布于 2022-05-26 08:09:49

也许您应该在reduce()flatmap()中提供更详细的代码信息。根据当前可用的信息,此问题可能是reduce中的状态未被重置。您可以检查reduce中的代码或是否正确地设置了ttl。

票数 0
EN
页面原文内容由Stack Overflow提供。腾讯云小微IT领域专用引擎提供翻译支持
原文链接:

https://stackoverflow.com/questions/72382115

复制
相关文章

相似问题

领券
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档