apache / apache/beam

[Bug]: Read from datastore with inequality filters

Open
#27,486 1 comment 2 reactions 0 assignees View on GitHub
awaiting triage bug dataflow io P2 python stale
Dominant language
Java
Stars
8.7k
Forks
4.7k
Avg merge
1d 20h
Merged PRs (30d)
196

Description

Apache Beam: 2.44.0
Runners: DataflowRunner, DirectRunner

I have this pipeline which reads events from datastore and I want to add a filter to the query.
```python
with Pipeline(options=self.options) as p:
read_event_query = Query(
kind=self.configs.datastore.events_kind,
project=self.configs.cloud.project_id,
namespace=self.configs.datastore.namespace,
)
start_timestamp = int(round(time.time() * 1000) - 6.048e+8)
read_event_query.filters = [('timestamp' , '>', start_timestamp)]

events = (
p
| 'Read from datastore' >> ReadFromDatastore(query=read_event_query, num_splits=1)
)

_ = (events
| Map(print))
```

It outputs this error:
```
2023-07-13 12:12:27,713: INFO: Unable to parallelize the given query: ', 1688634747019)],projection=(), order=(), distinct_on=(), limit=None)>
Traceback (most recent call last):
File "/home/kgiantsios/miniconda3/envs/*****/lib/python3.10/site-packages/apache_beam/io/gcp/datastore/v1new/datastoreio.py", line 184, in process
query_splitter.validate_split(query)
File "/home/kgiantsios/miniconda3/envs/****/lib/python3.10/site-packages/apache_beam/io/gcp/datastore/v1new/query_splitter.py", line 109, in validate_split
raise SplitNotPossibleError('Query cannot have any inequality filters.')
apache_beam.io.gcp.datastore.v1new.query_splitter.SplitNotPossibleError: Query cannot have any inequality filters.

```

Although in Docstring of ReadFromDatastore notes:
```
However, when the `query` is configured with a `limit` or if the
query contains inequality filters like `GREATER_THAN, LESS_THAN` etc., then
all the returned results will be read by a single worker in order to ensure
correct data. Since data is read from a single worker, this could have
significant impact on the performance of the job.
```

Shouldn't the operation proceed even with the limitation of the single worker?

### Issue Priority

Priority: 1 (data loss / total loss of function)

### Issue Components

- [X] Component: Python SDK
- [ ] Component: Java SDK
- [ ] Component: Go SDK
- [ ] Component: Typescript SDK
- [X] Component: IO connector
- [ ] Component: Beam examples
- [ ] Component: Beam playground
- [ ] Component: Beam katas
- [ ] Component: Website
- [ ] Component: Spark Runner
- [ ] Component: Flink Runner
- [ ] Component: Samza Runner
- [ ] Component: Twister2 Runner
- [ ] Component: Hazelcast Jet Runner
- [X] Component: Google Cloud Dataflow Runner

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.