Spring Cloud Stream集成Kafka

这篇具有很好参考价值的文章主要介绍了Spring Cloud Stream集成Kafka。希望对大家有所帮助。如果存在错误或未考虑完全的地方,请大家不吝赐教,您也可以点击"举报违法"按钮提交疑问。

Spring Cloud Stream是一个构建消息驱动微服务的框架,抽象了MQ的使用方式, 提供统一的API操作。Spring Cloud Stream通过Binder(绑定器)、inputs/outputs Channel完成应用程序和MQ的解耦。

  • Binder
    负责绑定应用程序和MQ中间件,即指定应用程序是和KafKa交互还是和RabbitMQ交互或者和其他的MQ中间件交互

  • inputs/outputs Channel
    inputs/outputs Channel抽象发布订阅消息的方式,即无论是什么类型的MQ应用程序都通过统一的方式发布订阅消息

我们已经搭建好了Kafka(参考Kafka单节点安装),本文主要介绍一下Spring Cloud Stream与Kafka进行集成实现消息的生产及消费。

项目创建

首先需要创建一个SpringBoot项目,命名为:spring-integration-kafka,在配置文件中导入相关的依赖。
项目情况为:

  • 构建工具:Gradle
  • SpringBoot版本:2.7.5
  • SpringBoot依赖管理版本:1.0.15.RELEASE
  • SpringCloud依赖管理版本:2021.0.5

项目依赖

配置文件build.gradle.kts的关键配置项如下:

plugins {
    id("org.springframework.boot") version "2.7.5"
    id("io.spring.dependency-management") version "1.0.15.RELEASE"
}

apply(plugin = "org.springframework.boot")
apply(plugin = "io.spring.dependency-management")
apply(plugin = "java")


extra["springCloudVersion"] = "2021.0.5"

dependencyManagement {
    imports {
        mavenBom("org.springframework.cloud:spring-cloud-dependencies:${property("springCloudVersion")}")
    }
}

dependencies {
	implementation("org.springframework.boot:spring-boot-starter-web")
	implementation("org.springframework.boot:spring-boot-starter-actuator")
	
    implementation("org.springframework.cloud:spring-cloud-starter-bootstrap")
    implementation("org.springframework.cloud:spring-cloud-starter-stream-kafka")
    
    implementation("io.springfox:springfox-boot-starter:3.0.0")
    implementation("com.github.xiaoymin:swagger-bootstrap-ui:1.9.6")
}

集成配置

定义配置文件application.yml,配置文件中主要配置Kafka的地址、以及Spring Colud Stream的Binder和inputs/outputs Channel,其中:kafkaChannel1用于向Kafka发送消息;kafkaChannel2用于消费Kafka的消息。

spring:
  kafka:
    bootstrap-servers: wux-labs-vm:9092  # 定义Kafka的地址
    producer:
      acks: 1
  cloud:
    stream:
      binders:
        kafkaBiner1:   # 定义一个Binder,名称随意
          type: kafka   # Binder的类型是 kafka
          environment:
            spring:
              kafka: ${spring.kafka}  # Binder的配置使用前面配置的Kafka的信息
      default-binder: kafkaBiner1    # 默认Binder,是前面配置的Binder的名称
      bindings:
        kafkaChannel1:      # 定义一个(作为outputs Channel)通道,名称随意,在代码中使用该通道名称即可
          binder: kafkaBiner1 # 使用kafkaBiner1
          destination: KafkaFirstTopic # 定义目标Topic的名称
        kafkaChannel2:       # 定义一个(作为inputs Channel)通道,名称随意,在代码中使用该通道名称即可
          binder: kafkaBiner1 # 使用kafkaBiner1
          destination: KafkaFirstTopic # 定义目标Topic的名称
          group: group0   # 作为消息的消费方,需要指定group

集成生产者

下面开发一个生产者,发送消息需要通过outputs Channel进行,使用kafkaChannel1发送消息到Kafka。

