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

RabbitMQ?Stream插件使用案例代碼

 更新時間:2024年04月17日 11:49:49   作者:Doker技術(shù)人的品牌  
這篇文章主要介紹了RabbitMQ?Stream插件使用案例代碼,2.4版為RabbitMQ流插件引入了對RabbitMQStream插件Java客戶端的初始支持,需要的朋友可以參考下

2.4版為RabbitMQ流插件引入了對RabbitMQStream插件Java客戶端的初始支持。

  • RabbitStreamTemplate
  • StreamListener容器

將spring rabbit流依賴項(xiàng)添加到項(xiàng)目中:

<dependency>
  <groupId>org.springframework.amqp</groupId>
  <artifactId>spring-rabbit-stream</artifactId>
  <version>3.1.4</version>
</dependency>

您可以使用RabbitAdmin bean,使用QueueBuilder.stream()方法指定隊(duì)列類型,正常地配置隊(duì)列。例如:

@Bean
Queue stream() {
    return QueueBuilder.durable("stream.queue1")
            .stream()
            .build();
}

然而,這僅在您還使用non-stream 組件(如SimpleMessageListenerContainer或DirectMessageListeneerContainer)時才有效,因?yàn)樵诖蜷_AMQP連接時會觸發(fā)管理員來聲明定義的bean。如果您的應(yīng)用程序僅使用流組件,或者您希望使用高級流配置功能,則應(yīng)改為配置StreamAdmin:

@Bean
StreamAdmin streamAdmin(Environment env) {
    return new StreamAdmin(env, sc -> {
        sc.stream("stream.queue1").maxAge(Duration.ofHours(2)).create();
        sc.stream("stream.queue2").create();
    });
}

一、Sending Messages

RabbitStreamTemplate提供RabbitTemplate(AMQP)功能的子集。

public interface RabbitStreamOperations extends AutoCloseable {
	CompletableFuture<Boolean> send(Message message);
	CompletableFuture<Boolean> convertAndSend(Object message);
	CompletableFuture<Boolean> convertAndSend(Object message, @Nullable MessagePostProcessor mpp);
	CompletableFuture<Boolean> send(com.rabbitmq.stream.Message message);
	MessageBuilder messageBuilder();
	MessageConverter messageConverter();
	StreamMessageConverter streamMessageConverter();
	@Override
	void close() throws AmqpException;
}

RabbitStreamTemplate實(shí)現(xiàn)具有以下構(gòu)造函數(shù)和屬性:

public RabbitStreamTemplate(Environment environment, String streamName) {
}
public void setMessageConverter(MessageConverter messageConverter) {
}
public void setStreamConverter(StreamMessageConverter streamConverter) {
}
public synchronized void setProducerCustomizer(ProducerCustomizer producerCustomizer) {
}

MessageConverter在convertAndSend方法中用于將對象轉(zhuǎn)換為Spring AMQP消息。

StreamMessageConverter用于將Spring AMQP消息轉(zhuǎn)換為本機(jī)流消息。

您也可以直接發(fā)送本機(jī)流消息;使用messageBuilder()方法提供對生產(chǎn)者的消息生成器的訪問。

ProducerCustomizer提供了一種機(jī)制,用于在生成生產(chǎn)者之前對其進(jìn)行自定義。

 二、Receiving Messages

異步消息接收由StreamListenerContainer(以及使用@RabbitListener時的StreamRabbitListerContainerFactory)提供。

偵聽器容器需要一個Environment以及一個流名稱。

您可以使用經(jīng)典的MessageListener接收Spring AMQP消息,也可以使用新接口接收本地流消息:

public interface StreamMessageListener extends MessageListener {
	void onStreamMessage(Message message, Context context);
}

有關(guān)支持的屬性的信息,請參閱消息偵聽器容器配置。

與模板類似,容器具有ConsumerCustomizer屬性。

有關(guān)自定義環(huán)境和使用者的信息,請參閱Java客戶端文檔。

使用@RabbitListener時,配置StreamRabbitListerContainerFactory;此時,大多數(shù)@RabbitListener屬性(并發(fā)等)將被忽略。僅支持id、隊(duì)列、autoStartup和containerFactory。此外,隊(duì)列只能包含一個流名稱。

