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

spring-cloud-stream的手動(dòng)消息確認(rèn)問(wèn)題

 更新時(shí)間:2023年05月25日 14:38:27   作者:l1161558158  
這篇文章主要介紹了spring-cloud-stream的手動(dòng)消息確認(rèn)問(wèn)題,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教

spring-cloud-stream的手動(dòng)消息確認(rèn)

對(duì)于kafka-binder來(lái)說(shuō),設(shè)置autoCommitOffset為false.然后在listen中手動(dòng)確認(rèn)

@StreamListener(Sink.INPUT)
void listen(@Header(KafkaHeaders.ACKNOWLEDGMENT) Acknowledgment acknowledgment){
? ? //...業(yè)務(wù)代碼
? ? acknowledgment.acknowledge();
}

需要注意的是autoCommitOffset的設(shè)置位置.

spring.cloud.stream.kafka.bindings.input.consumer.autoCommitOffset=false#應(yīng)該在這里設(shè)置
spring.cloud.stream.bindings.input.consumer.autoCommitOffset=false#這里設(shè)置是無(wú)效的,獲取Acknowledgment時(shí)會(huì)是null

springcloud的stream消息組件的使用@StreamListener

常見(jiàn)問(wèn)題(使用rabbitmq)

消息分組防止多實(shí)例重復(fù)消費(fèi)

在一個(gè)服務(wù)多實(shí)例場(chǎng)景下使用默認(rèn)使用@StreamListener監(jiān)聽(tīng)消息消費(fèi),yml中沒(méi)有特殊配置的話是會(huì)導(dǎo)致消息重復(fù)消費(fèi)的,原因是此時(shí)每個(gè)實(shí)例都是匿名在rabbitmq上注冊(cè)的隊(duì)列,需要給消費(fèi)者指定一個(gè)消費(fèi)組,讓消息在組里只被消費(fèi)一次;

spring.cloud.stream.bindings.xxx(消費(fèi)者隊(duì)列名).group=xxx(組名)

在springboot下在同一個(gè)服務(wù)(項(xiàng)目中)使用@input和@outPut時(shí)指定的隊(duì)列名是不可以重復(fù)的.會(huì)在啟動(dòng)編譯的時(shí)候報(bào)bean定義重復(fù)。需要在yml給生產(chǎn)者和消費(fèi)者指定同一個(gè)交換機(jī)。

spring:
  rabbitmq:
    host: xxx.xxx.xxx.xx
    port: 35672
    username: xxx
    password: xxx
    virtual-host: /xxx
  cloud:
    stream:
      bindings:
        in:
          #若消息系統(tǒng)是RabbitMQ,目的地(destination)就是指exchange,消息系統(tǒng)是Kafka,那么就是指topic
          destination: test
          #在多實(shí)例的時(shí)候需要制定一個(gè)消息分組,不然每個(gè)實(shí)例都是匿名方式把隊(duì)列注冊(cè)到rabbitmq上去,導(dǎo)致一個(gè)交換機(jī)下有多個(gè)隊(duì)列
          #并且默認(rèn)生成的交換機(jī)是topic類(lèi)型的,會(huì)導(dǎo)致重復(fù)消費(fèi)
          group: myIn
        out:
          destination: test

先上依賴(lài)

<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>
    <parent>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-parent</artifactId>
        <version>1.5.8.RELEASE</version>
        <relativePath/> <!-- lookup parent from repository -->
    </parent>
    <groupId>com.fchan</groupId>
    <artifactId>springcloudstream</artifactId>
    <version>0.0.1-SNAPSHOT</version>
    <name>springcloudstream</name>
    <description>Demo project for Spring Boot</description>
    <properties>
        <java.version>1.8</java.version>
    </properties>
    <dependencies>
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter</artifactId>
        </dependency>
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-web</artifactId>
        </dependency>
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-test</artifactId>
            <scope>test</scope>
        </dependency>
        <dependency>
            <groupId>org.springframework.cloud</groupId>
            <artifactId>spring-cloud-starter-stream-rabbit</artifactId>
