spring cloud stream使用笔记

一、背景

以前我们在spring boot中构建一个消息驱动的微服务应用,通常会使用rabbitMQ或是kafka来做消息中间件,应用中均需代码实现具体消息中间件的通信细节。此时如果再更换一个新的消息中间件,这会我们又需新增这些通信代码,写起来会比较繁琐,而stream出现就是为了简化这一过程。

二、简介


它是一个构建消息驱动的微服务应用的框架。通过一些抽象出来的基础概念,来简化消息中间件的使用。我们可以看下官网上的处理模型图:

image
  • 关键概念

Inputs 接收消息的通道
Output 发送消息的通道
Binder 可理解为一个抽象的中间件,应用通过在spring cloud stream中所注入的inputs,outputs通道来跟外界消息通信,而这些通道又是通过具体中间件的Binder实现来连接到消息队列的服务器上。有了Binder,甚至可以不改一行代码,就切换中间件的类型。目前Binder实现支持的具体中间件类型为:rabbitMQ 和 kfaka这俩

当然MQ中的消费组group 和 分区 partion的概念他也有,跟kafka里面的概念是一样的。
Group:消费组,一个消息到达一个消费组后,只能被这个消费组的其中一个实例消费掉;
Partion:消息分区,一个非常大的topic可以分布到多个broker(即服务器)上

三、使用步骤

以下以rabbitMQ为具体中间件作为示例:

1. 在pom.xml中添加依赖
 <dependency>
       <groupid >org.springframework.cloud</groupid>
       <artifactid>spring-cloud-starter-stream-rabbit</artifactid>
</dependency>
2.自定义通道的创建

自定义通道的创建有两种方式,一种是提前在代码里定义好的,另一种是在运行时通过读取完通道名来创建的。

2.1 方式一:提前定义好的通道
  • 定义生产者
    我们定义一个生产者类SampleSource,这类要完成2件事:
  • 2.1.1 自定义发送通道
  • 2.1.2 完成发送消息的功能
@EnableBinding(SampleSource.MultiOutputSource.class)
public class SampleSource {
   
       //自定义发送通道
    public interface MultiOutputSource {
        String OUTPUT1 = "output1";

        String OUTPUT2 = "output2";

        @Output(OUTPUT1)
        MessageChannel output1();

        @Output(OUTPUT2)
        MessageChannel output2();
    }
}

注意:要加上@EnableBinding 绑定通道才能够发出消息到mq的服务器。

  • 方式二:动态创建通道
@EnableBinding
@Controller
public class SourceWithDynamicDestination {

    @Autowired
    private BinderAwareChannelResolver resolver;

    @RequestMapping(path = "/{target}", method = POST, consumes = "*/*")
    @ResponseStatus(HttpStatus.ACCEPTED)
    public void handleRequest(@RequestBody String body, @PathVariable("target") target,
           @RequestHeader(HttpHeaders.CONTENT_TYPE) Object contentType) {
        sendMessage(body, target, contentType);
    }

    private void sendMessage(String body, String target, Object contentType) {
        resolver.resolveDestination(target).send(MessageBuilder.createMessage(body,
                new MessageHeaders(Collections.singletonMap(MessageHeaders.CONTENT_TYPE, contentType))));
    }
}

这是官网的示例代码,可以看到关键代码是这两句

@Autowired
private BinderAwareChannelResolver resolver;
//中间省略代码...
  resolver.resolveDestination(target).send(MessageBuilder.createMessage(body,
                new MessageHeaders(Collections.singletonMap(MessageHeaders.CONTENT_TYPE, contentType))));

resolveDestination里面的处理就是,先查找传入的target通道名,看下有没创建过,如果没有就会默认创建一个。

  • 生产者实现发送消息的函数
    在SampleSource这类里添加sendMessage函数
public class SampleSource{
 @Bean
    @InboundChannelAdapter(value = MultiOutputSource.OUTPUT1, poller = @Poller(fixedDelay = "1000", maxMessagesPerPoll = "1"))
    public synchronized MessageSource<String> messageSource1() {
        return new MessageSource<String>() {
            public Message<String> receive() {
                String message = "FromSource1";
                System.out.println("******************");
                System.out.println("From Source1");
                System.out.println("******************");
                System.out.println("Sending value: " + message);
                return new GenericMessage(message);
            }
        };
    }