import io.swagger.annotations.Api;
import io.swagger.annotations.ApiOperation;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.cloud.stream.function.StreamBridge;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;

@Api("生产者接口")
@RestController
@RequestMapping("/producer")
public class ProducerController {
    @Autowired
    private StreamBridge bridge;

    @ApiOperation("向Kafka发送数据")
    @PostMapping("/kafka")
    public String sendToKafka(String message) {
        boolean status = bridge.send("kafkaChannel1", message);
        return "发送消息:" + message + "=====>" + status;
    }
}

集成消费者

消费消息需要通过inputs Channel进行,定义一个Processor,指定订阅通道为kafkaChannel2,这个通道被用于进行消息消费,需要定义group。

import org.springframework.cloud.stream.annotation.Input;
import org.springframework.messaging.SubscribableChannel;
import org.springframework.stereotype.Component;

@Component
public interface ConsumerProcessor {
    @Input("kafkaChannel2")
    SubscribableChannel subscribableChannel();
}

启用通道并监听。

import org.springframework.cloud.stream.annotation.EnableBinding;
import org.springframework.cloud.stream.annotation.StreamListener;

@EnableBinding(ConsumerProcessor.class)
public class ConsumerProcessorImpl {
    @StreamListener("kafkaChannel2")
    public void kafkaStreamListener(Object message) {
        System.out.println("接收到Kafka消息:" + new String((byte[]) message));
    }
}

集成验证

生产者验证

首先启动一个Kafka自带的消费者,监听KafkaFirstTopic

springboot集成kafka stream,# SpringCloud,# Kafka,kafka,java,java-rabbitmq

接下来启动SpringBoot项目并发送消息。在消费者那里可以看到接收到的消息。

springboot集成kafka stream,# SpringCloud,# Kafka,kafka,java,java-rabbitmq

消费者验证

前面消息已经发送到了Kafka的Topic了,可以看到控制台直接打印出了监听到的消息。

springboot集成kafka stream,# SpringCloud,# Kafka,kafka,java,java-rabbitmq
至此,Spring Cloud Stream集成Kafka完成。文章来源地址https://www.toymoban.com/news/detail-643671.html

到了这里,关于Spring Cloud Stream集成Kafka的文章就介绍完了。如果您还想了解更多内容,请在右上角搜索TOY模板网以前的文章或继续浏览下面的相关文章,希望大家以后多多支持TOY模板网!

本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若转载,请注明出处: 如若内容造成侵权/违法违规/事实不符,请点击违法举报进行投诉反馈,一经查实,立即删除!

领支付宝红包 赞助服务器费用

