googleapis / googleapis/google-cloud-python
PubSub: support batching publish requests with asyncio
- Vorherrschende Sprache
- Python
- Sterne
- 5.4k
- Forks
- 1.8k
- Ø Merge
- 3 T. 4 Std.
- Gemergte PRs (30 T.)
- 122
Beschreibung
**Is your feature request related to a problem? Please describe.**
I have an `asyncio` application that needs to publish messages to PubSub, but I'm having issues because `google.cloud.pubsub.PublisherClient.publish`:
1. returns futures that aren't compatible with `await` or `asyncio.wrap_future`
1. returns futures that never complete if `Batch._commit` throws an uncaught exception (like in googleapis/google-cloud-python#7103 and googleapis/google-cloud-python#7071)
1. doesn't enforce a maximum number of threads, which is eating memory
**Describe the solution you'd like**
I wrote a new `google.cloud.pubsub_v1.publisher._batch.async.Batch` that implements `google.cloud.pubsub_v1.publisher._batch.base.Batch`. It uses `asyncio` to provide awaitable futures that automatically propagate exceptions. It uses a shared `concurrent.futures.ThreadPoolExecutor` in conjunction with `asyncio.wrap_future` to asynchronously call `Batch.client.publish` while enforcing a maximum number of workers. I specifically only wrapped `Batch.client.publish` in a thread because (if i understand correctly) it only blocks on exclusive access to the grpc channel, so it shouldn't create performance issues as seen in the first alternative below.
I would like to submit this as a pull request, but only if it would be useful.
**Describe alternatives you've considered**
* I tried patching `google.cloud.pubsub_v1.publisher._batch.thread.Batch` to use `concurrent.futures.ThreadPoolExecutor`. Unfortunately it had performance issues when all workers would reach a `time.sleep` and there wouldn't be any workers to check that not yet submitted tasks could be ready.
* I tried patching `google.cloud.pubsub_v1.futures.Future` to inherit from `concurrent.futures.Future`. This fixed compatiblity with `asyncio.wrap_future`, but not uncaught exceptions and unlimited thread spawning.
* I tried to patch `google.cloud.pubsub_v1.publisher._batch.thread.Batch` to join spawned threads, which would propagate uncaught exceptions, but I was unable to figure out a solution.
Beitragsleitfaden
Bewertung
Dieses Issue wurde noch nicht bewertet.