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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 11:37:14