相关文章

  • Spring cloud stream 结合 rabbitMq使用

    之前的开发主要是底层开发,没有深入涉及到消息方面。现在面对的是一个这样的场景: 假设公司项目A用了RabbitMQ,而项目B用了Kafka。这时候就会出现有两个消息框架,这两个消息框架可能编码有所不同,且结构也有所不同,而且之前甚至可能使用的是别的框架,造成了一个

    2024年02月04日
    浏览(46)
  • SpringBoot 如何使用 Spring Cloud Stream 处理事件

    在分布式系统中,事件驱动架构(Event-Driven Architecture,EDA)已经成为一种非常流行的架构模式。事件驱动架构将系统中的各个组件连接在一起,以便它们可以相互协作,响应事件并执行相应的操作。SpringBoot 也提供了一种方便的方式来处理事件——使用 Spring Cloud Stream。 Spr

    2024年02月10日
    浏览(48)
  • 《微服务实战》 第十六章 Spring cloud stream应用

    第十六章 Spring cloud stream应用 第十五章 RabbitMQ 延迟队列 第十四章 RabbitMQ应用 https://github.com/spring-cloud/spring-cloud-stream-binder-rabbit 官方定义Spring Cloud Stream是一个构建消息驱动微服务的框架。应用程序通过inputs或者outputs来与Spring Cloud Stream中binder对象交互。通过我们配置来bindin

    2024年02月06日
    浏览(37)
  • 实战:Spring Cloud Stream消息驱动框架整合rabbitMq

    相信很多同学都开发过WEB服务,在WEB服务的开发中一般是通过缓存、队列、读写分离、削峰填谷、限流降级等手段来提高服务性能和保证服务的正常投用。对于削峰填谷就不得不用到我们的MQ消息中间件,比如适用于大数据的kafka,性能较高支持事务活跃度高的rabbitmq等等,MQ的

    2024年02月08日
    浏览(45)
  • Spring Cloud Stream 4.0.4 rabbitmq 发送消息多function

    注意当多个消费者时,需要添加配置项:spring.cloud.function.definition 启动日志 交换机名称对应: spring.cloud.stream.bindings.demo-in-0.destination配置项的值 队列名称是交换机名称+分组名 http://localhost:8080/sendMsg?delay=10000name=zhangsan 问题总结 问题一 解决办法: 查看配置是否正确: spring

    2024年02月19日
    浏览(42)
  • Spring Cloud Stream解密:流式数据在微服务中的魔力

    欢迎来到我的博客,代码的世界里,每一行都是一个故事 在微服务的大舞台上,数据流就像一曲美妙的交响乐,而Spring Cloud Stream正是指挥家,将音符有序地传递给每个微服务。在这篇文章中,我们将揭开Spring Cloud Stream的神秘面纱,一起探索在微服务体系结构中如何通过流式

    2024年02月20日
    浏览(39)
  • 【Spring云原生系列】SpringBoot+Spring Cloud Stream:消息驱动架构(MDA)解析,实现异步处理与解耦合!

    🎉🎉 欢迎光临,终于等到你啦 🎉🎉 🏅我是 苏泽 ,一位对技术充满热情的探索者和分享者。🚀🚀 🌟持续更新的专栏 《Spring 狂野之旅:从入门到入魔》 🚀 本专栏带你从Spring入门到入魔   这是苏泽的个人主页可以看到我其他的内容哦👇👇 努力的苏泽 http://suzee.blog.

    2024年03月10日
    浏览(49)
  • 【微服务学习】spring-cloud-starter-stream 4.x 版本的使用(rocketmq 版)

    @[TOC](【微服务学习】spring-cloud-starter-stream 4.x 版本的使用(rocketmq 版)) 2.1 消息发送者 2.1.1 使用 StreamBridge streamBridge; 往指定信道发送消息 2.1.2 通过隐式绑定信道, 注册 Bean 发送消息 2.2 消息接收者 注意: 多个方法之间可以使用 “|” 间隔, 但是绑定时 多个需要按顺序写. 其中

    2024年02月03日
    浏览(37)
  • SpringCloud(17)之SpringCloud Stream

            Spring Cloud Stream是一个框架,用于构建与共享消息系统连接的高度可扩展的事件驱动微服务。该框架提供了一个灵活的编程模型,该模型建立在已经建立和熟悉的Spring习惯用法和最佳实践之上,包括对持久发布/子语义、使用者组和有状态分区的支持。         

    2024年03月12日
    浏览(43)
  • 消息驱动 —— SpringCloud Stream

    Spring Cloud Stream 是用于构建消息驱动的微服务应用程序的框架,提供了多种中间件的合理配置 Spring Cloud Stream 包含以下核心概念: Destination Binders:目标绑定器,目标指的是 Kafka 或者 RabbitMQ,绑定器就是封装了目标中间件的包,如果操作的是 Kafka,就使用 Kafka Binder,如果操作

    2024年02月08日
    浏览(40)

觉得文章有用就打赏一下文章作者

支付宝扫一扫打赏

博客赞助

微信扫一扫打赏

请作者喝杯咖啡吧~博客赞助

支付宝扫一扫领取红包,优惠每天领

二维码1

领取红包

二维码2

领红包