pyflink.datastream.window.ProcessingTimeSessionWindows#
- class ProcessingTimeSessionWindows(session_gap: int)[source]#
A WindowAssigner that windows elements into sessions based on the current processing time. Windows cannot overlap.
For example, the processing interval is set to 1 minutes:
>>> data_stream.key_by(lambda x: x[0], key_type=Types.STRING()) \ ... .window(ProcessingTimeSessionWindows.with_gap(Time.minutes(1)))
Methods
assign_windows
(element, timestamp, context)- param element
The element to which windows should be assigned.
get_default_trigger
(env)- param env
The StreamExecutionEnvironment used to compile the DataStream job.
get_window_serializer
()- return
A TypeSerializer for serializing windows that are assigned by this WindowAssigner.
is_event_time
()- return
True if elements are assigned to windows based on event time, false otherwise.
merge_windows
(windows, callback)Determines which windows (if any) should be merged.
with_dynamic_gap
(extractor)Creates a new SessionWindows WindowAssigner that assigns elements to sessions based on the element timestamp.
with_gap
(size)Creates a new SessionWindows WindowAssigner that assigns elements to sessions based on the element timestamp.