三、Examples

@Bean
RabbitStreamTemplate streamTemplate(Environment env) {
    RabbitStreamTemplate template = new RabbitStreamTemplate(env, "test.stream.queue1");
    template.setProducerCustomizer((name, builder) -> builder.name("test"));
    return template;
}
@Bean
RabbitListenerContainerFactory<StreamListenerContainer> rabbitListenerContainerFactory(Environment env) {
    return new StreamRabbitListenerContainerFactory(env);
}
@RabbitListener(queues = "test.stream.queue1")
void listen(String in) {
    ...
}
@Bean
RabbitListenerContainerFactory<StreamListenerContainer> nativeFactory(Environment env) {
    StreamRabbitListenerContainerFactory factory = new StreamRabbitListenerContainerFactory(env);
    factory.setNativeListener(true);
    factory.setConsumerCustomizer((id, builder) -> {
        builder.name("myConsumer")
                .offset(OffsetSpecification.first())
                .manualTrackingStrategy();
    });
    return factory;
}
@RabbitListener(id = "test", queues = "test.stream.queue2", containerFactory = "nativeFactory")
void nativeMsg(Message in, Context context) {
    ...
    context.storeOffset();
}
@Bean
Queue stream() {
    return QueueBuilder.durable("test.stream.queue1")
            .stream()
            .build();
}
@Bean
Queue stream() {
    return QueueBuilder.durable("test.stream.queue2")
            .stream()
            .build();
}

2.4.5版將adviceChain屬性添加到StreamListenerContainer(及其工廠)。還提供了一個新的工廠bean來創(chuàng)建一個無狀態(tài)重試攔截器,該攔截器帶有一個可選的StreamMessageRecoverer,用于在使用原始流消息時使用。

@Bean
public StreamRetryOperationsInterceptorFactoryBean sfb(RetryTemplate retryTemplate) {
    StreamRetryOperationsInterceptorFactoryBean rfb =
            new StreamRetryOperationsInterceptorFactoryBean();
    rfb.setRetryOperations(retryTemplate);
    rfb.setStreamMessageRecoverer((msg, context, throwable) -> {
        ...
    });
    return rfb;
}

四、Super Streams

超級流是分區(qū)流的抽象概念,通過將多個流隊(duì)列綁定到具有參數(shù)x-Super-Stream:true的交換來實(shí)現(xiàn)。

1、調(diào)配

為了方便起見,可以通過定義類型為SuperStream的單個bean來提供超級流。

@Bean
SuperStream superStream() {
    return new SuperStream("my.super.stream", 3);
}

RabbitAdmin檢測到這個bean,并將聲明交換(my.super.stream)和3個隊(duì)列(分區(qū))-my.super-stream-n,其中n是0,1,2,綁定的路由密鑰等于n。

如果您還希望通過AMQP向exchange 發(fā)布,您可以提供自定義路由密鑰:

@Bean
SuperStream superStream() {
    return new SuperStream("my.super.stream", 3, (q, i) -> IntStream.range(0, i)
					.mapToObj(j -> "rk-" + j)
					.collect(Collectors.toList()));
}

key 的數(shù)量必須等于分區(qū)的數(shù)量。

2、向超級流生產(chǎn)消息

你必須向 RabbitStreamTemplate 添加一個 superStreamRoutingFunction

@Bean
RabbitStreamTemplate streamTemplate(Environment env) {
    RabbitStreamTemplate template = new RabbitStreamTemplate(env, "stream.queue1");
    template.setSuperStreamRouting(message -> {
        // some logic to return a String for the client's hashing algorithm
    });
    return template;
}

你也可以通過AMQP發(fā)布,使用 RabbitTemplate。