<!--            <version>2.0.1.RELEASE</version>-->
        </dependency>
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-actuator</artifactId>
        </dependency>
    </dependencies>
    <dependencyManagement>
        <dependencies>
            <dependency>
                <groupId>org.springframework.cloud</groupId>
                <artifactId>spring-cloud-stream-dependencies</artifactId>
                <version>Ditmars.RELEASE</version>
                <type>pom</type>
                <scope>import</scope>
            </dependency>
        </dependencies>
    </dependencyManagement>
    <build>
        <plugins>
            <plugin>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-maven-plugin</artifactId>
            </plugin>
        </plugins>
    </build>
</project>

再上yml配置

spring:
  rabbitmq:
    host: xxx.xxx.xxx.xx
    port: 35672
    username: xxx
    password: xxx
    virtual-host: /xxx
  cloud:
    stream:
      bindings:
        in:
          #若消息系統(tǒng)是RabbitMQ,目的地(destination)就是指exchange,消息系統(tǒng)是Kafka,那么就是指topic
          destination: test
          #在多實(shí)例的時(shí)候需要制定一個(gè)消息分組,不然每個(gè)實(shí)例都是匿名方式把隊(duì)列注冊(cè)到rabbitmq上去,導(dǎo)致一個(gè)交換機(jī)下有多個(gè)隊(duì)列
          #并且默認(rèn)生成的交換機(jī)是topic類(lèi)型的,會(huì)導(dǎo)致重復(fù)消費(fèi)
          group: myIn
        out:
          destination: test

消息生產(chǎn)者

package com.fchan.springcloudstream.service;
import org.springframework.cloud.stream.annotation.Input;
import org.springframework.cloud.stream.annotation.Output;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.SubscribableChannel;
public interface MyMessageChannel {
    String out = "out";
    String in = "in";
    @Output(out)
    MessageChannel out();
    @Input(in)
    SubscribableChannel in();
}

發(fā)送消息

package com.fchan.springcloudstream.controller;
import com.fchan.springcloudstream.service.MyMessageChannel;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import javax.annotation.Resource;
import java.util.HashMap;
import java.util.Map;
@RestController
public class MessageController {
    @Resource
    private MyMessageChannel myMessageChannel;
    @RequestMapping("test")
    public String testMessage(){
        Map<String,Object> map = new HashMap<>();
        map.put("shopId", "123");
        myMessageChannel.out().send(MessageBuilder.withPayload(map).build());
        return "success";
    }
}

消息消費(fèi)者

package com.fchan.springcloudstream.service;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.cloud.stream.annotation.EnableBinding;
import org.springframework.cloud.stream.annotation.StreamListener;
import org.springframework.messaging.Message;
import org.springframework.stereotype.Component;
import java.util.Map;
@Component
@EnableBinding({MyMessageChannel.class})
public class MyConsumer {
    Logger log = LoggerFactory.getLogger(MyConsumer.class);
    @StreamListener(MyMessageChannel.in)
    public void input(Message<Map<String,Object>> message){
        log.info("收到消息:{}", message.getPayload());
    }
}

總結(jié)

以上為個(gè)人經(jīng)驗(yàn),希望能給大家一個(gè)參考,也希望大家多多支持腳本之家。