    @Bean
    @InboundChannelAdapter(value = MultiOutputSource.OUTPUT2, poller = @Poller(fixedDelay = "1000", maxMessagesPerPoll = "1"))
    public synchronized MessageSource<String> timerMessageSource() {
        return new MessageSource<String>() {
            public Message<String> receive() {
                String message = "FromSource2";
                System.out.println("******************");
                System.out.println("From Source2");
                System.out.println("******************");
                System.out.println("Sending value: " + message);
                return new GenericMessage(message);
            }
        };
    }

}

对应的配置application.yml:

spring:
  cloud.stream.bindings:
    output1:
      contentType: application/json #约定消息的内容编码格式
    output2:
      contentType: application/json #约定消息的内容编码格式

  rabbitmq:
    host: 127.0.0.1 
    port: 5672
    username: sa
    password: 123456
3. 消费者类的实现

消费者类要做的事也是相似的:
3.1.自定义接收通道
3.2消费消息的功能实现

@EnableBinding(SampleSink.MultiInputSink.class)
public class SampleSink {

    @StreamListener(MultiInputSink.INPUT1)
    public synchronized void receive1(String message) {
        System.out.println("******************");
        System.out.println("At Sink1");
        System.out.println("******************");
        System.out.println("Received message " + message);
    }

    @StreamListener(MultiInputSink.INPUT2)
    public synchronized void receive2(String message) {
        System.out.println("******************");
        System.out.println("At Sink2");
        System.out.println("******************");
        System.out.println("Received message " + message);
    }

    public interface MultiInputSink {
        String INPUT1 = "input1";

        String INPUT2 = "input2";

        @Input(INPUT1)
        SubscribableChannel input1();

        @Input(INPUT2)
        SubscribableChannel input2();
    }
}

对应的配置 application.yml (当然用application.propertities)

spring:
  cloud.stream:
    bindings:
      input1:
        group: inputGroup #加上group是为了持久化
      input2:
        group: inputGroup2
rabbit:
      host: 127.0.0.1 
      port: 5672
      username: sa
      password: 123456
4.死信队列设置

问题列表:

  • 消息在什么条件下进入死信队列?发送失败后,如何设置重试次数 or TTL?

  • 进入死信队列之后,若又需要该消息重新回到原队列进行处理,该怎么办


转入死信队列的条件:

配置死信队列及消息消费失败重试次数(application.yml):

image
  • 配置消息消费重试次数

两种方式
1) 如果允许重试一定次数:如上图配置所示,设置max_attempt ,大于1即可
2)如果不允许重发,消费失败了就进入死信队列,在配置中添加requeueRejected设为true

spring:
  cloud.stream:
    bindings:
      input1:
        group: inputGroup1
    rabbit:
      bindings:
        input1:
          consumer:
            autoBindDlq: true #启用死信队列,默认会生成一个DLX EXCHANGE,当消息重复消费失败后
            dlqDeadLetterExchange: input-deadLetter.DLX  #如果该列声明,那么deadLetterExchange也要声明,这个保持一致
            deadLetterExchange: input-deadLetter.DLX #与dlqDeadLetterExchange保持一致
            requeueRejected: true
      host: 127.0.0.1   
      port: 5672
      username: sa
      password: 123456
  • 在代码中将某消息转入死信队列,另可见官网示例

5、消息的格式:

(1)消息头,包含了以下字段:

image

翻看源码,MessageHeader是个Map结构,左边是字段名,右边是字段内容。可以在创建MessageHeader的时候传入已经初始化的Map,注意我们可以在这指定body的contentType。contentType能填什么内容,查找下表即可(官网上找到的):

image

(2)至于消息体(payLoad),它支持自定义结构,格式自定。


6、如何保证消息的可靠性?

一般是通过具体的消息中间件来保证。

配置组就可以来保证消息可靠性。见官网描述:消息的持久性

设置持久性的属性:durableSubscription

image

四、相关链接