到此這篇關(guān)于RabbitMQ Stream插件使用詳解的文章就介紹到這了,更多相關(guān)RabbitMQ Stream插件內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • 淺談JAVA并發(fā)之ReentrantLock

    淺談JAVA并發(fā)之ReentrantLock

    本文主要介紹了基于AQS實(shí)現(xiàn)的ReentrantLock(重入鎖),感興趣的同學(xué),可以參考下。
    2021-06-06
  • mybatis批量update時報(bào)錯multi-statement not allow的問題

    mybatis批量update時報(bào)錯multi-statement not allow的問題

    這篇文章主要介紹了mybatis批量update時報(bào)錯multi-statement not allow的問題及解決方案,具有很好的參考價(jià)值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教
    2023-10-10
  • 詳解Java使用super和this來重載構(gòu)造方法

    詳解Java使用super和this來重載構(gòu)造方法

    這篇文章主要介紹了詳解Java使用super和this來重載構(gòu)造方法的相關(guān)資料,這里提供實(shí)例來幫助大家理解這部分內(nèi)容,需要的朋友可以參考下
    2017-08-08
  • 通過使用Byte?Buddy便捷創(chuàng)建Java?Agent

    通過使用Byte?Buddy便捷創(chuàng)建Java?Agent

    這篇文章主要為大家介紹了如何通過使用Byte?Buddy便捷創(chuàng)建Java?Agent的使用說明,有需要的朋友可以借鑒參考下希望能夠有所幫助,祝大家多多進(jìn)步
    2022-03-03
  • 詳解Spring Boot 事務(wù)的使用

    詳解Spring Boot 事務(wù)的使用

    spring Boot 使用事務(wù)非常簡單,首先使用注解 @EnableTransactionManagement 開啟事務(wù)支持后,然后在訪問數(shù)據(jù)庫的Service方法上添加注解 @Transactional 便可。接下來通過本文重點(diǎn)給大家介紹spring boot事務(wù)的使用,需要的的朋友參考下吧
    2017-04-04
  • Java Socket編程心跳包創(chuàng)建實(shí)例解析

    Java Socket編程心跳包創(chuàng)建實(shí)例解析

    這篇文章主要介紹了Java Socket編程心跳包創(chuàng)建實(shí)例解析,具有一定借鑒價(jià)值,需要的朋友可以參考下
    2017-12-12
  • Java實(shí)現(xiàn)HTTPS連接的示例代碼

    Java實(shí)現(xiàn)HTTPS連接的示例代碼

    現(xiàn)在的網(wǎng)絡(luò)世界,安全性是大家都非常關(guān)注的問題,特別是對于咱們這些程序員來說,所以,理解并實(shí)現(xiàn)HTTPS連接,對于保護(hù)咱們的數(shù)據(jù)安全是極其重要的,下面我們就來學(xué)習(xí)一下在Java中如何實(shí)現(xiàn)HTTPS連接吧
    2023-12-12
  • Servlet+MyBatis項(xiàng)目轉(zhuǎn)Spring Cloud微服務(wù),多數(shù)據(jù)源配置修改建議

    Servlet+MyBatis項(xiàng)目轉(zhuǎn)Spring Cloud微服務(wù),多數(shù)據(jù)源配置修改建議

    今天小編就為大家分享一篇關(guān)于Servlet+MyBatis項(xiàng)目轉(zhuǎn)Spring Cloud微服務(wù),多數(shù)據(jù)源配置修改建議,小編覺得內(nèi)容挺不錯的,現(xiàn)在分享給大家,具有很好的參考價(jià)值,需要的朋友一起跟隨小編來看看吧
    2019-01-01
  • java堆排序概念原理介紹

    java堆排序概念原理介紹

    在本篇文章里我們給大家分享了關(guān)于java堆排序的概念原理相關(guān)知識點(diǎn)內(nèi)容,有需要的朋友們可以學(xué)習(xí)下。
    2018-10-10
  • SpringBoot多種場景傳參模式

    SpringBoot多種場景傳參模式

    傳參是非常常見的,本文主要介紹了SpringBoot多種場景傳參模式,具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2021-07-07

最新評論

蓬安县| 常州市| 岳阳县| 门源| 荔浦县| 区。| 康平县| 鹤峰县| 平塘县| 酉阳| 清水县| 平凉市| 普洱| 措美县| 锦屏县| 沈丘县| 潜山县| 九江市| 衢州市| 柏乡县| 黄浦区| 临桂县| 禹州市| 多伦县| 昆明市| 安徽省| 普兰县| 卫辉市| 唐山市| 古丈县| 多伦县| 分宜县| 宜章县| 平泉县| 武清区| 怀集县| 黎平县| 玛曲县| 农安县| 衡山县| 金塔县|