Skip to content

Commit 28cc7fa

Browse files
committed
Fill in callbacks for TriggerAfterProcessingTime
1 parent 0ac0732 commit 28cc7fa

2 files changed

Lines changed: 61 additions & 16 deletions

File tree

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

Lines changed: 52 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -573,26 +573,69 @@ func (t *TriggerDefault) String() string {
573573
return "Default"
574574
}
575575

576-
// TriggerAfterProcessingTime fires once after a specified amount of processing time
577-
// has passed since an element was first seen.
578-
type TriggerAfterProcessingTime struct {
576+
// TimestampTransform is the engine's representation of a processing time transform.
577+
type TimestampTransform struct {
579578
Delay time.Duration
580579
AlignToPeriod time.Duration
581580
AlignToOffset time.Duration
582581
}
583582

584-
func (t *TriggerAfterProcessingTime) onElement(input triggerInput, state *StateData) {}
583+
// TriggerAfterProcessingTime fires once after a specified amount of processing time
584+
// has passed since an element was first seen.
585+
// Uses the extra state field to track if the time has been reached.
586+
type TriggerAfterProcessingTime struct {
587+
Transforms []TimestampTransform
588+
}
589+
590+
func (t *TriggerAfterProcessingTime) onElement(input triggerInput, state *StateData) {
591+
ts := state.getTriggerState(t)
592+
if ts.finished {
593+
return
594+
}
595+
596+
if ts.extra == nil {
597+
ts.extra = mtime.Now()
598+
}
599+
600+
state.setTriggerState(t, ts)
601+
}
602+
603+
func (t *TriggerAfterProcessingTime) applyTimestampTransforms(start mtime.Time) mtime.Time {
604+
ret := start
605+
for _, transform := range t.Transforms {
606+
ret = ret + mtime.Time(transform.Delay/time.Millisecond)
607+
if transform.AlignToPeriod > 0 {
608+
// Formula from https://cloud.google.com/blog/products/data-analytics/windowing-and-triggering-in-apache-beam
609+
// timestamp - (timestamp % period) + period
610+
// And with an offset, we adjust before and after.
611+
tsMs := ret
612+
periodMs := mtime.Time(transform.AlignToPeriod / time.Millisecond)
613+
offsetMs := mtime.Time(transform.AlignToOffset / time.Millisecond)
614+
615+
adjustedMs := tsMs - offsetMs
616+
alignedMs := adjustedMs - (adjustedMs % periodMs) + periodMs + offsetMs
617+
ret = alignedMs
618+
}
619+
}
620+
return ret
621+
}
585622

586623
func (t *TriggerAfterProcessingTime) shouldFire(state *StateData) bool {
587-
return false
624+
ts := state.getTriggerState(t)
625+
if ts.extra == nil {
626+
return false
627+
}
628+
startTime := ts.extra.(mtime.Time)
629+
firingTime := t.applyTimestampTransforms(startTime)
630+
return mtime.Now() > firingTime
588631
}
589632

590633
func (t *TriggerAfterProcessingTime) onFire(state *StateData) {}
591634

592-
func (t *TriggerAfterProcessingTime) reset(state *StateData) {}
635+
func (t *TriggerAfterProcessingTime) reset(state *StateData) {
636+
delete(state.Trigger, t)
637+
}
593638

594639
func (t *TriggerAfterProcessingTime) String() string {
595-
return fmt.Sprintf("AfterProcessingTime[Delay: %v, AlignToPeriod: %v, AlignToOffset: %v]", t.Delay, t.AlignToPeriod, t.AlignToOffset)
640+
return fmt.Sprintf("AfterProcessingTime[%v]", t.Transforms)
596641
}
597-
598-
// TODO https://github.com/apache/beam/issues/31438 Handle TriggerAfterProcessingTime

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

Lines changed: 9 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -458,22 +458,24 @@ func buildTrigger(tpb *pipepb.Trigger) engine.Trigger {
458458
case *pipepb.Trigger_Repeat_:
459459
return &engine.TriggerRepeatedly{Repeated: buildTrigger(at.Repeat.GetSubtrigger())}
460460
case *pipepb.Trigger_AfterProcessingTime_:
461-
var delay, period, offset time.Duration
462-
// TODO: support multiple transforms.
463-
if len(at.AfterProcessingTime.GetTimestampTransforms()) > 0 {
464-
ts := at.AfterProcessingTime.GetTimestampTransforms()[0]
461+
var transforms []engine.TimestampTransform
462+
for _, ts := range at.AfterProcessingTime.GetTimestampTransforms() {
463+
var delay, period, offset time.Duration
465464
if d := ts.GetDelay(); d != nil {
466465
delay = time.Duration(d.GetDelayMillis()) * time.Millisecond
467466
}
468467
if align := ts.GetAlignTo(); align != nil {
469468
period = time.Duration(align.GetPeriod()) * time.Millisecond
470469
offset = time.Duration(align.GetOffset()) * time.Millisecond
471470
}
471+
transforms = append(transforms, engine.TimestampTransform{
472+
Delay: delay,
473+
AlignToPeriod: period,
474+
AlignToOffset: offset,
475+
})
472476
}
473477
return &engine.TriggerAfterProcessingTime{
474-
Delay: delay,
475-
AlignToPeriod: period,
476-
AlignToOffset: offset,
478+
Transforms: transforms,
477479
}
478480
case *pipepb.Trigger_AfterSynchronizedProcessingTime_:
479481
panic(fmt.Sprintf("unsupported trigger: %v", prototext.Format(tpb)))

0 commit comments

Comments
 (0)