When trying to run a simple streaming pipeline with a Splittable DoFn (SDF) based generator, python DirectRunner fails with AssertionError: wrong timestamp for StageNode<inputs=['ref_PCollection_PCollection_3_split'],side_inputs=[]. after generating exactly one input from the SDF.
This is a minimal reproducible example with apache_beam.transforms.periodicsequence.ImpulseSeqGenDoFn as the SDF:
pipeline_options = PipelineOptions(pipeline_args, streaming=True, save_main_session=True)
start = Timestamp.now()
stop = start + Duration(seconds=(24 * 60 * 60))
interval = 5.0
pipeline = Pipeline(options=pipeline_options)
(pipeline
| "Create" >> Create([(start, stop, interval)])
| "Generator" >> ParDo(ImpulseSeqGenDoFn())
| "Print" >> ParDo(lambda x: logging.info(f"GEN: {x}")))
pipeline.run().wait_until_finish()
This fails when executed with python3 generator.py --runner DirectRunner:
...
INFO:root:GEN: 1632697296.421499
Traceback (most recent call last):
File "generator.py", line 157, in <module>
r = pipeline.run()
File "/usr/local/lib/python3.6/dist-packages/apache_beam/pipeline.py", line 564, in run
return self.runner.run_pipeline(self, self._options)
File "/usr/local/lib/python3.6/dist-packages/apache_beam/runners/direct/direct_runner.py", line 131, in run_pipeline
return runner.run_pipeline(pipeline, options)
File "/usr/local/lib/python3.6/dist-packages/apache_beam/runners/portability/fn_api_runner/fn_runner.py", line 196, in run_pipeline
pipeline.to_runner_api(default_environment=self._default_environment))
File "/usr/local/lib/python3.6/dist-packages/apache_beam/runners/portability/fn_api_runner/fn_runner.py", line 206, in run_via_runner_api
return self.run_stages(stage_context, stages)
File "/usr/local/lib/python3.6/dist-packages/apache_beam/runners/portability/fn_api_runner/fn_runner.py", line 385, in run_stages
runner_execution_context, bundle_context_manager)
File "/usr/local/lib/python3.6/dist-packages/apache_beam/runners/portability/fn_api_runner/fn_runner.py", line 668, in _run_stage
bundle_context_manager.stage.name))
AssertionError: wrong timestamp for StageNode<inputs=['ref_PCollection_PCollection_3_split'],side_inputs=[].
The exact same pipeline runs as expected on Dataflow runner, generating messages in correct intervals.
Since the error happens after the first input from the SDF, my though was that the error happens when deferring the remainder of the work for the SDF. The actual assertion from /apache_beam/runners/portability/fn_api_runner/fn_runner.py:668 refers to a watermark.
assert (runner_execution_context.watermark_manager.get_stage_node(
bundle_context_manager.stage.name).output_watermark()
< timestamp.MAX_TIMESTAMP), (
'wrong timestamp for %s. '
% runner_execution_context.watermark_manager.get_stage_node(
bundle_context_manager.stage.name))
The SDF apache_beam.transforms.periodicsequence.ImpulseSeqGenDoFn does not explicitly use a watermark. Neither is this SDF decorated with @beam.DoFn.unbounded_per_element(), suggesting maybe that it is not intended for streaming pipelines. Running the same job with streaming=False yields the same error.
When adding a step that expands to iobase.Sink, this issue no longer occurs and the pipeline behaves correctly, with the SDF generating outputs in the correct intervals. Following is an example of working pipeline:
pipeline_options = PipelineOptions(pipeline_args, streaming=True, save_main_session=True)
start = Timestamp.now()
stop = start + Duration(seconds=(24 * 60 * 60))
interval = 5.0
pipeline = Pipeline(options=pipeline_options)
(pipeline
| "Create" >> Create([(start, stop, interval)])
| "Generator" >> ParDo(ImpulseSeqGenDoFn())
| "Print" >> ParDo(lambda x: logging.info(f"GEN: {x}"))
| 'Write' >> WriteToText("generated.out"))
pipeline.run().wait_until_finish()
I am assuming that because Dataflow executes the pipeline correctly, this may be an issue with DirectRunner and justifies a bug report. But to be sure, has anyone encountered a similar issue with DirectRunner and SDF or is this not a correct way to build a streaming pipeline based on SDF?