【发布时间】:2018-04-14 08:33:16
【问题描述】:
在我的场景中,我正在使用 PubSub 安排任务。这是多达 2.000 条 PubSub 消息,这些消息由在 Google Compute Engine 内的 Docker 容器内运行的 Python 脚本使用。该脚本使用 PubSub 消息。
每条消息的处理时间约为 30 秒到 5 分钟。因此,确认截止日期为 600 秒(10 分钟)。
from google.cloud import pubsub_v1
from google.cloud.pubsub_v1.subscriber.message import Message
def handle_message(message: Message):
# do your stuff here (max. 600sec)
message.ack()
return
def receive_messages(project, subscription_name):
subscriber = pubsub_v1.SubscriberClient()
subscription_path = subscriber.subscription_path(project, subscription_name)
flow_control = pubsub_v1.types.FlowControl(max_messages=5)
subscription = subscriber.subscribe(subscription_path, flow_control=flow_control)
future = subscription.open(handle_message)
# Blocks the thread while messages are coming in through the stream. Any
# exceptions that crop up on the thread will be set on the future.
# The subscriber is non-blocking, so we must keep the main thread from
# exiting to allow it to process messages in the background.
try:
future.result()
except Exception as e:
# do some logging
raise
因为我正在处理这么多 PubSub 消息,所以我是 creating a template 的计算引擎,它以这两种方式之一使用 auto-scaling:
gcloud compute instance-groups managed create my-worker-group \
--zone=europe-west3-a \
--template=my-worker-template \
--size=0
gcloud beta compute instance-groups managed set-autoscaling my-worker-group \
--zone=europe-west3-a \
--max-num-replicas=50 \
--min-num-replicas=0 \
--target-cpu-utilization=0.4
gcloud beta compute instance-groups managed set-autoscaling my-worker-group \
--zone=europe-west3-a \
--max-num-replicas=50 \
--min-num-replicas=0 \
--update-stackdriver-metric=pubsub.googleapis.com/subscription/num_undelivered_messages \
--stackdriver-metric-filter="resource.type = pubsub_subscription AND resource.label.subscription_id = my-pubsub-subscription" \
--stackdriver-metric-single-instance-assignment=10
到目前为止,一切都很好。选项一可扩展到大约 8 个实例,而第二个选项将启动最大数量的实例。现在我发现发生了一些奇怪的事情,这就是我在这里发帖的原因。也许你能帮帮我?!
消息重复:似乎每个实例中的 PubSub 服务(计算引擎中 docker 容器内的 Python 脚本)读取一批消息(约 10 条),有点像缓冲区并提供给它们到我的代码。看起来所有同时启动的实例都将读取所有相同的消息(2.000 条中的前 10 条),并将开始处理相同的内容。在我的日志中,我看到大多数消息由不同的机器处理 3 次。我期待 PubSub 知道某个订阅者是否缓冲了 10 条消息,以便另一个订阅者缓冲 10 条不同的消息而不是相同的消息。
确认截止日期:由于缓冲进入缓冲区末尾的消息(比如说消息 8 或 9)必须在缓冲区中等待,直到前面的消息(消息 1 到 7 ) 已处理。该等待时间加上它自己的处理时间之和可能会达到 600 秒的超时时间。
负载平衡:因为每台机器都缓存了这么多消息,负载只被少数实例消耗,而其他实例完全空闲。这发生在使用 PubSub 堆栈驱动程序指标的缩放选项 2 中。
人们告诉我,我需要使用 Cloud SQL 或其他方式实现手动同步服务,其中每个实例都指示它在哪个消息上工作,这样其他实例就不会以相同的方式启动。但我觉得这不可能是真的——因为那时我不明白 PubSub 到底是什么。
更新:我找到了nice explanation by Gregor Hohpe,他是 2015 年的《企业集成模式》一书的合著者。实际上我的观察是错误的,但观察到的副作用是真实的。
Google Cloud Pub/Sub API 实际上同时实现了 发布-订阅频道和竞争消费者模式。在 Cloud Pub/Sub 的核心是一个经典的发布-订阅频道, 它将发布给它的单个消息传递给多个 订户。这种模式的一个优点是添加订阅者 没有副作用,这是发布-订阅频道的原因之一 有时被认为比点对点耦合更松散 通道,将消息传递给一个订阅者。添加 消费者到点对点渠道导致竞争消费者 因此有很强的副作用。
我观察到的副作用是关于每个订阅者(订阅相同订阅,点对点 == 竞争消费者)中的消息缓冲和消息流控制。当前版本的 Python Client Lib 包装了 PubSub REST API(和 RPC)。如果使用该包装器,则无法控制:
- 在一个虚拟机上启动了多少个容器;如果 CPU 尚未充分利用,可能会启动多个容器
- 一次从订阅中拉出多少条消息(缓冲);完全无法控制
- 在容器内部启动了多少线程来处理拉取的消息;如果值低于固定值,则 flow_control(max_messages) 无效。
我们观察到的副作用是:
- 一个消费者一次提取大量消息(大约 100 到 1.000 条)并将它们排队到其客户端缓冲区中。因此所有其他根据自动缩放规则启动的虚拟机都不会收到任何消息,因为所有消息都在前几个虚拟机的队列中
- 如果消息在确认截止日期前运行,则会将消息重新传递到同一 VM 或任何其他 VM(或 docker 容器)。因此,您需要在处理消息时修改确认截止日期。处理开始时,截止日期计数器开始计时。
- 假设消息的处理是一项长时间运行的任务(例如机器学习),您可以
- 预先确认消息,但如果没有更多消息等待,这将导致 VM 被自动缩放规则关闭。该规则不关心 CPU 利用率是否仍然很高并且处理尚未完成。
- 处理后确认消息。在这种情况下,您需要在处理该消息时修改该特定消息的确认截止日期。自上次修改以来,不得有一个代码块违反了截止日期。
尚未研究的可能解决方案:
- 使用 Java 客户端库,因为它可以更好地控制拉取和使用消息
- 使用 Python 客户端库的底层 API 调用和类
- 构建协调竞争消费者的同步存储
【问题讨论】:
-
我不知道你为什么要处理消息拉取调用
open()方法而不是使用subscriber.subscribe()作为唯一的未来,传递给它你想用每条消息调用的回调函数.该文档提供了一个特定的example of pull subscription usingFlowControl,它不需要使用open()。然后,您可以指定max_messages的下限以提高吞吐量。 -
好点。我会试试的。我认为
open()是为了阻塞主线程。但是future.result()无论如何都会阻止。 -
另外,我知道
result()方法会在消息进入时阻止执行,而实际上您希望订阅者是非阻塞的,以便它可以在后台处理消息,所以我也会删除那部分代码。您使用 (1)open()和 (2)result()有什么具体原因吗?我认为这个问题可能与您如何在回调函数中处理消息处理的后台执行有关,这应该是非阻塞的。最后,您使用的是哪个版本的库?最新的是0.34.0。 -
我不确定
open()方法的作用,因为我找不到任何参考或文档,但从您的用例中我知道您需要非阻塞调用。如果我误解了这种情况,请纠正我。 -
我们使用的是最新版本。今天刚更新。但看起来 API 最近发生了变化。在以前的版本之一中,没有办法给
subscribe方法提供回调,所以你had to useopen。不过我会检查新的 API。future.result()用于将任何全局异常捕获为described here。
标签: python docker google-cloud-platform google-compute-engine google-cloud-pubsub