0
votes

I am trying to develop dataflow pipeline using sliding window on bounded and streaming dataset using Apache Beam Python sdk. The pipeline is as follows:

  • Reading data
  • Assigning timestamps
  • Windowing using SlidingWindows() of 3 size and 1 period
  • Grouping the elements in the windows
  • Printing the output

The sample Data:

  data = [{'serverID': 'server_1', 'CPU_Utilization': 0, 'timestamp': 1},
          {'serverID': 'server_1', 'CPU_Utilization': 1, 'timestamp': 2},
          {'serverID': 'server_1', 'CPU_Utilization': 2, 'timestamp': 3},
          {'serverID': 'server_1', 'CPU_Utilization': 3, 'timestamp': 4}]

The beam pipeline:

    import apache_beam as beam
    from apache_beam.options.pipeline_options import PipelineOptions
    from apache_beam.options.pipeline_options import SetupOptions

    pipeline_options = PipelineOptions()
    pipeline_options.view_as(SetupOptions).save_main_session = True 

    p = beam.Pipeline(options=pipeline_options)

    with beam.Pipeline(options=pipeline_options) as p:
         events =  (p | 'Create Events' >> beam.Create(data) \
                      | 'Add Timestamps' >> beam.Map(lambda x: beam.window.TimestampedValue(x, 
                                              x['timestamp'])) \
                      | 'PairWithOne' >> beam.Map(lambda x: (None, x)) \
                      |'Sliding Window'>>beam.WindowInto(beam.window.SlidingWindows(3,1))
                      | 'Group by key' >> beam.GroupByKey()
                      |beam.Map(print))

The output I got:

    (None, [{'serverID': 'server_1', 'CPU_Utilization': 0, 'timestamp': 1}, {'serverID': 
    'server_1', 'CPU_Utilization': 1, 'timestamp': 2}, {'serverID': 'server_1', 
    'CPU_Utilization': 2, 'timestamp': 3}])
    (None, [{'serverID': 'server_1', 'CPU_Utilization': 0, 'timestamp': 1}, {'serverID': 
    'server_1', 'CPU_Utilization': 1, 'timestamp': 2}])
    (None, [{'serverID': 'server_1', 'CPU_Utilization': 0, 'timestamp': 1}])
    (None, [{'serverID': 'server_1', 'CPU_Utilization': 1, 'timestamp': 2}, {'serverID': 
    'server_1', 'CPU_Utilization': 2, 'timestamp': 3}, {'serverID': 'server_1', 
    'CPU_Utilization': 3, 'timestamp': 4}])
    (None, [{'serverID': 'server_1', 'CPU_Utilization': 2, 'timestamp': 3}, {'serverID': 
    'server_1', 'CPU_Utilization': 3, 'timestamp': 4}])
    (None, [{'serverID': 'server_1', 'CPU_Utilization': 3, 'timestamp': 4}])

The expected output (i.e to discard data when starting window timestamp is less than the timestamp of 1st row or element in data):

    (None, [{'serverID': 'server_1', 'CPU_Utilization': 0, 'timestamp': 1}, {'serverID': 
    'server_1', 'CPU_Utilization': 1, 'timestamp': 2}, {'serverID': 'server_1', 
    'CPU_Utilization': 2, 'timestamp': 3}])
    (None, [{'serverID': 'server_1', 'CPU_Utilization': 1, 'timestamp': 2}, {'serverID': 
    'server_1', 'CPU_Utilization': 2, 'timestamp': 3}, {'serverID': 'server_1', 
    'CPU_Utilization': 3, 'timestamp': 4}])
    (None, [{'serverID': 'server_1', 'CPU_Utilization': 2, 'timestamp': 3}, {'serverID': 
    'server_1', 'CPU_Utilization': 3, 'timestamp': 4}])
    (None, [{'serverID': 'server_1', 'CPU_Utilization': 3, 'timestamp': 4}])

I have also tried AfterCount(n) trigger but this trigger does not consider the data when number of data points is less than n.
Any help on this would be really appreciated.

1

1 Answers

0
votes

It's probably working as intended. Could you try validating the output by writing it into a file instead of using print? It could be that there are multiple panes fired and printed but the accumulation_mode is discarding.

Also, you can try using Beam Notebooks when debugging streaming pipelines. Use something like show(pcoll, include_window_info=True) to visualize the event time, windows and pane info of a PCollection instead of printing them.