【问题标题】:Why Airflow PubSubPullOperator didn't pull max messages?为什么 Airflow PubSubPullOperator 没有提取最大消息?
【发布时间】:2022-07-14 22:14:17
【问题描述】:

我在气流中使用 PubSubPullOperator 从 gcp 订阅中提取消息。

pull_messages_task = PubSubPullOperator(
        task_id="pull_messages",
        ack_messages=True,
        project_id=GCP_PROJECT_ID,
        subscription="k8s-sub",
        gcp_conn_id=GCP_CONN_ID,
        max_messages=50
    )

从订阅中提取消息并保存在 Xcom 中工作正常。 我的问题是为什么 PubSubPullOperator 每次都无法提取等于 max_messages 的消息数?

例如,我向 GCP 主题发布了 250 条消息。 My Dag 每分钟运行一次,每次拉取 50 条消息。

以下是气流的过程日志:

[2022-05-17 14:53:04,630] {pubsub.py:536} INFO - Pulling max 50 messages from subscription (path) projects/production-1/subscriptions/k8s-sub
[2022-05-17 14:53:06,661] {pubsub.py:550} INFO - Pulled 50 messages from subscription (path) projects/production-1/subscriptions/k8s-sub

[2022-05-17 14:54:04,312] {pubsub.py:536} INFO - Pulling max 50 messages from subscription (path) projects/production-1/subscriptions/k8s-sub
[2022-05-17 14:54:06,239] {pubsub.py:550} INFO - Pulled 16 messages from subscription (path) projects/production-1/subscriptions/k8s-sub

[2022-05-17 14:55:04,055] {pubsub.py:536} INFO - Pulling max 50 messages from subscription (path) projects/production-1/subscriptions/k8s-sub
[2022-05-17 14:55:05,259] {pubsub.py:550} INFO - Pulled 4 messages from subscription (path) projects/production-1/subscriptions/k8s-sub

[2022-05-17 14:56:04,590] {pubsub.py:536} INFO - Pulling max 50 messages from subscription (path) projects/production-1/subscriptions/k8s-sub
[2022-05-17 14:56:06,527] {pubsub.py:550} INFO - Pulled 20 messages from subscription (path) projects/production-1/subscriptions/k8s-sub

[2022-05-17 14:57:04,083] {pubsub.py:536} INFO - Pulling max 50 messages from subscription (path) projects/production-1/subscriptions/k8s-sub
[2022-05-17 14:57:07,428] {pubsub.py:550} INFO - Pulled 38 messages from subscription (path) projects/production-1/subscriptions/k8s-sub

[2022-05-17 14:58:05,561] {pubsub.py:536} INFO - Pulling max 50 messages from subscription (path) projects/production-1/subscriptions/k8s-sub
[2022-05-17 14:58:07,431] {pubsub.py:550} INFO - Pulled 50 messages from subscription (path) projects/production-1/subscriptions/k8s-sub

[2022-05-17 14:59:04,348] {pubsub.py:536} INFO - Pulling max 50 messages from subscription (path) projects/production-1/subscriptions/k8s-sub
[2022-05-17 14:59:05,462] {pubsub.py:550} INFO - Pulled 50 messages from subscription (path) projects/production-1/subscriptions/k8s-sub

[2022-05-17 15:00:06,882] {pubsub.py:536} INFO - Pulling max 50 messages from subscription (path) projects/production-1/subscriptions/k8s-sub
[2022-05-17 15:00:08,710] {pubsub.py:550} INFO - Pulled 2 messages from subscription (path) projects/production-1/subscriptions/k8s-sub

[2022-05-17 15:01:03,519] {pubsub.py:536} INFO - Pulling max 50 messages from subscription (path) projects/production-1/subscriptions/k8s-sub
[2022-05-17 15:01:03,688] {pubsub.py:550} INFO - Pulled 20 messages from subscription (path) projects/production-1/subscriptions/k8s-sub

我很确定每个 dag 运行时间不到 1 分钟。并且 50 条消息大小不超过 Xcom 限制(48KB)。

有人对这种情况有任何想法吗?或者有谁知道 Operator 是如何决定要拉多少条消息的?

非常感谢。

【问题讨论】:

  • 这是使用PubSubPullOperator 的正常行为,因为此运算符是非阻塞任务。如果你想要每 50 条消息提取一次的东西,你可以使用 PubSubPullSensor
  • @JoseGutierrezPaliza 感谢您的回复。我将 PubSubPullOperator 更改为 PubSubPullSensor。但结果保持不变:(唯一不同的是,如果主题中没有消息 PubSubPullOperator 将通过但 PubSubPullSensor 将等待。

标签: google-cloud-platform airflow publish-subscribe


【解决方案1】:

此功能似乎源自 google 的 pubsub 客户端库,而不是气流运算符本身的功能/问题。

来自谷歌documentation

  The maximum number of messages to return for this request. Must be a positive integer. The Pub/Sub system may return fewer than the number specified.

运营商依赖 PubSubPullOperator 使用 PubSubHook 使用 SubscriberClient

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2012-04-01
    • 1970-01-01
    • 1970-01-01
    • 2011-02-17
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多