You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用GCP PubSub源的Flink作业反复出现DEADLINE_EXCEEDED错误求助

问题描述

使用GCP Pub/Sub作为Flink作业的数据源时,当订阅中没有消息,每隔15秒会重复抛出DEADLINE_EXCEEDED错误,导致作业从RUNNING切换为FAILED。具体错误日志如下:

Source: pub-sub-source -> filter -> (Sink: Docstore Sink, Map -> Sink: pinot-kafka-sink) (2/2)#3 (b7790c8ab117377fb8d85b1af23b1d11) switched from RUNNING to FAILED.
io.grpc.StatusRuntimeException: DEADLINE_EXCEEDED: deadline exceeded after 14.999620292s. [remote_addr=pubsub.googleapis.com/123.123.12.12:443]
    at io.grpc.stub.ClientCalls.toStatusRuntimeException(ClientCalls.java:262)
    at io.grpc.stub.ClientCalls.getUnchecked(ClientCalls.java:243)
    at io.grpc.stub.ClientCalls.blockingUnaryCall(ClientCalls.java:156)
    at com.google.pubsub.v1.SubscriberGrpc$SubscriberBlockingStub.pull(SubscriberGrpc.java:1641)
    at org.apache.flink.streaming.connectors.gcp.pubsub.BlockingGrpcPubSubSubscriber.pull(BlockingGrpcPubSubSubscriber.java:73)
    at org.apache.flink.streaming.connectors.gcp.pubsub.BlockingGrpcPubSubSubscriber.pull(BlockingGrpcPubSubSubscriber.java:77)
    at org.apache.flink.streaming.connectors.gcp.pubsub.BlockingGrpcPubSubSubscriber.pull(BlockingGrpcPubSubSubscriber.java:77)
    at org.apache.flink.streaming.connectors.gcp.pubsub.BlockingGrpcPubSubSubscriber.pull(BlockingGrpcPubSubSubscriber.java:77)
    at org.apache.flink.streaming.connectors.gcp.pubsub.BlockingGrpcPubSubSubscriber.pull(BlockingGrpcPubSubSubscriber.java:67)
    at org.apache.flink.streaming.connectors.gcp.pubsub.PubSubSource.run(PubSubSource.java:128)
    at org.apache.flink.streaming.api.operators.StreamSource.run(StreamSource.java:110)
    at org.apache.flink.streaming.api.operators.StreamSource.run(StreamSource.java:66)
    at org.apache.flink.streaming.runtime.tasks.SourceStreamTask$LegacySourceFunctionThread.run(StreamSourceTask.java:267)

当前初始化PubSubSource的代码为:

PubSubSource pubSubSource =
        PubSubSource.newBuilder()
            .withDeserializationSchema(new SomeDeserializer())
            .withProjectName("some-project")
            .withSubscriptionName("some-subscription")
            .withCredentials(GoogleCredentials.fromStream(file))
            .build();

已尝试将检查点间隔调整为小于Pub/Sub订阅的确认超时,但问题仍存在。

解决方案

针对订阅无消息时的超时问题,需要配置以下参数:

  • 延长Pull请求超时时间:默认gRPC Pull请求的超时时间为15秒,恰好对应报错间隔。通过withPullTimeout调长超时时间(如30秒),让请求在无消息时能等待更久再返回,避免触发超时异常:

    .withPullTimeout(Duration.ofSeconds(30))
    
  • 增加Pull请求重试次数:默认重试逻辑在超时后直接抛出异常终止作业,通过withMaxRetries设置更大的重试次数,让连接器在超时后重试而非直接失败:

    .withMaxRetries(5)
    
  • 使用异步流式Pull替代阻塞式Pull:如果Flink版本支持,通过withSubscriberFactory指定AsyncGrpcPubSubSubscriber,它会保持长连接监听消息,避免频繁短连接的超时问题:

    .withSubscriberFactory(AsyncGrpcPubSubSubscriber::new)
    
  • 设置无消息时的Pull间隔:通过withIdleBetweenPulls控制无消息时两次Pull请求的间隔,避免过于频繁发起请求导致超时:

    .withIdleBetweenPulls(Duration.ofSeconds(5))
    
调整后的示例代码
PubSubSource pubSubSource =
        PubSubSource.newBuilder()
            .withDeserializationSchema(new SomeDeserializer())
            .withProjectName("some-project")
            .withSubscriptionName("some-subscription")
            .withCredentials(GoogleCredentials.fromStream(file))
            .withPullTimeout(Duration.ofSeconds(30))
            .withMaxRetries(5)
            .withIdleBetweenPulls(Duration.ofSeconds(5))
            .withSubscriberFactory(AsyncGrpcPubSubSubscriber::new)
            .build();

内容的提问来源于stack exchange,提问作者Amrith M

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.06 10:44:56