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

報(bào)錯(cuò)日志
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。。。
問(wèn)題所在
yml配置文件中定義的兩個(gè)通道:evad_input和devilvan_input,卻共用了一個(gè)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:
# 一個(gè)通道只能唯一對(duì)應(yīng)一個(gè)topic,否則會(huì)報(bào)binder
destination: Evad
binder: kafka
# 一個(gè)消費(fèi)者組可以被多個(gè)通道使用
group: consumer-group-evad
content-type: text/plain
devilvan_output:
destination: Evad
binder: kafka
content-type: text/plain解決方法
新定義一個(gè)Topic:Evad05,使devilvan通道對(duì)應(yīng)topic,區(qū)別于evad通道對(duì)應(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:
# 一個(gè)通道只能唯一對(duì)應(yīng)一個(gè)topic,否則會(huì)報(bào)binder
destination: Evad05
binder: kafka
# 一個(gè)消費(fèi)者組可以被多個(gè)通道使用
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(通過(guò)通道發(fā)送消息)
@Slf4j
@Component
public class EvadMessageSender {
@Autowired
private EvadChannel channel;
/**
* 消息發(fā)送到默認(rèn)通道:缺省通道對(duì)應(yīng)缺省主題
*
* @param message
*/
public void sendEvadMessage(String message) {
channel.sendEvadMessage().send(MessageBuilder.withPayload(message).build());
}
/**
* 消息發(fā)送到默認(rèn)通道:缺省通道對(duì)應(yīng)缺省主題
*
* @param message
*/
public void sendDevilvanMessage(String message) {
channel.sendDevilvanMessage().send(MessageBuilder.withPayload(message).build());
}
}
4. EvadReceiveListener(訂閱/消費(fèi)者)
@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é)
以上為個(gè)人經(jīng)驗(yàn),希望能給大家一個(gè)參考,也希望大家多多支持腳本之家。
相關(guān)文章
springboot中項(xiàng)目啟動(dòng)時(shí)實(shí)現(xiàn)初始化方法加載參數(shù)
這篇文章主要介紹了springboot中項(xiàng)目啟動(dòng)時(shí)實(shí)現(xiàn)初始化方法加載參數(shù),具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2021-12-12
詳解SpringBoot中添加@ResponseBody注解會(huì)發(fā)生什么
這篇文章主要介紹了詳解SpringBoot中添加@ResponseBody注解會(huì)發(fā)生什么,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧2020-11-11
Java如何把int類(lèi)型轉(zhuǎn)換成byte
這篇文章主要介紹了Java如何把int類(lèi)型轉(zhuǎn)換成byte,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下2020-02-02
Springboot集成spring data elasticsearch過(guò)程詳解
這篇文章主要介紹了springboot集成spring data elasticsearch過(guò)程詳解,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下2020-04-04
Java?中很好用的數(shù)據(jù)結(jié)構(gòu)(你絕對(duì)沒(méi)用過(guò))
今天跟大家介紹的就是?java.util.EnumMap,也是?java.util?包下面的一個(gè)集合類(lèi),同樣的也有對(duì)應(yīng)的的?java.util.EnumSet,對(duì)java數(shù)據(jù)結(jié)構(gòu)相關(guān)知識(shí)感興趣的朋友一起看看吧2022-05-05
SpringBoot運(yùn)用AOP來(lái)實(shí)現(xiàn)分布式鎖的示例代碼
本文主要介紹了通過(guò)注解和AOP實(shí)現(xiàn)分布式鎖的方案,包含鎖過(guò)期時(shí)間、等待超時(shí)設(shè)置及自動(dòng)續(xù)約功能,利用定時(shí)任務(wù)監(jiān)控鎖狀態(tài)并延長(zhǎng)有效期,感興趣的可以了解一下2025-09-09
關(guān)于@DS注解切換數(shù)據(jù)源失敗的原因?qū)崙?zhàn)記錄
項(xiàng)目配置了多個(gè)數(shù)據(jù)源,需要使用@DS注解來(lái)切換數(shù)據(jù)源,但是卻遇到了問(wèn)題,下面這篇文章主要給大家介紹了關(guān)于@DS注解切換數(shù)據(jù)源失敗原因的相關(guān)資料,需要的朋友可以參考下2023-05-05

