Flink消费Google PubSub触发DEADLINE_EXCEEDED异常求助
我正在开展一项POC,目标是从Google PubSub主题获取消息并通过Flink作业处理。目前已成功通过Flink的PubSubSource拉取到消息,但持续触发以下gRPC超时异常,导致作业不断重启,本地环境及正式Flink集群均存在该问题:
io.grpc.StatusRuntimeException: DEADLINE_EXCEEDED: deadline exceeded after 14.999285491s. [remote_addr=pubsub.googleapis.com/142.250.125.95: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 org.apache.flink.streaming.examples.gcp.pubsub.shaded.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:67) at org.apache.flink.streaming.runtime.tasks.SourceStreamTask$LegacySourceFunctionThread.run(SourceStreamTask.java:333)
针对该问题,可尝试以下几种解决方案:
延长gRPC拉取超时时间
Flink PubSub连接器默认的gRPC拉取超时(15秒)可能无法满足网络延迟较高或消息批量较大的场景。通过配置pubsub.subscriber.grpc.deadline.ms参数延长超时,例如设置为30秒:config.setString("pubsub.subscriber.grpc.deadline.ms", "30000");调整PubSub订阅的ACK超时
检查Google PubSub订阅的ackDeadline配置,若该值过短,可能导致拉取过程中未完成处理就触发超时。可在Google Cloud Console中将订阅的ackDeadline延长至60秒以上。优化作业资源与并行度
若作业并行度过高,会导致每个Source实例拉取压力过大;或资源(CPU、内存)不足导致处理缓慢,间接引发gRPC超时。适当降低并行度,确保每个TaskManager分配到足够的计算资源。限制递归重试次数
从栈trace可见BlockingGrpcPubSubSubscriber.pull存在多次递归调用,超时后的无限重试会加剧问题。通过pubsub.subscriber.max.retries参数限制重试次数:config.setString("pubsub.subscriber.max.retries", "3");升级Flink PubSub连接器版本
旧版本连接器可能存在gRPC超时处理的bug,建议升级至对应Flink版本的最新稳定连接器,新版本通常会修复这类超时相关问题。
内容的提问来源于stack exchange,提问作者Jaumzera