最后编辑于
©著作权归作者所有,转载或内容合作请联系作者
  • 序言:七十年代末,一起剥皮案震惊了整个滨河市,随后出现的几起案子,更是在滨河造成了极大的恐慌,老刑警刘岩,带你破解...
    沈念sama阅读 199,711评论 5 468
  • 序言:滨河连续发生了三起死亡事件,死亡现场离奇诡异,居然都是意外死亡,警方通过查阅死者的电脑和手机,发现死者居然都...
    沈念sama阅读 83,932评论 2 376
  • 文/潘晓璐 我一进店门,熙熙楼的掌柜王于贵愁眉苦脸地迎上来,“玉大人,你说我怎么就摊上这事。” “怎么了?”我有些...
    开封第一讲书人阅读 146,770评论 0 330
  • 文/不坏的土叔 我叫张陵,是天一观的道长。 经常有香客问我,道长,这世上最难降的妖魔是什么? 我笑而不...
    开封第一讲书人阅读 53,799评论 1 271
  • 正文 为了忘掉前任,我火速办了婚礼,结果婚礼上,老公的妹妹穿的比我还像新娘。我一直安慰自己,他们只是感情好,可当我...
    茶点故事阅读 62,697评论 5 359
  • 文/花漫 我一把揭开白布。 她就那样静静地躺着,像睡着了一般。 火红的嫁衣衬着肌肤如雪。 梳的纹丝不乱的头发上,一...
    开封第一讲书人阅读 48,069评论 1 276
  • 那天,我揣着相机与录音,去河边找鬼。 笑死,一个胖子当着我的面吹牛,可吹牛的内容都是我干的。 我是一名探鬼主播,决...
    沈念sama阅读 37,535评论 3 390
  • 文/苍兰香墨 我猛地睁开眼,长吁一口气:“原来是场噩梦啊……” “哼!你这毒妇竟也来了?” 一声冷哼从身侧响起,我...
    开封第一讲书人阅读 36,200评论 0 254
  • 序言:老挝万荣一对情侣失踪,失踪者是张志新(化名)和其女友刘颖,没想到半个月后,有当地人在树林里发现了一具尸体,经...
    沈念sama阅读 40,353评论 1 294
  • 正文 独居荒郊野岭守林人离奇死亡,尸身上长有42处带血的脓包…… 初始之章·张勋 以下内容为张勋视角 年9月15日...
    茶点故事阅读 35,290评论 2 317
  • 正文 我和宋清朗相恋三年,在试婚纱的时候发现自己被绿了。 大学时的朋友给我发了我未婚夫和他白月光在一起吃饭的照片。...
    茶点故事阅读 37,331评论 1 329
  • 序言:一个原本活蹦乱跳的男人离奇死亡,死状恐怖,灵堂内的尸体忽然破棺而出,到底是诈尸还是另有隐情,我是刑警宁泽,带...
    沈念sama阅读 33,020评论 3 315
  • 正文 年R本政府宣布,位于F岛的核电站,受9级特大地震影响,放射性物质发生泄漏。R本人自食恶果不足惜,却给世界环境...
    茶点故事阅读 38,610评论 3 303
  • 文/蒙蒙 一、第九天 我趴在偏房一处隐蔽的房顶上张望。 院中可真热闹,春花似锦、人声如沸。这庄子的主人今日做“春日...
    开封第一讲书人阅读 29,694评论 0 19
  • 文/苍兰香墨 我抬头看了看天上的太阳。三九已至,却和暖如春,着一层夹袄步出监牢的瞬间,已是汗流浃背。 一阵脚步声响...
    开封第一讲书人阅读 30,927评论 1 255
  • 我被黑心中介骗来泰国打工, 没想到刚下飞机就差点儿被人妖公主榨干…… 1. 我叫王不留,地道东北人。 一个月前我还...
    沈念sama阅读 42,330评论 2 346
  • 正文 我出身青楼,却偏偏与公主长得像,于是被迫代替她去往敌国和亲。 传闻我的和亲对象是个残疾皇子,可洞房花烛夜当晚...
    茶点故事阅读 41,904评论 2 341

推荐阅读更多精彩内容