惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

美团技术团队
人人都是产品经理
人人都是产品经理
月光博客
月光博客
V
V2EX
WordPress大学
WordPress大学
酷 壳 – CoolShell
酷 壳 – CoolShell
Last Week in AI
Last Week in AI
博客园 - 三生石上(FineUI控件)
小众软件
小众软件
Hugging Face - Blog
Hugging Face - Blog
V
Visual Studio Blog
宝玉的分享
宝玉的分享
雷峰网
雷峰网
钛媒体:引领未来商业与生活新知
钛媒体:引领未来商业与生活新知
博客园 - Franky
博客园 - 聂微东
博客园 - 司徒正美
博客园 - 【当耐特】
爱范儿
爱范儿
奇客Solidot–传递最新科技情报
奇客Solidot–传递最新科技情报
OSCHINA 社区最新新闻
OSCHINA 社区最新新闻
大猫的无限游戏
大猫的无限游戏
博客园 - 叶小钗
阮一峰的网络日志
阮一峰的网络日志

博客园 - 蜗牛无敌

管理心得-协调 项目git 提交 hooks excel 公式支持 ~ - flowable-ui本地部署 mqtt实战 自动生成提示 修改表中数据时候,关联查询,例如删除一张表中 小于最大时间的sql, rocketmq-spring-boot-starter的使用 MYSQL ORDER BY 自定义排序 java程序cpu飘高排查 window canal的本地部署 是用hutool工具实现动态定时任务 python使用selenium框架模拟登录 获取token,并打包成exe springboot mybatisplus 使用拦截器改写sql语句 spring-cloud-stream-rocketmq实现消息的顺序消息 springboot 集成es 使用resthightlevel springboot 多租户环境下 由于业务需要搜集所有租户的基本数据 spring boot rockmq 生产者和消费者 实战 springboot为基础的项目 由于线上日志打印的级别较低,无法打印日志的方法 window系统关闭指定端口的应用方法
spring-cloud-stater-stream-rocketmq
蜗牛无敌 · 2025-05-30 · via 博客园 - 蜗牛无敌
   <dependency>
            <groupId>com.alibaba.cloud</groupId>
            <artifactId>spring-cloud-starter-stream-rocketmq</artifactId>
            <version>2021.0.5.0</version>
        </dependency>

1、配置

spring:
  application:
    # 应用名称
    name: jnpf-ftb
  cloud:
    stream:
      rocketmq:
        binder:
          name-server: ${ROCKETMQ_HOST:192.168.3.22:30094}
        bindings:
          output:
            producer:
              sync: true
              group: jnpf-group1
      bindings:
        permission-output:  #生产
          content-type: text/json
          destination: permission-topic
          group: jnpf-group1
        permission-input: #消费
          content-type: text/json
          destination: permission-topic
          group: jnpf-ftb-consumer

2、通道

package jnpf.qualifications.consummer;

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 PositionGradeConsumerSource {
    String OUTPUT = "permission-output";
    String INPUT = "permission-input";
    @Output(OUTPUT)
    MessageChannel output();
    @Input(INPUT)
    SubscribableChannel input();
}

 3、监听消息,注意header可以通过调整查询key是什么,因为每个版本可能不一样

package jnpf.qualifications.consummer;

import lombok.extern.slf4j.Slf4j;
import org.springframework.cloud.stream.annotation.EnableBinding;
import org.springframework.cloud.stream.annotation.StreamListener;
import org.springframework.stereotype.Component;
import org.springframework.messaging.Message;

@Slf4j
@Component
@EnableBinding(PositionGradeConsumerSource.class)
public class PositionGradeConsumer {
    @StreamListener(target = PositionGradeConsumerSource.INPUT,condition = "headers['ROCKET_TAGS'] == 'TAG_GRADE'")
    public void receive(Message<String> message) {
        // 获取消息体
        String payload = message.getPayload();

        // 获取 headers
        var headers = message.getHeaders();

        log.error("接收到消息:{},Tags={}", payload, tags);
    }
}