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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 04:32:31