Skip to content

Commit 06315f5

Browse files
committed
Construct after-processing-time trigger from proto and define trigger callbacks.
1 parent 2beb75c commit 06315f5

3 files changed

Lines changed: 199 additions & 4 deletions

File tree

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

Lines changed: 85 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -79,8 +79,9 @@ func (ws WinStrat) String() string {
7979

8080
// triggerInput represents a Key + window + stage's trigger conditions.
8181
type triggerInput struct {
82-
newElementCount int // The number of new elements since the last check.
83-
endOfWindowReached bool // Whether or not the end of the window has been reached.
82+
newElementCount int // The number of new elements since the last check.
83+
endOfWindowReached bool // Whether or not the end of the window has been reached.
84+
emNow mtime.Time // The current processing time in the runner.
8485
}
8586

8687
// Trigger represents a trigger for a windowing strategy. A trigger determines when
@@ -573,4 +574,85 @@ func (t *TriggerDefault) String() string {
573574
return "Default"
574575
}
575576

576-
// TODO https://github.com/apache/beam/issues/31438 Handle TriggerAfterProcessingTime
577+
// TimestampTransform is the engine's representation of a processing time transform.
578+
type TimestampTransform struct {
579+
Delay time.Duration
580+
AlignToPeriod time.Duration
581+
AlignToOffset time.Duration
582+
}
583+
584+
// TriggerAfterProcessingTime fires once after a specified amount of processing time
585+
// has passed since an element was first seen.
586+
// Uses the extra state field to track the processing time of the first element.
587+
type TriggerAfterProcessingTime struct {
588+
Transforms []TimestampTransform
589+
}
590+
591+
type afterProcessingTimeState struct {
592+
emNow mtime.Time
593+
firingTime mtime.Time
594+
}
595+
596+
func (t *TriggerAfterProcessingTime) onElement(input triggerInput, state *StateData) {
597+
ts := state.getTriggerState(t)
598+
if ts.finished {
599+
return
600+
}
601+
602+
if ts.extra == nil {
603+
ts.extra = afterProcessingTimeState{
604+
emNow: input.emNow,
605+
firingTime: t.applyTimestampTransforms(input.emNow),
606+
}
607+
} else {
608+
s, _ := ts.extra.(afterProcessingTimeState)
609+
s.emNow = input.emNow
610+
ts.extra = s
611+
}
612+
613+
state.setTriggerState(t, ts)
614+
}
615+
616+
func (t *TriggerAfterProcessingTime) applyTimestampTransforms(start mtime.Time) mtime.Time {
617+
ret := start
618+
for _, transform := range t.Transforms {
619+
ret = ret + mtime.Time(transform.Delay/time.Millisecond)
620+
if transform.AlignToPeriod > 0 {
621+
// timestamp - (timestamp % period) + period
622+
// And with an offset, we adjust before and after.
623+
tsMs := ret
624+
periodMs := mtime.Time(transform.AlignToPeriod / time.Millisecond)
625+
offsetMs := mtime.Time(transform.AlignToOffset / time.Millisecond)
626+
627+
adjustedMs := tsMs - offsetMs
628+
alignedMs := adjustedMs - (adjustedMs % periodMs) + periodMs + offsetMs
629+
ret = alignedMs
630+
}
631+
}
632+
return ret
633+
}
634+
635+
func (t *TriggerAfterProcessingTime) shouldFire(state *StateData) bool {
636+
ts := state.getTriggerState(t)
637+
if ts.extra == nil {
638+
return false
639+
}
640+
s := ts.extra.(afterProcessingTimeState)
641+
return s.emNow >= s.firingTime
642+
}
643+
644+
func (t *TriggerAfterProcessingTime) onFire(state *StateData) {
645+
ts := state.getTriggerState(t)
646+
if ts.finished {
647+
return
648+
}
649+
triggerClearAndFinish(t, state)
650+
}
651+
652+
func (t *TriggerAfterProcessingTime) reset(state *StateData) {
653+
delete(state.Trigger, t)
654+
}
655+
656+
func (t *TriggerAfterProcessingTime) String() string {
657+
return fmt.Sprintf("AfterProcessingTime[%v]", t.Transforms)
658+
}

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

Lines changed: 93 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -401,6 +401,99 @@ func TestTriggers_isReady(t *testing.T) {
401401
{triggerInput{newElementCount: 1, endOfWindowReached: true}, false},
402402
{triggerInput{newElementCount: 1, endOfWindowReached: true}, true}, // Late
403403
},
404+
}, {
405+
name: "afterProcessingTime_Delay_Exact",
406+
trig: &TriggerAfterProcessingTime{
407+
Transforms: []TimestampTransform{
408+
{Delay: 3 * time.Second},
409+
},
410+
},
411+
inputs: []io{
412+
{triggerInput{emNow: 0}, false},
413+
{triggerInput{emNow: 1000}, false},
414+
{triggerInput{emNow: 2000}, false},
415+
{triggerInput{emNow: 3000}, true},
416+
{triggerInput{emNow: 4000}, false},
417+
{triggerInput{emNow: 5000}, false},
418+
{triggerInput{emNow: 6000}, false},
419+
{triggerInput{emNow: 7000}, false},
420+
},
421+
}, {
422+
name: "afterProcessingTime_Delay_Late",
423+
trig: &TriggerAfterProcessingTime{
424+
Transforms: []TimestampTransform{
425+
{Delay: 3 * time.Second},
426+
},
427+
},
428+
inputs: []io{
429+
{triggerInput{emNow: 0}, false},
430+
{triggerInput{emNow: 1000}, false},
431+
{triggerInput{emNow: 2000}, false},
432+
{triggerInput{emNow: 3001}, true}, // a little after the expected firing time
433+
{triggerInput{emNow: 4000}, false},
434+
},
435+
}, {
436+
name: "afterProcessingTime_AlignToPeriodOnly",
437+
trig: &TriggerAfterProcessingTime{
438+
Transforms: []TimestampTransform{
439+
{AlignToPeriod: 5 * time.Second},
440+
},
441+
},
442+
inputs: []io{
443+
{triggerInput{emNow: 1500}, false},
444+
{triggerInput{emNow: 2000}, false},
445+
{triggerInput{emNow: 4999}, false},
446+
{triggerInput{emNow: 5000}, true}, // 1.5 is aligned to 5
447+
{triggerInput{emNow: 5001}, false},
448+
},
449+
}, {
450+
name: "afterProcessingTime_AlignToPeriodAndOffset",
451+
trig: &TriggerAfterProcessingTime{
452+
Transforms: []TimestampTransform{
453+
{AlignToPeriod: 5 * time.Second, AlignToOffset: 200 * time.Millisecond},
454+
},
455+
},
456+
inputs: []io{
457+
{triggerInput{emNow: 1500}, false},
458+
{triggerInput{emNow: 2000}, false},
459+
{triggerInput{emNow: 5119}, false},
460+
{triggerInput{emNow: 5200}, true}, // 1.5 is aligned to 5.2
461+
{triggerInput{emNow: 5201}, false},
462+
},
463+
}, {
464+
name: "afterProcessingTime_TwoTransforms",
465+
trig: &TriggerAfterProcessingTime{
466+
Transforms: []TimestampTransform{
467+
{AlignToPeriod: 5 * time.Second, AlignToOffset: 200 * time.Millisecond},
468+
{Delay: 1 * time.Second},
469+
},
470+
},
471+
inputs: []io{
472+
{triggerInput{emNow: 1500}, false},
473+
{triggerInput{emNow: 2000}, false},
474+
{triggerInput{emNow: 5119}, false},
475+
{triggerInput{emNow: 5200}, false},
476+
{triggerInput{emNow: 5201}, false},
477+
{triggerInput{emNow: 6119}, false},
478+
{triggerInput{emNow: 6200}, true}, // 1.5 is aligned to 6.2
479+
{triggerInput{emNow: 6201}, false},
480+
},
481+
}, {
482+
name: "afterProcessingTime_Repeated", trig: &TriggerRepeatedly{
483+
&TriggerAfterProcessingTime{
484+
Transforms: []TimestampTransform{
485+
{Delay: 3 * time.Second},
486+
}}},
487+
inputs: []io{
488+
{triggerInput{emNow: 0}, false},
489+
{triggerInput{emNow: 1000}, false},
490+
{triggerInput{emNow: 2000}, false},
491+
{triggerInput{emNow: 3000}, true}, // first the first time
492+
{triggerInput{emNow: 4000}, false}, // trigger firing time is set again
493+
{triggerInput{emNow: 5000}, false},
494+
{triggerInput{emNow: 6000}, false},
495+
{triggerInput{emNow: 7000}, true}, // trigger firing again
496+
},
404497
}, {
405498
name: "default",
406499
trig: &TriggerDefault{},

sdks/go/pkg/beam/runners/prism/internal/execute.go

Lines changed: 21 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -461,7 +461,27 @@ func buildTrigger(tpb *pipepb.Trigger) engine.Trigger {
461461
}
462462
case *pipepb.Trigger_Repeat_:
463463
return &engine.TriggerRepeatedly{Repeated: buildTrigger(at.Repeat.GetSubtrigger())}
464-
case *pipepb.Trigger_AfterProcessingTime_, *pipepb.Trigger_AfterSynchronizedProcessingTime_:
464+
case *pipepb.Trigger_AfterProcessingTime_:
465+
var transforms []engine.TimestampTransform
466+
for _, ts := range at.AfterProcessingTime.GetTimestampTransforms() {
467+
var delay, period, offset time.Duration
468+
if d := ts.GetDelay(); d != nil {
469+
delay = time.Duration(d.GetDelayMillis()) * time.Millisecond
470+
}
471+
if align := ts.GetAlignTo(); align != nil {
472+
period = time.Duration(align.GetPeriod()) * time.Millisecond
473+
offset = time.Duration(align.GetOffset()) * time.Millisecond
474+
}
475+
transforms = append(transforms, engine.TimestampTransform{
476+
Delay: delay,
477+
AlignToPeriod: period,
478+
AlignToOffset: offset,
479+
})
480+
}
481+
return &engine.TriggerAfterProcessingTime{
482+
Transforms: transforms,
483+
}
484+
case *pipepb.Trigger_AfterSynchronizedProcessingTime_:
465485
panic(fmt.Sprintf("unsupported trigger: %v", prototext.Format(tpb)))
466486
default:
467487
return &engine.TriggerDefault{}

0 commit comments

Comments
 (0)