I'm building a pipeline with apache beam (GCP Dataflow) and python, my pipeline looks like this:
...
with beam.Pipeline(options=self.pipeline_options) as pipeline:
somepipeline = (
pipeline
| "ReadPubSubMessage" >> ReadFromPubSub(
subscription=self.custom_options.some_subscription)
| "Windowing" >> beam.WindowInto(beam.window.FixedWindows(30))
| "DecodePubSubMessage" >> beam.ParDo(DecodePubSubMessage()).with_outputs(ERROR_OUTPUT_NAME, main=MAIN_OUTPUT_NAME)
| "Geting and sorting listings" >> beam.ParDo(SortByCompletion())
| "Batching listings" >> beam.BatchElements(min_batch_size=3,max_batch_size=3)
| "Print logs" >> beam.Map(logging.info)
)
...
And everything works as expected when I run pipeline via DirectRunner (you can see 1 batch with 3 elements inside):

But when I run the same code with a DataflowRunner I'm getting this result (you can see 3 batches with 1 element inside of each batch):

This happens even when I'm running this pipeline in parallel (in two terminal windows). Both were run with a streaming flag. Messages were sent to pubsub via python script one by one immediately.
Question: What can cause this problem with DataflowRunner (my assumption was the number of workers in dataflow but when I checked it there was only 1 worker in this job) and how I can get the same result as via DirrectRunner.
Thank you!