apache / apache/beam

Error running Top transform in Dataflow

Open
#21,229 2 comments 0 reactions 0 assignees View on GitHub
bug core dataflow P3 python runners
Dominant language
Java
Stars
8.7k
Forks
4.7k
Avg merge
1d 20h
Merged PRs (30d)
196

Description

When running the following transform in Dataflow (the problem does not happen with the direct running)... I is a streaming pipeline where I am using a SlidingWindow.
```

beam.combiners.Top.Of(n=10, key=lambda item: item[1]).without_defaults() 
```

 
I am getting an error
 
```

Error message from worker: generic::unknown: Traceback (most recent call last): File "/usr/local/lib/python3.7/site-packages/apache_beam/runners/worker/sdk_worker.py",
line 284, in _execute response = task() File "/usr/local/lib/python3.7/site-packages/apache_beam/runners/worker/sdk_worker.py",
line 357, in lambda: self.create_worker().do_instruction(request), request) File "/usr/local/lib/python3.7/site-packages/apache_beam/runners/worker/sdk_worker.py",
line 602, in do_instruction getattr(request, request_type), request.instruction_id) File "/usr/local/lib/python3.7/site-packages/apache_beam/runners/worker/sdk_worker.py",
line 639, in process_bundle bundle_processor.process_bundle(instruction_id)) File "/usr/local/lib/python3.7/site-packages/apache_beam/runners/worker/bundle_processor.py",
line 994, in process_bundle element.data) File "/usr/local/lib/python3.7/site-packages/apache_beam/runners/worker/bundle_processor.py",
line 222, in process_encoded self.output(decoded_value) File "apache_beam/runners/worker/operations.py",
line 351, in apache_beam.runners.worker.operations.Operation.output File "apache_beam/runners/worker/operations.py",
line 353, in apache_beam.runners.worker.operations.Operation.output File "apache_beam/runners/worker/operations.py",
line 215, in apache_beam.runners.worker.operations.SingletonConsumerSet.receive File "apache_beam/runners/worker/operations.py",
line 921, in apache_beam.runners.worker.operations.CombineOperation.process File "apache_beam/runners/worker/operations.py",
line 925, in apache_beam.runners.worker.operations.CombineOperation.process File "/usr/local/lib/python3.7/site-packages/apache_beam/transforms/combiners.py",
line 835, in extract_only return self.combine_fn.extract_output(accumulator) File "/usr/local/lib/python3.7/site-packages/apache_beam/transforms/combiners.py",
line 502, in extract_output heap.sort(reverse=True) File "apache_beam/transforms/cy_combiners.py", line
389, in apache_beam.transforms.cy_combiners.ComparableValue.__lt__ AssertionError passed through: ==>
dist_proc/dax/workflow/worker/fnapi_service_impl.cc:644
```

 

`What might be the cause of this issue?`

 

Imported from Jira [BEAM-12847](https://issues.apache.org/jira/browse/BEAM-12847). Original Jira may contain additional context.
Reported by: mirene.

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.