googleapis / googleapis/google-cloud-python

PubSub: support batching publish requests with asyncio

未关闭
#15,654 7 条评论 5 个 reaction 已指派 1 人 已被 @abbrowne126 认领 在 GitHub 查看
api: pubsub type: feature request
主要语言
Python
星标
5.4k
派生
1.8k
平均合并
2 天 23 小时
30 天内合并 PR
123

描述

**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.

贡献指南

打开贡献指南

评估

这个 Issue 还没有评估数据。

把新 issue 发到你的邮箱

精选适合新手参与的 GitHub issue 摘要。