Skip to content

Commit 0ac0732

Browse files
committed
Construct trigger struct based on pipeline proto.
1 parent 258b8d1 commit 0ac0732

3 files changed

Lines changed: 42 additions & 2 deletions

File tree

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

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -573,4 +573,26 @@ 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 {
579+
Delay time.Duration
580+
AlignToPeriod time.Duration
581+
AlignToOffset time.Duration
582+
}
583+
584+
func (t *TriggerAfterProcessingTime) onElement(input triggerInput, state *StateData) {}
585+
586+
func (t *TriggerAfterProcessingTime) shouldFire(state *StateData) bool {
587+
return false
588+
}
589+
590+
func (t *TriggerAfterProcessingTime) onFire(state *StateData) {}
591+
592+
func (t *TriggerAfterProcessingTime) reset(state *StateData) {}
593+
594+
func (t *TriggerAfterProcessingTime) String() string {
595+
return fmt.Sprintf("AfterProcessingTime[Delay: %v, AlignToPeriod: %v, AlignToOffset: %v]", t.Delay, t.AlignToPeriod, t.AlignToOffset)
596+
}
597+
576598
// TODO https://github.com/apache/beam/issues/31438 Handle TriggerAfterProcessingTime

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

Lines changed: 19 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -457,7 +457,25 @@ func buildTrigger(tpb *pipepb.Trigger) engine.Trigger {
457457
}
458458
case *pipepb.Trigger_Repeat_:
459459
return &engine.TriggerRepeatedly{Repeated: buildTrigger(at.Repeat.GetSubtrigger())}
460-
case *pipepb.Trigger_AfterProcessingTime_, *pipepb.Trigger_AfterSynchronizedProcessingTime_:
460+
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]
465+
if d := ts.GetDelay(); d != nil {
466+
delay = time.Duration(d.GetDelayMillis()) * time.Millisecond
467+
}
468+
if align := ts.GetAlignTo(); align != nil {
469+
period = time.Duration(align.GetPeriod()) * time.Millisecond
470+
offset = time.Duration(align.GetOffset()) * time.Millisecond
471+
}
472+
}
473+
return &engine.TriggerAfterProcessingTime{
474+
Delay: delay,
475+
AlignToPeriod: period,
476+
AlignToOffset: offset,
477+
}
478+
case *pipepb.Trigger_AfterSynchronizedProcessingTime_:
461479
panic(fmt.Sprintf("unsupported trigger: %v", prototext.Format(tpb)))
462480
default:
463481
return &engine.TriggerDefault{}

sdks/go/pkg/beam/runners/prism/internal/jobservices/management.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -316,7 +316,7 @@ func (s *Server) Prepare(ctx context.Context, req *jobpb.PrepareJobRequest) (_ *
316316
func hasUnsupportedTriggers(tpb *pipepb.Trigger) bool {
317317
unsupported := false
318318
switch at := tpb.GetTrigger().(type) {
319-
case *pipepb.Trigger_AfterProcessingTime_, *pipepb.Trigger_AfterSynchronizedProcessingTime_:
319+
case *pipepb.Trigger_AfterSynchronizedProcessingTime_:
320320
return true
321321
case *pipepb.Trigger_AfterAll_:
322322
for _, st := range at.AfterAll.GetSubtriggers() {

0 commit comments

Comments
 (0)