【发布时间】:2021-04-21 23:13:23
【问题描述】:
我在 python3.6 中使用 Google Pub/Sub 客户端 v2.2.0 作为订阅者。
我希望我的应用程序在确认它已收到的所有消息后正常关闭。
来自 Google 指南的订阅者的示例代码,稍作更改将显示我的问题:
from concurrent.futures import TimeoutError
from google.cloud import pubsub_v1
from time import sleep
# TODO(developer)
# project_id = "your-project-id"
# subscription_id = "your-subscription-id"
# Number of seconds the subscriber should listen for messages
# timeout = 5.0
subscriber = pubsub_v1.SubscriberClient()
# The `subscription_path` method creates a fully qualified identifier
# in the form `projects/{project_id}/subscriptions/{subscription_id}`
subscription_path = subscriber.subscription_path(project_id, subscription_id)
def callback(message):
print(f"Received {message}.")
sleep(30)
message.ack()
print("Acked")
streaming_pull_future = subscriber.subscribe(subscription_path, callback=callback)
print(f"Listening for messages on {subscription_path}..\n")
# Wrap subscriber in a 'with' block to automatically call close() when done.
with subscriber:
sleep(10)
streaming_pull_future.cancel()
streaming_pull_future.result()
来自https://cloud.google.com/pubsub/docs/pull
我希望这段代码停止拉取消息并完成正在运行的消息然后退出。
实际上,这段代码停止拉取消息并完成执行正在运行的消息,但它不确认消息。 .ack() 发生但服务器没有收到 ack,所以下次运行相同的消息再次返回。
1.为什么服务端收不到ack?
2。如何优雅地关闭订阅者?
3. .cancel() 的预期行为是什么?
【问题讨论】:
-
我看了一下库,停止进程(取消)等待所有线程结束。我想到了别的东西:您的订阅确认截止日期是什么时候?
-
@guillaumeblaquiere 我的确认截止日期是默认的 600 秒
-
@JohnHanley 即使睡了 60 秒,确认仍然没有发生。
-
SIGTERM 发生在一个更复杂的代码中,所以我做了一个没有它的简单示例。
-
在我的实际应用程序中,我使用 sigterm 处理程序来调用 .cancel()。在这里,使用没有 sigterm(处理或调用)的更简单的代码,我观察到取消后未确认的消息的相同行为。在问题中写 sigterm 令人困惑,我将其删除。
标签: python google-cloud-platform publish-subscribe google-cloud-pubsub pypubsub