如何处理Apache Camel拉取Google Pubsub数据时的超时异常?
Apache Camel Google Pubsub 拉取数据时的DEADLINE_EXCEEDED错误处理方案
我开发了一个基于Apache Camel的项目,用于拉取Google Pubsub主题的数据,当前路由配置如下:
from("google-pubsub:xProjectx:xTopicx?maxMessagesPerPoll=20&concurrentConsumers=10&ackMode=AUTO") .doTry() .to("direct:someRoute") .endDoTry() .doCatch(Exception.class) .to("direct:errorHandle") .end();
运行过程中偶尔会收到Google Pubsub的超时错误,且现有错误处理机制无法捕获该异常,错误栈信息如下:
com.google.api.gax.rpc.DeadlineExceededException: io.grpc.StatusRuntimeException: DEADLINE_EXCEEDED: deadline exceeded after 40.786965000s. [closed=[], open=[[buffered_nanos=362483400, remote_addr=pubsub.googleapis.com/216.58.212.10:443]]] at com.google.api.gax.rpc.ApiExceptionFactory.createException(ApiExceptionFactory.java:94) ~[gax-2.18.1.jar!/:2.18.1] at com.google.api.gax.rpc.ApiExceptionFactory.createException(ApiExceptionFactory.java:41) ~[gax-2.18.1.jar!/:2.18.1] at com.google.api.gax.grpc.GrpcApiExceptionFactory.create(GrpcApiExceptionFactory.java:86) ~[gax-grpc-2.18.1.jar!/:2.18.1] at com.google.api.gax.grpc.GrpcApiExceptionFactory.create(GrpcApiExceptionFactory.java:66) ~[gax-grpc-2.18.1.jar!/:2.18.1] at com.google.api.gax.grpc.GrpcExceptionCallable$ExceptionTransformingFuture.onFailure(GrpcExceptionCallable.java:97) ~[gax-grpc-2.18.1.jar!/:2.18.1] at com.google.api.core.ApiFutures$1.onFailure(ApiFutures.java:67) ~[api-common-2.2.0.jar!/:na] at com.google.common.util.concurrent.Futures$4.run(Futures.java:1123) ~[guava-20.0.jar!/:na] at com.google.common.util.concurrent.MoreExecutors$DirectExecutor.execute(MoreExecutors.java:435) ~[guava-20.0.jar!/:na] at com.google.common.util.concurrent.AbstractFuture.executeListener(AbstractFuture.java:900) ~[guava-20.0.jar!/:na] at com.google.common.util.concurrent.AbstractFuture.complete(AbstractFuture.java:811) ~[guava-20.0.jar!/:na] at com.google.common.util.concurrent.AbstractFuture.setException(AbstractFuture.java:675) ~[guava-20.0.jar!/:na] at io.grpc.stub.ClientCalls$GrpcFuture.setException(ClientCalls.java:572) ~[grpc-stub-1.46.0.jar!/:1.46.0] at io.grpc.stub.ClientCalls$UnaryStreamToFuture.onClose(ClientCalls.java:542) ~[grpc-stub-1.46.0.jar!/:1.46.0] at io.grpc.PartialForwardingClientCallListener.onClose(PartialForwardingClientCallListener.java:39) ~[grpc-api-1.46.0.jar!/:1.46.0] at io.grpc.ForwardingClientCallListener.onClose(ForwardingClientCallListener.java:23) ~[grpc-api-1.46.0.jar!/:1.46.0] at io.grpc.ForwardingClientCallListener$SimpleForwardingClientCallListener.onClose(ForwardingClientCallListener.java:40) ~[grpc-api-1.46.0.jar!/:1.46.0] at com.google.api.gax.grpc.ChannelPool$ReleasingClientCall$1.onClose(ChannelPool.java:535) ~[gax-grpc-2.18.1.jar!/:2.18.1] at io.grpc.internal.ClientCallImpl.closeObserver(ClientCallImpl.java:562) ~[grpc-core-1.46.0.jar!/:1.46.0] at io.grpc.internal.ClientCallImpl.access$300(ClientCallImpl.java:70) ~[grpc-core-1.46.0.jar!/:1.46.0] at io.grpc.internal.ClientCallImpl$ClientStreamListenerImpl$1StreamClosed.runInternal(ClientCallImpl.java:743) ~[grpc-core-1.46.0.jar!/:1.46.0] at io.grpc.internal.ClientCallImpl$ClientStreamListenerImpl$1StreamClosed.runInContext(ClientCallImpl.java:722) ~[grpc-core-1.46.0.jar!/:1.46.0] at io.grpc.internal.ContextRunnable.run(ContextRunnable.java:37) ~[grpc-core-1.46.0.jar!/:1.46.0] at io.grpc.internal.SerializingExecutor.run(SerializingExecutor.java:133) ~[grpc-core-1.46.0.jar!/:1.46.0] at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:539) ~[na:na] at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) ~[na:na] at java.base/java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:304) ~[na:na] at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136) ~[na:na] at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635) ~[na:na] at java.base/java.lang.Thread.run(Thread.java:833) ~[na:na] Caused by: io.grpc.StatusRuntimeException: DEADLINE_EXCEEDED: deadline exceeded after 40.786965000s. [closed=[], open=[[buffered_nanos=362483400, remote_addr=pubsub.googleapis.com/216.58.212.10:443]]] at io.grpc.Status.asRuntimeException(Status.java:535) ~[grpc-api-1.46.0.jar!/:1.46.0] ... 17 common frames omitted
我已经尝试过使用bridgeErrorHandler、exceptionHandler参数,以及onException代码块,但均无法处理该错误。以下是可行的解决建议:
可行解决建议
1. 调整Pubsub端点的gRPC超时与重试参数
这个超时是gRPC层面的请求超时,默认配置可能无法满足网络或负载需求,可直接在端点URL中添加参数:
pullEndpointTimeout:设置拉取请求的超时时间(例如调整为60000毫秒,即60秒)retrySettings:配置针对超时异常的重试策略,示例:from("google-pubsub:xProjectx:xTopicx?maxMessagesPerPoll=20&concurrentConsumers=10&ackMode=AUTO&pullEndpointTimeout=60000&retrySettings=maxAttempts=5,initialRetryDelay=1000,retryDelayMultiplier=2.0") .to("direct:someRoute");
2. 配置组件级别的异常处理器
该DeadlineExceededException是在Camel Pubsub组件的消费线程中抛出的,路由级别的错误处理器无法捕获。需要为Pubsub组件单独设置异常处理器:
// 获取Pubsub组件实例 GooglePubsubComponent pubsubComponent = context.getComponent("google-pubsub", GooglePubsubComponent.class); // 设置自定义异常处理器 pubsubComponent.setExceptionHandler(new ExceptionHandler() { @Override public void handleException(Throwable exception, Exchange exchange) { if (exception instanceof DeadlineExceededException || exception instanceof StatusRuntimeException) { // 执行错误处理逻辑:记录日志、触发告警、重试等 log.error("Pubsub拉取超时,错误信息:{}", exception.getMessage(), exception); // 可根据需求决定是否标记Exchange为已处理 exchange.setException(null); } } });
3. 优化并发与拉取负载参数
当前配置的concurrentConsumers=10和maxMessagesPerPoll=20可能导致并发请求过多,超出Pubsub的处理能力,引发超时。可尝试:
- 降低
concurrentConsumers数量至5左右 - 减少
maxMessagesPerPoll为10,减轻单次拉取的消息负载
4. 精准指定全局异常捕获类型
避免泛用Exception.class捕获,明确指定要处理的异常类型,并标记为已处理:
// 全局异常捕获配置 onException(DeadlineExceededException.class, StatusRuntimeException.class) .handled(true) // 标记异常已处理,阻止向上传播 .log("捕获Pubsub超时异常:${exception.message}") .to("direct:errorHandle"); // 路由定义 from("google-pubsub:xProjectx:xTopicx?maxMessagesPerPoll=20&concurrentConsumers=10&ackMode=AUTO") .to("direct:someRoute");
5. 排查网络与Pubsub配额问题
- 检查服务器到
pubsub.googleapis.com的网络连通性,确认无防火墙、代理导致的延迟 - 登录Google Cloud控制台,查看Pubsub的配额使用情况,确认是否达到请求速率限制或资源配额上限
内容的提问来源于stack exchange,提问作者Sedat Göz
相关产品推荐
相关产品推荐

