[
https://issues.apache.org/jira/browse/FLINK-33524?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
ASF GitHub Bot updated FLINK-33524:
-----------------------------------
Labels: pull-request-available (was: )
> IntervalJoinOperator 's judgment on late data has bug
> -----------------------------------------------------
>
> Key: FLINK-33524
> URL: https://issues.apache.org/jira/browse/FLINK-33524
> Project: Flink
> Issue Type: Bug
> Components: API / DataStream
> Affects Versions: 1.18.0
> Environment: Due to the Watermark calculation method :
> {code:java}
> public void onPeriodicEmit(WatermarkOutput output) {
> output.emitWatermark(new Watermark(maxTs - delayTime - 1L));
> }{code}
> data that was delayed by 1 millisecond in this method was incorrectly
> determined
> {code:java}
> private boolean isLate(long timestamp) {
> long currentWatermark = internalTimerService.currentWatermark();
> return timestamp < currentWatermark;
> } {code}
>
> Reporter: ZhangTao
> Priority: Minor
> Labels: pull-request-available
>
> package:
> org.apache.flink.streaming.api.operators.co;
>
> class: IntervalJoinOperator
>
> method:
> isLate
>
> When data with a 1-millisecond delay enters the judgment, an incorrect value
> will be returned
> {code:java}
> private boolean isLate(long timestamp) {
> long currentWatermark = internalTimerService.currentWatermark();
> return timestamp < currentWatermark;
> } {code}
--
This message was sent by Atlassian Jira
(v8.20.10#820010)