apache / apache/beam

[Feature Request]: documentation how to override connector logic to be reflected when executing pipeline

Open
#23,457 2 comments 0 reactions 0 assignees View on GitHub
io mongodb new feature P2 python
Dominant language
Java
Stars
8.7k
Forks
4.7k
Avg merge
1d 20h
Merged PRs (30d)
196

Description

### What would you like to happen?

Hi I am using mongodbio connector and needed to overwrite its logic of write https://github.com/apache/beam/blob/3c7a3d40ce12eef5f4d361c67f1286b487847f65/sdks/python/apache_beam/io/mongodbio.py#L792 to not use ReplaceOne but UpdateOne.

In the main function which has the pipeline logic I added sth like that bellow:

```python
beam.io.mongodbio._MongoSink.write = my_custom_function
```

and this is my pipeline logic:
```python
import apach_beam as beam

def custom_write(self, documents):
"""
beam.io.WriteToMongoDB._MongoSink update to merge data not override
:param self:
:param documents:
:return:
"""
if self.client is None:
self.client = MongoClient(host=self.uri, **self.spec)
requests = []
for doc in documents:
# match document based on _id field, if not found in current collection,
# insert new one, otherwise overwrite it.
requests.append(
UpdateOne(
filter={"_id": doc.get("_id", None)},
update={"$set": doc},
upsert=True))
resp = self.client[self.db][self.coll].bulk_write(requests)
beam.io.mongodbio._LOGGER.debug(
"BulkWrite to MongoDB result in nModified:%d, nUpserted:%d, "
"nMatched:%d, Errors:%s" % (
resp.modified_count,
resp.upserted_count,
resp.matched_count,
resp.bulk_api_result.get("writeErrors"),
))

beam.io.mongodbio._MongoSink.write = custom_write

def main(argv=None):
parser = argparse.ArgumentParser()

known_args, pipeline_args = parser.parse_known_args(argv)

pipeline_options = PipelineOptions(pipeline_args)

with beam.Pipeline(options=pipeline_options) as pipeline:

(pipeline
| f"QueryTable" >> beam.io.ReadFromBigQuery(
query="")
| f"TransformingData" >> beam.Map()
| f'WritingToMongo' >> beam.io.WriteToMongoDB(uri=uri,
db=db,
coll=coll,
batch_size=limit)
)

if __name__ == '__main__':
logging.getLogger().setLevel(logging.INFO)
main()
```

This works when I use DirectRunner but it doesn't work when using DataflowRunner

I am using my own docker image with dependencies already installed.

My question is how I can make a change to existing connector so it is reflected when executing pipeline ? Like how I can dynamically replace some logic or functions ?

### Issue Priority

Priority: 2

### Issue Component

Component: io-py-mongodb

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.