GoogleCloudPlatform / GoogleCloudPlatform/training-data-analyst

Issue with 7_Advanced_Streaming_Analytics/streaming_minute_traffic_pipeline.py

Open
#1,908 1 comment 1 reaction 0 assignees View on GitHub
Dominant language
Jupyter Notebook
Stars
8.6k
Forks
6.1k
Avg merge
4h 44m
Merged PRs (30d)
2

Description

Line 112 `trigger=AfterProcessingTime(120)` leads to the following issue:

`python3 streaming_minute_traffic_pipeline.py --project=${PROJECT_ID} --region=${REGION} --staging_location=${PIPELINE_FOLDER}/staging --temp_location=${PIPELINE_FOLDER}/temp --runner=${RUNNER} --input_topic=${PUBSUB_TOPIC} --window_duration=${WINDOW_DURATION} --allowed_lateness=${ALLOWED_LATENESS} --table_name=${OUTPUT_TABLE_NAME} --dead_letter_bucket=${DEADLETTER_BUCKET}
streaming_minute_traffic_pipeline.py:114: FutureWarning: WriteToFiles is experimental.
| 'WriteUnparsedToGCS' >> fileio.WriteToFiles(output_path, shards=1, max_writers_per_bundle=0)
/home/jupyter/project/training-data-analyst/quests/dataflow_python/7_Advanced_Streaming_Analytics/lab/df-env/lib/python3.7/site-packages/apache_beam/io/fileio.py:550: BeamDeprecationWarning: options is deprecated since First stable release. References to .options will not be supported
p.options.view_as(GoogleCloudOptions).temp_location or
Traceback (most recent call last):
File "streaming_minute_traffic_pipeline.py", line 138, in
run()
File "streaming_minute_traffic_pipeline.py", line 114, in run
| 'WriteUnparsedToGCS' >> fileio.WriteToFiles(output_path, shards=1, max_writers_per_bundle=0)
File "/home/jupyter/project/training-data-analyst/quests/dataflow_python/7_Advanced_Streaming_Analytics/lab/df-env/lib/python3.7/site-packages/apache_beam/pvalue.py", line 137, in __or__
return self.pipeline.apply(ptransform, self)
File "/home/jupyter/project/training-data-analyst/quests/dataflow_python/7_Advanced_Streaming_Analytics/lab/df-env/lib/python3.7/site-packages/apache_beam/pipeline.py", line 652, in apply
transform.transform, pvalueish, label or transform.label)
File "/home/jupyter/project/training-data-analyst/quests/dataflow_python/7_Advanced_Streaming_Analytics/lab/df-env/lib/python3.7/site-packages/apache_beam/pipeline.py", line 662, in apply
return self.apply(transform, pvalueish)
File "/home/jupyter/project/training-data-analyst/quests/dataflow_python/7_Advanced_Streaming_Analytics/lab/df-env/lib/python3.7/site-packages/apache_beam/pipeline.py", line 708, in apply
pvalueish_result = self.runner.apply(transform, pvalueish, self._options)
File "/home/jupyter/project/training-data-analyst/quests/dataflow_python/7_Advanced_Streaming_Analytics/lab/df-env/lib/python3.7/site-packages/apache_beam/runners/dataflow/dataflow_runner.py", line 141, in apply
return super().apply(transform, input, options)
File "/home/jupyter/project/training-data-analyst/quests/dataflow_python/7_Advanced_Streaming_Analytics/lab/df-env/lib/python3.7/site-packages/apache_beam/runners/runner.py", line 185, in apply
return m(transform, input, options)
File "/home/jupyter/project/training-data-analyst/quests/dataflow_python/7_Advanced_Streaming_Analytics/lab/df-env/lib/python3.7/site-packages/apache_beam/runners/runner.py", line 215, in apply_PTransform
return transform.expand(input)
File "/home/jupyter/project/training-data-analyst/quests/dataflow_python/7_Advanced_Streaming_Analytics/lab/df-env/lib/python3.7/site-packages/apache_beam/io/fileio.py", line 577, in expand
| beam.ParDo(
File "/home/jupyter/project/training-data-analyst/quests/dataflow_python/7_Advanced_Streaming_Analytics/lab/df-env/lib/python3.7/site-packages/apache_beam/pvalue.py", line 137, in __or__
return self.pipeline.apply(ptransform, self)
File "/home/jupyter/project/training-data-analyst/quests/dataflow_python/7_Advanced_Streaming_Analytics/lab/df-env/lib/python3.7/site-packages/apache_beam/pipeline.py", line 652, in apply
transform.transform, pvalueish, label or transform.label)
File "/home/jupyter/project/training-data-analyst/quests/dataflow_python/7_Advanced_Streaming_Analytics/lab/df-env/lib/python3.7/site-packages/apache_beam/pipeline.py", line 662, in apply
return self.apply(transform, pvalueish)
File "/home/jupyter/project/training-data-analyst/quests/dataflow_python/7_Advanced_Streaming_Analytics/lab/df-env/lib/python3.7/site-packages/apache_beam/pipeline.py", line 708, in apply
pvalueish_result = self.runner.apply(transform, pvalueish, self._options)
File "/home/jupyter/project/training-data-analyst/quests/dataflow_python/7_Advanced_Streaming_Analytics/lab/df-env/lib/python3.7/site-packages/apache_beam/runners/dataflow/dataflow_runner.py", line 141, in apply
return super().apply(transform, input, options)
File "/home/jupyter/project/training-data-analyst/quests/dataflow_python/7_Advanced_Streaming_Analytics/lab/df-env/lib/python3.7/site-packages/apache_beam/runners/runner.py", line 185, in apply
return m(transform, input, options)
File "/home/jupyter/project/training-data-analyst/quests/dataflow_python/7_Advanced_Streaming_Analytics/lab/df-env/lib/python3.7/site-packages/apache_beam/runners/dataflow/dataflow_runner.py", line 830, in apply_GroupByKey
return transform.expand(pcoll)
File "/home/jupyter/project/training-data-analyst/quests/dataflow_python/7_Advanced_Streaming_Analytics/lab/df-env/lib/python3.7/site-packages/apache_beam/transforms/core.py", line 2609, in expand
raise ValueError(msg)
ValueError: GroupRecordsByDestinationAndShard: Unsafe trigger: `AfterProcessingTime(delay=300)` may lose data. Reason: MAY_FINISH. This can be overriden with the --allow_unsafe_triggers flag.
`

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.