相關(guān)文章

  • vue3將頁(yè)面生成pdf導(dǎo)出的操作指南

    vue3將頁(yè)面生成pdf導(dǎo)出的操作指南

    最近工作中有需要將一些前端頁(yè)面(如報(bào)表頁(yè)面等)導(dǎo)出為pdf的需求,下面這篇文章主要給大家介紹了關(guān)于vue3 如何將頁(yè)面生成 pdf 導(dǎo)出,文中通過(guò)實(shí)例代碼介紹的非常詳細(xì),需要的朋友可以參考下
    2023-07-07
  • Vue自定義render統(tǒng)一項(xiàng)目組彈框功能

    Vue自定義render統(tǒng)一項(xiàng)目組彈框功能

    這篇文章主要介紹了Vue自定義render統(tǒng)一項(xiàng)目組彈框功能,本文通過(guò)實(shí)例代碼給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2020-06-06
  • Vue3.3?+?TS4構(gòu)建實(shí)現(xiàn)ElementPlus功能的組件庫(kù)示例

    Vue3.3?+?TS4構(gòu)建實(shí)現(xiàn)ElementPlus功能的組件庫(kù)示例

    Vue.js?是目前最盛行的前端框架之一,而?TypeScript?則是一種靜態(tài)類(lèi)型言語(yǔ),它能夠讓開(kāi)發(fā)人員在編寫(xiě)代碼時(shí)愈加平安和高效,本文將引見(jiàn)如何運(yùn)用?Vue.js?3.3?和?TypeScript?4?構(gòu)建一個(gè)自主打造媲美?ElementPlus?的組件庫(kù)
    2023-10-10
  • Vue 自定義指令詳解

    Vue 自定義指令詳解

    本文介紹了如何在Vue中定義和使用自定義指令,包括指令的注冊(cè)、鉤子函數(shù)、參數(shù)以及常見(jiàn)指令的封裝,如v-copy、v-longpress等,自定義指令在處理某些底層DOM操作時(shí)非常便捷,感興趣的朋友一起看看吧
    2025-01-01
  • vue + typescript + 極驗(yàn)登錄驗(yàn)證的實(shí)現(xiàn)方法

    vue + typescript + 極驗(yàn)登錄驗(yàn)證的實(shí)現(xiàn)方法

    這篇文章主要介紹了vue + typescript + 極驗(yàn) 登錄驗(yàn)證的實(shí)現(xiàn)方法,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧
    2019-06-06
  • laravel5.3 vue 實(shí)現(xiàn)收藏夾功能實(shí)例詳解

    laravel5.3 vue 實(shí)現(xiàn)收藏夾功能實(shí)例詳解

    這篇文章主要介紹了laravel5.3 vue 實(shí)現(xiàn)收藏夾功能,本文通過(guò)實(shí)例代碼給大家介紹的非常詳細(xì),需要的朋友可以參考下
    2018-01-01
  • vue axios 在頁(yè)面切換時(shí)中斷請(qǐng)求方法 ajax

    vue axios 在頁(yè)面切換時(shí)中斷請(qǐng)求方法 ajax

    下面小編就為大家分享一篇vue axios 在頁(yè)面切換時(shí)中斷請(qǐng)求方法 ajax,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。一起跟隨小編過(guò)來(lái)看看吧
    2018-03-03
  • vue treeselect獲取當(dāng)前選中項(xiàng)的label實(shí)例

    vue treeselect獲取當(dāng)前選中項(xiàng)的label實(shí)例

    這篇文章主要介紹了vue treeselect獲取當(dāng)前選中項(xiàng)的label實(shí)例,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。一起跟隨小編過(guò)來(lái)看看吧
    2020-08-08
  • vue-devtools的安裝與使用教程

    vue-devtools的安裝與使用教程

    vue-devtools是一款基于chrome游覽器的插件,用于調(diào)試vue應(yīng)用,這可以極大地提高我們的調(diào)試效率,這篇文章主要介紹了vue-devtools的安裝與使用教程,需要的朋友可以參考下
    2023-03-03
  • 詳解vue3.2中setup語(yǔ)法糖<script?lang="ts"?setup>

    詳解vue3.2中setup語(yǔ)法糖<script?lang="ts"?setup>

    Vue 3.2 引入了語(yǔ)法,這是一種稍微不那么冗長(zhǎng)的聲明組件的方式,下面這篇文章主要介紹了詳解vue3.2中setup語(yǔ)法糖<script?lang="ts"setup>的相關(guān)資料,文中通過(guò)實(shí)例代碼介紹的非常詳細(xì),需要的朋友可以參考下
    2023-01-01

最新評(píng)論

呼图壁县| 灵川县| 临泉县| 大埔区| 抚远县| 兰考县| 承德市| 长岭县| 高邑县| 棋牌| 康平县| 盐池县| 手游| 杭锦后旗| 武隆县| 宣恩县| 霍城县| 永济市| 聂拉木县| 遵义县| 邛崃市| 万全县| 西和县| 宁陕县| 浦江县| 车致| 栖霞市| 冕宁县| 南汇区| 元谋县| 巴里| 林周县| 建德市| 海宁市| 大余县| 遵义市| 额尔古纳市| 定南县| 揭西县| 那坡县| 沙洋县|