解決SpringCloudStream整合Kafka,兩個通道對應(yīng)同一個topic報錯的情況
總結(jié)
- 一個通道(如:evad_input)只能唯一對應(yīng)一個topic,否則會報錯
- 消費者組則可以被多個通道共同使用

報錯日志
2022-05-25 14:46:03.697 ERROR 17108 --- [ask-scheduler-1] o.s.cloud.stream.binding.BindingService : Failed to create consumer binding; retrying in 30 seconds
。。。
org.springframework.cloud.stream.binder.BinderException: Exception thrown while starting consumer:
。。。
Caused by: org.springframework.beans.factory.support.BeanDefinitionOverrideException: Invalid bean definition with name 'Evad.consumer-group-evad.errors.recoverer' defined in null。。。
問題所在
yml配置文件中定義的兩個通道:evad_input和devilvan_input,卻共用了一個topic:Evad,導(dǎo)致綁定失敗。
配置文件
spring:
application:
name: devilvan-kafka
cloud:
stream:
default-binder: kafka
bindings:
evad_input:
destination: Evad
binder: kafka
group: consumer-group-evad
content-type: text/plain
evad_output:
destination: Evad
binder: kafka
content-type: text/plain
devilvan_input:
# 一個通道只能唯一對應(yīng)一個topic,否則會報binder
destination: Evad
binder: kafka
# 一個消費者組可以被多個通道使用
group: consumer-group-evad
content-type: text/plain
devilvan_output:
destination: Evad
binder: kafka
content-type: text/plain解決方法
新定義一個Topic:Evad05,使devilvan通道對應(yīng)topic,區(qū)別于evad通道對應(yīng)的topic
修改后
spring:
application:
name: devilvan-kafka
cloud:
stream:
default-binder: kafka
bindings:
evad_input:
destination: Evad
binder: kafka
group: consumer-group-evad
content-type: text/plain
evad_output:
destination: Evad
binder: kafka
content-type: text/plain
devilvan_input:
# 一個通道只能唯一對應(yīng)一個topic,否則會報binder
destination: Evad05
binder: kafka
# 一個消費者組可以被多個通道使用
group: consumer-group-evad
content-type: text/plain
devilvan_output:
destination: Evad05
binder: kafka
content-type: text/plain代碼
1. XXXController(生產(chǎn)消息的控制器)
@PostMapping("sendEvadMessage")
public ResultMessage<String> sendEvadMessage(@RequestBody String message) {
ResultMessage<String> resultMessage = new ResultMessage<>();
sender.sendEvadMessage(message);
resultMessage.setData(message);
return resultMessage.success();
}
@PostMapping("sendDevilvanMessage")
public ResultMessage<String> sendDevilvanMessage(@RequestBody String message) {
ResultMessage<String> resultMessage = new ResultMessage<>();
sender.sendDevilvanMessage(message);
resultMessage.setData(message);
return resultMessage.success();
}
2. 自定義通道
public interface EvadChannel {
String EVAD_INPUT = "evad_input";
String EVAD_OUTPUT = "evad_output";
String DEVILVAN_INPUT = "devilvan_input";
String DEVILVAN_OUTPUT = "devilvan_output";
/**
* 缺省接收消息通道
* @return channel 返回缺省信息接收通道
*/
@Input(EVAD_INPUT)
MessageChannel receiveEvadMessage();
/**
* 缺省發(fā)送消息通道
* @return channel 返回缺省信息發(fā)送通道
*/
@Output(EVAD_OUTPUT)
MessageChannel sendEvadMessage();
/**
* 缺省接收消息通道
* @return channel 返回缺省信息接收通道
*/
@Input(DEVILVAN_INPUT)
MessageChannel receiveDevilvanMessage();
/**
* 缺省發(fā)送消息通道
* @return channel 返回缺省信息發(fā)送通道
*/
@Output(DEVILVAN_OUTPUT)
MessageChannel sendDevilvanMessage();
}
3. EvadMessageSender(通過通道發(fā)送消息)
@Slf4j
@Component
public class EvadMessageSender {
@Autowired
private EvadChannel channel;
/**
* 消息發(fā)送到默認通道:缺省通道對應(yīng)缺省主題
*
* @param message
*/
public void sendEvadMessage(String message) {
channel.sendEvadMessage().send(MessageBuilder.withPayload(message).build());
}
/**
* 消息發(fā)送到默認通道:缺省通道對應(yīng)缺省主題
*
* @param message
*/
public void sendDevilvanMessage(String message) {
channel.sendDevilvanMessage().send(MessageBuilder.withPayload(message).build());
}
}
4. EvadReceiveListener(訂閱/消費者)
@Slf4j
@Configuration
@EnableBinding(value = EvadChannel.class)
public class EvadReceiveListener {
@StreamListener(EvadChannel.EVAD_INPUT)
public void receiveEvadMessage(Message<String> message) {
log.info("{} 訂閱消息:通道 = " + EvadChannel.EVAD_INPUT + ",data = {}",
DateUtil.now(), message.getPayload());
}
@StreamListener(EvadChannel.DEVILVAN_INPUT)
public void receiveDevilvanMessage(Message<String> message) {
log.info("{} 訂閱消息:通道 = " + EvadChannel.DEVILVAN_INPUT + ",data = {}",
DateUtil.now(), message.getPayload());
}
}
總結(jié)
以上為個人經(jīng)驗,希望能給大家一個參考,也希望大家多多支持腳本之家。
相關(guān)文章
springboot中項目啟動時實現(xiàn)初始化方法加載參數(shù)
這篇文章主要介紹了springboot中項目啟動時實現(xiàn)初始化方法加載參數(shù),具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教2021-12-12
詳解SpringBoot中添加@ResponseBody注解會發(fā)生什么
這篇文章主要介紹了詳解SpringBoot中添加@ResponseBody注解會發(fā)生什么,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧2020-11-11
Springboot集成spring data elasticsearch過程詳解
這篇文章主要介紹了springboot集成spring data elasticsearch過程詳解,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友可以參考下2020-04-04
Java?中很好用的數(shù)據(jù)結(jié)構(gòu)(你絕對沒用過)
今天跟大家介紹的就是?java.util.EnumMap,也是?java.util?包下面的一個集合類,同樣的也有對應(yīng)的的?java.util.EnumSet,對java數(shù)據(jù)結(jié)構(gòu)相關(guān)知識感興趣的朋友一起看看吧2022-05-05
SpringBoot運用AOP來實現(xiàn)分布式鎖的示例代碼
本文主要介紹了通過注解和AOP實現(xiàn)分布式鎖的方案,包含鎖過期時間、等待超時設(shè)置及自動續(xù)約功能,利用定時任務(wù)監(jiān)控鎖狀態(tài)并延長有效期,感興趣的可以了解一下2025-09-09
關(guān)于@DS注解切換數(shù)據(jù)源失敗的原因?qū)崙?zhàn)記錄
項目配置了多個數(shù)據(jù)源,需要使用@DS注解來切換數(shù)據(jù)源,但是卻遇到了問題,下面這篇文章主要給大家介紹了關(guān)于@DS注解切換數(shù)據(jù)源失敗原因的相關(guān)資料,需要的朋友可以參考下2023-05-05

