Skip to content

Commit 1670089

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

2 files changed

Lines changed: 50 additions & 11 deletions

File tree

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

Lines changed: 27 additions & 7 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

@@ -634,7 +637,7 @@ func (t *TriggerAfterProcessingTime) applyTimestampTransforms(start mtime.Time)
634637

635638
func (t *TriggerAfterProcessingTime) shouldFire(state *StateData) bool {
636639
ts := state.getTriggerState(t)
637-
if ts.extra == nil {
640+
if ts.extra == nil || ts.finished {
638641
return false
639642
}
640643
s := ts.extra.(afterProcessingTimeState)
@@ -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+
state.setTriggerState(t, ts)
650656
}
651657

652658
func (t *TriggerAfterProcessingTime) reset(state *StateData) {
653-
delete(state.Trigger, t)
659+
ts := state.getTriggerState(t)
660+
if ts.extra != nil {
661+
if ts.extra.(afterProcessingTimeState).endOfWindowReached {
662+
delete(state.Trigger, t)
663+
return
664+
}
665+
}
666+
667+
// Not reaching the end of window yet.
668+
// We keep the state (especially the next possible firing time) in case the trigger is called again
669+
ts.finished = false
670+
s := ts.extra.(afterProcessingTimeState)
671+
s.firingTime = t.applyTimestampTransforms(s.firingTime) // compute next possible firing time
672+
ts.extra = s
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: 23 additions & 4 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
492-
{triggerInput{emNow: 4000}, false}, // trigger firing time is set again
491+
{triggerInput{emNow: 3000}, true}, // firing the first time, trigger set again
492+
{triggerInput{emNow: 4000}, false},
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)