Apache Camel与Google PubSub集成求助:消息收发及错误处理实现
Apache Camel 集成 Google PubSub 实现指南
一、环境依赖配置
首先在Maven pom.xml中添加必要依赖:
<dependencies> <!-- Camel Google PubSub 组件 --> <dependency> <groupId>org.apache.camel</groupId> <artifactId>camel-google-pubsub</artifactId> <version>3.20.2</version> <!-- 替换为对应Camel稳定版本 --> </dependency> <!-- GCP PubSub 客户端依赖 --> <dependency> <groupId>com.google.cloud</groupId> <artifactId>google-cloud-pubsub</artifactId> <version>1.123.11</version> </dependency> </dependencies>
GCP认证配置:
- 下载Service Account密钥文件,设置环境变量:
export GOOGLE_APPLICATION_CREDENTIALS="/path/to/your-key.json" - 也可在Camel路由中通过
credentialsLocation参数直接指定密钥路径。
二、核心功能实现
1. 消息发布到PubSub主题
创建Camel路由,将指定负载发送到目标主题,示例:
import org.apache.camel.builder.RouteBuilder; public class PubSubPublishRoute extends RouteBuilder { @Override public void configure() throws Exception { // 可替换为实际业务数据源(如HTTP接口、数据库事件等) from("direct:publish-message") .setBody(constant("{\"orderId\":\"1001\",\"status\":\"created\"}")) // 你的指定消息负载 .to("google-pubsub:your-gcp-project:your-topic-name") .log("消息已发布到主题: ${body}"); } }
组件URL格式:google-pubsub:gcp-project-id:topic-name,支持添加自定义客户端参数(如?credentialsLocation=/path/to/key.json)。
2. 订阅接收+队列存储
用Camelseda组件实现本地内存队列(可替换为ActiveMQ等持久化队列),先将PubSub订阅消息存入队列再异步处理:
import org.apache.camel.builder.RouteBuilder; public class PubSubSubscribeRoute extends RouteBuilder { @Override public void configure() throws Exception { // 从PubSub订阅拉取消息,存入本地队列 from("google-pubsub:your-gcp-project:your-subscription-name?maxMessages=10&pullInterval=1000") .log("接收到PubSub消息: ${body}") .to("seda:message-processing-queue?concurrentConsumers=3"); // 处理队列中的消息 from("seda:message-processing-queue") .process(exchange -> { String message = exchange.getIn().getBody(String.class); // 这里添加业务逻辑(如解析JSON、写入数据库等) System.out.println("处理消息: " + message); }) .log("消息处理完成"); } }
3. 错误处理实现
利用CamelonException机制结合PubSub死信主题,处理消息消费失败场景:
import org.apache.camel.builder.RouteBuilder; import org.apache.camel.model.dataformat.JsonLibrary; public class PubSubErrorHandlingRoute extends RouteBuilder { @Override public void configure() throws Exception { // 全局异常处理规则 onException(Exception.class) .maximumRedeliveries(3) // 重试3次 .redeliveryDelay(2000) // 重试间隔2秒 .log("消息重试失败,转入死信主题: ${body}") .to("google-pubsub:your-gcp-project:your-dead-letter-topic"); // 订阅消息并处理 from("google-pubsub:your-gcp-project:your-subscription-name") .unmarshal().json(JsonLibrary.Jackson) // 解析JSON负载 .process(exchange -> { // 模拟业务异常 String orderStatus = exchange.getIn().getBody(String.class); if ("error".equals(orderStatus)) { throw new RuntimeException("业务处理失败"); } }) .log("消息处理成功"); } }
也可将失败消息存储到本地文件用于后续排查:
onException(Exception.class) .log("消息处理失败,存储到本地") .to("file:/path/to/error-logs?fileName=error-${date:now:yyyyMMddHHmmssSSS}.txt");
三、启动项目
创建主类启动Camel上下文:
import org.apache.camel.main.Main; public class PubSubIntegrationApp { public static void main(String[] args) throws Exception { Main main = new Main(); main.addRouteBuilder(new PubSubPublishRoute()); main.addRouteBuilder(new PubSubSubscribeRoute()); main.addRouteBuilder(new PubSubErrorHandlingRoute()); main.run(args); } }
内容的提问来源于stack exchange,提问作者kilobyte
相关产品推荐
相关产品推荐

