使用Apache Camel从Kafka消费Protobuf消息遇错误,求指导
Apache Camel消费Kafka Protobuf消息报错解决
问题场景
需通过Apache Camel从Kafka主题消费Protobuf格式消息,转换为实体类后持久化到数据库,但消费环节触发报错,相关代码及错误信息如下:
测试代码
@Test void produceMessage() { DataChangedEvent event = DataChangedEvent.newBuilder().setId("123").build(); template.sendBody("direct:kafkaEventTopic", event); }
路由代码
@Component public class EventRoute extends RouteBuilder { @Override public void configure() { from("kafkaEventTopic") .routeId("eventRoute") .log(LoggingLevel.INFO, "About to process event: ${body}") // .unmarshal().protobuf(ProtobufLibrary.Jackson, JsonNode.class). .unmarshal().protobuf(new JacksonProtobufDataFormat(DataChangedEvent.class)) .log(LoggingLevel.INFO, "out: ${body}"); } }
报错信息
原因:java.lang.IllegalArgumentException: 无法找到属性defaultInstance的合适setter方法,因为不存在相同类型的setter方法:org.apache.camel.component.jackson.protobuf.JacksonProtobufDataFormat,也无法进行类型转换:setDefaultInstance的参数应为com.google.protobuf.Message的子类
解决方案
报错核心是JacksonProtobufDataFormat初始化方式错误:你传入的是DataChangedEvent.class,但该类构造方法需要的是Protobuf消息的defaultInstance实例,而非Class对象。
修正后的路由代码:
@Component public class EventRoute extends RouteBuilder { @Override public void configure() { from("kafkaEventTopic") .routeId("eventRoute") .log(LoggingLevel.INFO, "About to process event: ${body}") // 传入DataChangedEvent的defaultInstance实例 .unmarshal().protobuf(new JacksonProtobufDataFormat(DataChangedEvent.getDefaultInstance())) .log(LoggingLevel.INFO, "out: ${body}"); } }
额外注意事项:
- 确保项目已引入
camel-jackson-protobuf依赖 - Kafka端点需配置为二进制消息格式,避免额外序列化干扰
- 若要转换为自定义实体类,可在unmarshal后添加
bean()处理器完成对象转换,再执行数据库持久化操作
内容的提问来源于stack exchange,提问作者TheCoder
相关产品推荐
相关产品推荐

