Skip to content

Commit d772ece

Browse files
committed
Handle the case when after-processing-time trigger is called repeated.
1 parent f8c6594 commit d772ece

2 files changed

Lines changed: 48 additions & 9 deletions

File tree

sdks/go/pkg/beam/runners/prism/internal/engine/strategy.go

Lines changed: 26 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -589,8 +589,9 @@ type TriggerAfterProcessingTime struct {
589589
}
590590

591591
type afterProcessingTimeState struct {
592-
emNow mtime.Time
593-
firingTime mtime.Time
592+
emNow mtime.Time
593+
firingTime mtime.Time
594+
endOfWindowReached bool
594595
}
595596

596597
func (t *TriggerAfterProcessingTime) onElement(input triggerInput, state *StateData) {
@@ -601,12 +602,14 @@ func (t *TriggerAfterProcessingTime) onElement(input triggerInput, state *StateD
601602

602603
if ts.extra == nil {
603604
ts.extra = afterProcessingTimeState{
604-
emNow: input.emNow,
605-
firingTime: t.applyTimestampTransforms(input.emNow),
605+
emNow: input.emNow,
606+
firingTime: t.applyTimestampTransforms(input.emNow),
607+
endOfWindowReached: input.endOfWindowReached,
606608
}
607609
} else {
608610
s, _ := ts.extra.(afterProcessingTimeState)
609611
s.emNow = input.emNow
612+
s.endOfWindowReached = input.endOfWindowReached
610613
ts.extra = s
611614
}
612615

@@ -646,11 +649,28 @@ func (t *TriggerAfterProcessingTime) onFire(state *StateData) {
646649
if ts.finished {
647650
return
648651
}
649-
triggerClearAndFinish(t, state)
652+
653+
// We don't reset the state here, only mark it as finished
654+
ts.finished = true
655+
s := ts.extra.(afterProcessingTimeState)
656+
s.firingTime = t.applyTimestampTransforms(s.firingTime) // compute next possible firing time
657+
ts.extra = s
658+
state.setTriggerState(t, ts)
650659
}
651660

652661
func (t *TriggerAfterProcessingTime) reset(state *StateData) {
653-
delete(state.Trigger, t)
662+
ts := state.getTriggerState(t)
663+
if ts.extra != nil {
664+
if ts.extra.(afterProcessingTimeState).endOfWindowReached {
665+
delete(state.Trigger, t)
666+
return
667+
}
668+
}
669+
670+
// Not reaching the end of window yet.
671+
// We keep the state (especially the next possible firing time) in case the trigger is called again
672+
ts.finished = false
673+
state.setTriggerState(t, ts)
654674
}
655675

656676
func (t *TriggerAfterProcessingTime) String() string {

sdks/go/pkg/beam/runners/prism/internal/engine/strategy_test.go

Lines changed: 22 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -488,11 +488,30 @@ func TestTriggers_isReady(t *testing.T) {
488488
{triggerInput{emNow: 0}, false},
489489
{triggerInput{emNow: 1000}, false},
490490
{triggerInput{emNow: 2000}, false},
491-
{triggerInput{emNow: 3000}, true}, // first the first time
491+
{triggerInput{emNow: 3000}, true}, // firing the first time, trigger set again
492492
{triggerInput{emNow: 4000}, false}, // trigger firing time is set again
493493
{triggerInput{emNow: 5000}, false},
494-
{triggerInput{emNow: 6000}, false},
495-
{triggerInput{emNow: 7000}, true}, // trigger firing again
494+
{triggerInput{emNow: 6000}, true}, // firing the second time
495+
},
496+
}, {
497+
name: "afterProcessingTime_Repeated_AcrossWindows", trig: &TriggerRepeatedly{
498+
&TriggerAfterProcessingTime{
499+
Transforms: []TimestampTransform{
500+
{Delay: 3 * time.Second},
501+
}}},
502+
inputs: []io{
503+
{triggerInput{emNow: 0}, false},
504+
{triggerInput{emNow: 1000}, false},
505+
{triggerInput{emNow: 2000}, false},
506+
{triggerInput{emNow: 3000}, true}, // fire the first time, trigger is set again
507+
{triggerInput{emNow: 4000}, false},
508+
{triggerInput{emNow: 5000}, false},
509+
{triggerInput{emNow: 6000,
510+
endOfWindowReached: true}, true}, // fire the second time, reach end of window and start over
511+
{triggerInput{emNow: 7000}, false}, // trigger firing time is set to 7s + 3s = 10s
512+
{triggerInput{emNow: 8000}, false},
513+
{triggerInput{emNow: 9000}, false},
514+
{triggerInput{emNow: 10000}, true}, // fire in the new window
496515
},
497516
}, {
498517
name: "default",

0 commit comments

Comments
 (0)