如何在Apache Camel中结合AVRO协议声明定义AVRO路由?
没问题,我来一步步教你怎么在Apache Camel里结合你提供的AVRO协议定义路由。先帮你把协议整理了下(补全了get接口的响应类型,看起来是输入时的小笔误):
你的AVRO协议(整理后):
{ "namespace": "org.apache.camel.avro.generated", "protocol": "KeyValueProtocol", "types": [ { "name": "Key", "type": "record", "fields": [{"name": "key", "type": "string"}] }, { "name": "Value", "type": "record", "fields": [{"name": "value", "type": "string"}] } ], "messages": { "put": { "request": [{"name": "key", "type": "Key"}, {"name": "value", "type": "Value"}], "response": "null" }, "get": { "request": [{"name": "key", "type": "Key"}], "response": "Value" } } }
接下来分三个核心步骤实现路由:
步骤1:生成AVRO Java类
Camel需要借助AVRO生成的Java类来序列化/反序列化消息,推荐用Maven插件自动生成:
在你的pom.xml中添加AVRO插件配置:
<build> <plugins> <plugin> <groupId>org.apache.avro</groupId> <artifactId>avro-maven-plugin</artifactId> <version>1.11.3</version> <!-- 使用和你项目兼容的最新版本 --> <executions> <execution> <phase>generate-sources</phase> <goals> <goal>protocol</goal> </goals> <configuration> <sourceDirectory>${project.basedir}/src/main/resources/avro</sourceDirectory> <outputDirectory>${project.build.directory}/generated-sources/avro</outputDirectory> </configuration> </execution> </executions> </plugin> </plugins> </build>
把你的AVRO协议文件(命名为keyvalue.avpr)放到src/main/resources/avro目录,执行mvn generate-sources就能生成org.apache.camel.avro.generated包下的所有Java类(包括Key、Value、KeyValueProtocol等)。
步骤2:添加Camel AVRO依赖
确保项目引入Camel的AVRO组件依赖:
<dependency> <groupId>org.apache.camel</groupId> <artifactId>camel-avro</artifactId> <version>3.20.7</version> <!-- 版本要和你使用的Camel核心版本一致 --> </dependency>
步骤3:定义Camel路由
分**服务端(暴露AVRO接口)和客户端(调用AVRO接口)**两种场景实现:
场景1:作为服务端,提供AVRO服务
如果要让Camel接收外部的put/get请求并处理,可以这样写路由(以Spring Boot为例):
首先注册AVRO协议Bean:
import org.apache.avro.Protocol; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import java.io.File; import java.io.IOException; @Configuration public class AvroConfig { @Bean public Protocol keyValueProtocol() throws IOException { return Protocol.parse(new File("src/main/resources/avro/keyvalue.avpr")); } }
然后编写核心路由逻辑:
import org.apache.camel.builder.RouteBuilder; import org.apache.camel.component.avro.AvroConstants; import org.springframework.stereotype.Component; @Component public class AvroServiceRoute extends RouteBuilder { private final KeyValueService keyValueService; // 构造注入业务服务 public AvroServiceRoute(KeyValueService keyValueService) { this.keyValueService = keyValueService; } @Override public void configure() throws Exception { // 暴露AVRO服务端点,监听8080端口 from("avro:http://localhost:8080/avro/KeyValueProtocol?protocol=#keyValueProtocol") .choice() // 处理put请求 .when(header(AvroConstants.AVRO_MESSAGE_NAME).isEqualTo("put")) .bean(keyValueService, "put") // 处理get请求 .when(header(AvroConstants.AVRO_MESSAGE_NAME).isEqualTo("get")) .bean(keyValueService, "get") .otherwise() .log("收到未知AVRO消息:${header." + AvroConstants.AVRO_MESSAGE_NAME + "}"); } }
对应的业务服务示例:
import org.apache.camel.avro.generated.Key; import org.apache.camel.avro.generated.Value; import org.springframework.stereotype.Service; import java.util.HashMap; import java.util.Map; @Service public class KeyValueService { private final Map<String, String> dataStore = new HashMap<>(); public void put(Key key, Value value) { dataStore.put(key.getKey(), value.getValue()); } public Value get(Key key) { String valueStr = dataStore.get(key.getKey()); return valueStr != null ? new Value(valueStr) : null; } }
场景2:作为客户端,调用远程AVRO服务
如果要让Camel主动调用上面的AVRO服务,可以这样写路由:
import org.apache.camel.builder.RouteBuilder; import org.apache.camel.component.avro.AvroConstants; import org.apache.camel.avro.generated.Key; import org.apache.camel.avro.generated.Value; import org.springframework.stereotype.Component; @Component public class AvroClientRoute extends RouteBuilder { @Override public void configure() throws Exception { // 调用put接口的路由 from("direct:sendPut") .setHeader(AvroConstants.AVRO_MESSAGE_NAME, constant("put")) .setBody(exchange -> { // 构造AVRO请求参数 Key key = new Key("demo-key"); Value value = new Value("demo-value"); return org.apache.camel.avro.generated.KeyValueProtocol.PutRequest.newBuilder() .setKey(key) .setValue(value) .build(); }) .to("avro:http://localhost:8080/avro/KeyValueProtocol") .log("Put请求发送成功"); // 调用get接口的路由 from("direct:sendGet") .setHeader(AvroConstants.AVRO_MESSAGE_NAME, constant("get")) .setBody(exchange -> { Key key = new Key("demo-key"); return org.apache.camel.avro.generated.KeyValueProtocol.GetRequest.newBuilder() .setKey(key) .build(); }) .to("avro:http://localhost:8080/avro/KeyValueProtocol") .log("Get响应结果:${body.value}"); } }
小提示
- 务必保证Camel版本和AVRO版本兼容,避免依赖冲突
- 纯Java项目(非Spring Boot)需要手动初始化Camel上下文并注册相关Bean
- Camel会自动处理AVRO消息的序列化/反序列化,无需手动操作
内容的提问来源于stack exchange,提问作者THM
相关产品推荐
相关产品推荐

