apache_beam.transforms.trigger module¶
Support for Dataflow triggers.
Triggers control when in processing time windows get emitted.
-
class
apache_beam.transforms.trigger.
AccumulationMode
[source]¶ Bases:
future.types.newobject.newobject
Controls what to do with data when a trigger fires multiple times.
-
DISCARDING
= 1¶
-
ACCUMULATING
= 2¶
-
-
class
apache_beam.transforms.trigger.
TriggerFn
[source]¶ Bases:
future.types.newobject.newobject
A TriggerFn determines when window (panes) are emitted.
See https://beam.apache.org/documentation/programming-guide/#triggers
-
on_element
(element, window, context)[source]¶ Called when a new element arrives in a window.
Parameters: - element – the element being added
- window – the window to which the element is being added
- context – a context (e.g. a TriggerContext instance) for managing state and setting timers
-
on_merge
(to_be_merged, merge_result, context)[source]¶ Called when multiple windows are merged.
Parameters: - to_be_merged – the set of windows to be merged
- merge_result – the window into which the windows are being merged
- context – a context (e.g. a TriggerContext instance) for managing state and setting timers
-
should_fire
(time_domain, timestamp, window, context)[source]¶ Whether this trigger should cause the window to fire.
Parameters: - time_domain – WATERMARK for event-time timers and REAL_TIME for processing-time timers.
- timestamp – for time_domain WATERMARK, it represents the watermark: (a lower bound on) the watermark of the system and for time_domain REAL_TIME, it represents the trigger: timestamp of the processing-time timer.
- window – the window whose trigger is being considered
- context – a context (e.g. a TriggerContext instance) for managing state and setting timers
Returns: whether this trigger should cause a firing
-
on_fire
(watermark, window, context)[source]¶ Called when a trigger actually fires.
Parameters: - watermark – (a lower bound on) the watermark of the system
- window – the window whose trigger is being fired
- context – a context (e.g. a TriggerContext instance) for managing state and setting timers
Returns: whether this trigger is finished
-
-
class
apache_beam.transforms.trigger.
DefaultTrigger
[source]¶ Bases:
apache_beam.transforms.trigger.TriggerFn
Semantically Repeatedly(AfterWatermark()), but more optimized.
-
class
apache_beam.transforms.trigger.
AfterProcessingTime
(delay=0)[source]¶ Bases:
apache_beam.transforms.trigger.TriggerFn
Fire exactly once after a specified delay from processing time.
AfterProcessingTime is experimental. No backwards compatibility guarantees.
Initialize a processing time trigger with a delay in seconds.
-
class
apache_beam.transforms.trigger.
AfterWatermark
(early=None, late=None)[source]¶ Bases:
apache_beam.transforms.trigger.TriggerFn
Fire exactly once when the watermark passes the end of the window.
Parameters: - early – if not None, a speculative trigger to repeatedly evaluate before the watermark passes the end of the window
- late – if not None, a speculative trigger to repeatedly evaluate after the watermark passes the end of the window
-
LATE_TAG
= CombiningValueStateTag(is_late, CallableWrapperCombineFn(<built-in function any>))¶
-
class
apache_beam.transforms.trigger.
AfterCount
(count)[source]¶ Bases:
apache_beam.transforms.trigger.TriggerFn
Fire when there are at least count elements in this window pane.
AfterCount is experimental. No backwards compatibility guarantees.
-
COUNT_TAG
= CombiningValueStateTag(count, <apache_beam.transforms.combiners.CountCombineFn object>)¶
-
-
class
apache_beam.transforms.trigger.
Repeatedly
(underlying)[source]¶ Bases:
apache_beam.transforms.trigger.TriggerFn
Repeatedly invoke the given trigger, never finishing.
-
class
apache_beam.transforms.trigger.
AfterAny
(*triggers)[source]¶ Bases:
apache_beam.transforms.trigger._ParallelTriggerFn
Fires when any subtrigger fires.
Also finishes when any subtrigger finishes.
-
combine_op
()¶ any(iterable) -> bool
Return True if bool(x) is True for any x in the iterable. If the iterable is empty, return False.
-
-
class
apache_beam.transforms.trigger.
AfterAll
(*triggers)[source]¶ Bases:
apache_beam.transforms.trigger._ParallelTriggerFn
Fires when all subtriggers have fired.
Also finishes when all subtriggers have finished.
-
combine_op
()¶ all(iterable) -> bool
Return True if bool(x) is True for all values x in the iterable. If the iterable is empty, return True.
-
-
class
apache_beam.transforms.trigger.
AfterEach
(*triggers)[source]¶ Bases:
apache_beam.transforms.trigger.TriggerFn
-
INDEX_TAG
= CombiningValueStateTag(index, CallableWrapperCombineFn(<function <lambda>>))¶
-