使用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
相关产品推荐
相关产品推荐

