最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

解決SpringCloudStream整合Kafka,兩個(gè)通道對(duì)應(yīng)同一個(gè)topic報(bào)錯(cuò)的情況

 更新時(shí)間:2026年05月06日 09:09:58   作者:加把勁騎士RideOn  
文章指出通道需唯一對(duì)應(yīng)topic,否則會(huì)報(bào)錯(cuò),因兩個(gè)通道共用一個(gè)topic導(dǎo)致綁定失敗,通過(guò)修改配置文件,為不同通道分配不同的topic解決

總結(jié)

  1. 一個(gè)通道(如:evad_input)只能唯一對(duì)應(yīng)一個(gè)topic,否則會(huì)報(bào)錯(cuò)
  2. 消費(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)文章

  • 什么是RESTful?API,有什么作用

    什么是RESTful?API,有什么作用

    提到RESTful?API大家勢(shì)必或多或少聽(tīng)說(shuō)過(guò),但是什么是RESTful?API??如何理解RESTful?API?呢?今天咱們就來(lái)聊聊這個(gè)RESTful?API
    2023-11-11
  • springboot中項(xiàng)目啟動(dòng)時(shí)實(shí)現(xiàn)初始化方法加載參數(shù)

    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ā)生什么

    這篇文章主要介紹了詳解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

    這篇文章主要介紹了Java如何把int類(lèi)型轉(zhuǎn)換成byte,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下
    2020-02-02
  • SpringBoot中時(shí)間格式化的五種方法匯總

    SpringBoot中時(shí)間格式化的五種方法匯總

    時(shí)間格式化在項(xiàng)目中使用頻率是非常高的,這篇文章主要給大家介紹了關(guān)于SpringBoot中時(shí)間格式化的五種方法,文中通過(guò)示例代碼介紹的非常詳細(xì),需要的朋友可以參考下
    2021-07-07
  • Springboot集成spring data elasticsearch過(guò)程詳解

    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?中很好用的數(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)分布式鎖的示例代碼

    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
  • java中@NotBlank限制屬性不能為空

    java中@NotBlank限制屬性不能為空

    在實(shí)體類(lèi)的對(duì)應(yīng)屬性上添 @NotBlank注解,可以實(shí)現(xiàn)對(duì)空置的限制,本文就來(lái)介紹一下java中@NotBlank限制屬性不能為空,感興趣的可以了解一下
    2024-01-01
  • 關(guān)于@DS注解切換數(shù)據(jù)源失敗的原因?qū)崙?zhàn)記錄

    關(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

最新評(píng)論

体育| 峨山| 竹溪县| 收藏| 乳山市| 全州县| 广德县| 洛扎县| 定南县| 安康市| 横峰县| 宣恩县| 辽宁省| 遂宁市| 宽甸| 平潭县| 博白县| 安龙县| 巴彦县| 永定县| 闸北区| 拜泉县| 新疆| 拜泉县| 忻州市| 曲阳县| 武强县| 瓦房店市| 临朐县| 京山县| 长春市| 哈密市| 聊城市| 民县| 板桥市| 鹿邑县| 翼城县| 萨嘎县| 丹巴县| 临潭县| 台山市|