CentOS 迁移 TencentOS Server 指引

操作系统版本 停止维护时间 使用者影响
CentOS 8 2022年01月01日 停止维护后将无法获得包括问题修复和功能更新在内的任何软件维护和支持。
CentOS 7 2024年06月30日

针对以上情况,若您需新购云服务器实例,建议选择使用 TencentOS Server 镜像。若您正在使用 CentOS 实例,则可参考本文替换为 TencentOS Server。

版本说明

源端主机支持操作系统版本:
支持 CentOS 7系列操作系统版本:
CentOS 7.2 64位、CentOS 7.3 64位、CentOS 7.4 64位、CentOS 7.5 64位、CentOS 7.6 64位、CentOS 7.7 64位、CentOS 7.8 64位、CentOS 7.9 64位。
支持 CentOS 8系列操作系统版本:
CentOS 8.0 64位、CentOS 8.2 64位、CentOS 8.3 64位、CentOS 8.4 64位、CentOS 8.5 64位。
目标主机建议操作系统版本:
CentOS 7系列建议迁移至 TencentOS Server 2.4 (TK4)。
CentOS 8系列建议迁移至 TencentOS Server 3.1 (TK4)。

注意事项

以下情况可能会影响业务在迁移后无法正常运行:
业务程序安装且依赖了第三方的 rpm 包。
迁移后的目标版本是 tkernel4,基于5.4的内核。该版本较 CentOS 7及 CentOS 8的内核版本更新,一些较旧的特性在新版本可能会发生变化。建议强依赖于内核的用户了解所依赖的特性,或可咨询 在线客服。
业务程序依赖某个固定的 gcc 版本。 目前 TencentOS Server 2.4默认安装 gcc 4.8.5,TencentOS Server 3.1默认安装 gcc 8.5。
迁移结束后,需重启才能进入TencentOS Server 内核。
迁移不影响数据盘,仅 OS 层面的升级,不会对数据盘进行任何操作。
注意:
操作系统迁移会将内核升级为基于 5.4 版本的 tkernel4 内核,因此可能下列情况的系统可能受到影响:

  1. 业务程序依赖于某个固定的内核版本,或者自行编译了内核模块,如 GPU 机型迁移后内核需要重新安装 GPU 驱动;
  2. 原操作系统的某个模块由 rpm 包提供,在迁移后此 rpm 包可能无法为新的内核提供模块,如:
    xpmem-modules-2.6.3-2.54310.kver.3.10.0_1160.108.1.el7.x86_64.x86_64 为 kernel-3.10.0-1160.108.1.el7.x86_64 提供 ko 文件, 但无法为迁移后的 tkernel4 内核提供。这种情况下用户可获取源码重新编译安装该模块。

    资源要求

    空闲内存大于500MB。
    系统盘剩余空间大于10GB。
    若/boot挂载分区,该分区空间需要大于500MB。

    操作步骤

    迁移准备

  3. 迁移操作不可逆,为保障业务数据安全,强烈建议您在执行迁移前通过 创建快照 备份系统盘数据。
  4. 操作系统迁移需要用户具有 root 权限。

    执行迁移

    CentOS 7系列迁移至 TencentOS Server 2.4(TK4)

    1. 登录目标云服务器,详情请参见 使用标准登录方式登录 Linux 实例。

    2. 执行以下命令,获取迁移工具。

    注意:
    若您的系统安装了旧版本的迁移工具,请卸载后再安装新的工具包。

    wget https://mirrors.cloud.tencent.com/tencentos/2.4/tlinux/x86_64/RPMS/migrate2tencentos-1.07-6.tl2.x86_64.rpm

    3. 执行以下命令,安装迁移工具。

    rpm -ivh migrate2tencentos-1.07-6.tl2.x86_64.rpm

    4. 执行以下命令,开始迁移。

    4.1 通过下面命令之一进行迁移
    4.1.1 全量迁移

    将 CentOS 发行版的用户态软件包替换为 TencentOS 发行版,为系统安装 TencentOS 自主研发的 tkernel4,基于5.4的内核。

    /usr/local/bin/EasyMigration -d remote -k

    image.png
    image.png
    最后 reboot 重启

4.1.2 minimal 软件组迁移

将系统的核心组件包迁移成 TencentOS 发行版,为系统安装 TencentOS 自主研发的 tkernel4,基于5.4的内核。
该模式下迁移的用户态软件包规模较小,系统上其他非核心组件的软件仍然保留为 CentOS 发行版。

/usr/local/bin/EasyMigration -d remote -k -g minimal

minimal 软件组默认列表参考 附录一。
迁移需要一定时间,请耐心等待。脚本执行完成后,输出如下图所示信息,表示已完成迁移。

5. 重启实例,详情请参见 重启实例。

6. 检查迁移结果。

6.1 执行以下命令,检查 os-release。

cat /etc/os-release

返回如下图所示信息:

6.2 执行以下命令,检查内核。

uname -r

返回如下图所示信息:

说明:
内核默认为 yum 最新版本,请以您的实际返回结果为准,本文以图示版本为例。

6.3 执行以下命令,检查 yum。

yum makecache

返回如下图所示信息:

若您在迁移过程中遇到问题,或对迁移有更多需求,请联系 在线客服。
Minimal 软件组列表见下表格:

序号 名称
1 audit
2 basesystem
3 bash
4 btrfs-progs
5 coreutils
6 cronie
7 curl
8 dhclient
9 e2fsprogs
10 filesystem
11 firewalld
12 glibc
13 hostname
14 initscripts
15 iproute
16 iprutils
17 iptables
18 iputils
19 irqbalance
20 kbd
21 kexec-tools
22 less
23 man-db
24 ncurses
25 openssh-clients
26 openssh-server
27 parted
28 passwd
29 plymouth
30 policycoreutils
31 procps-ng
32 rootfiles
33 rpm
34 rsyslog
35 selinux-policy-targeted
36 setup
37 shadow-utils
38 sudo
39 systemd
40 tar
41 tuned
42 util-linux
43 vim-minimal
44 xfsprogs
45 yum
46 NetworkManager
47 NetworkManager-team
48 NetworkManager-tui
49 aic94xx-firmware
50 alsa-firmware
51 biosdevname
52 dracut-config-rescue
53 ivtv-firmware
54 iwl100-firmware
55 iwl1000-firmware
56 iwl105-firmware
57 iwl135-firmware
58 iwl2000-firmware
59 iwl2030-firmware
60 iwl3160-firmware
61 iwl3945-firmware
62 iwl4965-firmware
63 iwl5000-firmware
64 iwl5150-firmware
65 iwl6000-firmware
66 iwl6000g2a-firmware
67 iwl6000g2b-firmware
68 iwl6050-firmware
69 iwl7260-firmware
70 kernel-tools
71 libsysfs
72 linux-firmware
73 lshw
74 microcode_ctl
75 postfix
76 sg3_utils
77 sg3_utils-libs
78 dracut-config-generic
79 dracut-fips
80 dracut-fips-aesni
81 dracut-network
82 initial-setup
83 openssh-keycat
84 rdma-core
85 selinux-policy-mls
86 tboot
87 gdb
88 kexec-tools
89 latrace
90 libreport-cli
91 strace
92 systemtap-runtime
93 abrt-addon-ccpp
94 abrt-addon-python
95 abrt-cli
96 crash
97 crash-gcore-command
98 crash-ptdump-command
99 crash-trace-command
100 elfutils
101 kernel-tools
102 libreport-plugin-mailx
103 ltrace
104 memstomp
105 ps_mem
106 trace-cmd
107 valgrind
108 abrt-java-connector
109 gdb-gdbserver
110 glibc-utils
111 memtest86+
112 systemtap-client
113 systemtap-initscrip

CentOS 8系列迁移至 TencentOS 3.1(TK4)

  1. 登录目标云服务器,详情请参见 使用标准登录方式登录 Linux 实例。

  2. 执行以下命令,获取迁移工具。
    注意:
    若您的系统曾经安装了旧版本的迁移工具,请卸载后再安装新的工具包。

    wget https://mirrors.cloud.tencent.com/tlinux/3.1/extras/x86_64/os/Packages/migrate2tencentos-1.07-6.tl3.x86_64.rpm
  3. 执行以下命令,安装迁移工具。

    rpm -ivh migrate2tencentos-1.07-6.tl3.x86_64.rpm
  4. 执行以下命令,开始迁移。
    4.1 通过下面命令之一进行迁移
    4.1.1 全量迁移
    将 CentOS 发行版的用户态软件包替换为 TencentOS 发行版,为系统安装 TencentOS 自主研发的 tkernel4,基于5.4的内核。

    /usr/local/bin/EasyMigration -d remote -k

    4.1.2 minimal 软件组迁移
    将系统的核心组件包迁移成 TencentOS 发行版,为系统安装 TencentOS 自主研发的 tkernel4,基于5.4的内核。
    该模式下迁移的用户态软件包规模较小,系统上其他非核心组件的软件仍然保留为 CentOS 发行版。

    /usr/local/bin/EasyMigration -d remote -k -g minimal

    minimal 软件组默认列表参考附录一。
    迁移需要一定时间,请耐心等待。脚本执行完成后,输出如下图所示信息,表示已完成迁移。

  5. 重启实例,详情请参见 重启实例。

  6. 检查迁移结果。
    6.1 执行以下命令,检查 os-release。

    cat /etc/os-release

    返回如下图所示信息:

    6.2 执行以下命令,检查内核。

uname -r

返回如下图所示信息:

说明:
内核默认为 yum 最新版本,请以您的实际返回结果为准,本文以图示版本为例。
6.3 执行以下命令,检查 yum。

yum makecache

返回如下图所示信息:

若您在迁移过程中遇到问题,或对迁移有更多需求,请联系 在线客服。
Minimal 软件组列表见下表格:

序号 名称
1 audit
2 basesystem
3 bash
4 btrfs-progs
5 coreutils
6 cronie
7 curl
8 dhclient
9 e2fsprogs
10 filesystem
11 firewalld
12 glibc
13 hostname
14 initscripts
15 iproute
16 iprutils
17 iptables
18 iputils
19 irqbalance
20 kbd
21 kexec-tools
22 less
23 man-db
24 ncurses
25 openssh-clients
26 openssh-server
27 parted
28 passwd
29 plymouth
30 policycoreutils
31 procps-ng
32 rootfiles
33 rpm
34 rsyslog
35 selinux-policy-targeted
36 setup
37 shadow-utils
38 sudo
39 systemd
40 tar
41 tuned
42 util-linux
43 vim-minimal
44 xfsprogs
45 yum
46 NetworkManager
47 NetworkManager-team
48 NetworkManager-tui
49 aic94xx-firmware
50 alsa-firmware
51 biosdevname
52 dracut-config-rescue
53 ivtv-firmware
54 iwl100-firmware
55 iwl1000-firmware
56 iwl105-firmware
57 iwl135-firmware
58 iwl2000-firmware
59 iwl2030-firmware
60 iwl3160-firmware
61 iwl3945-firmware
62 iwl4965-firmware
63 iwl5000-firmware
64 iwl5150-firmware
65 iwl6000-firmware
66 iwl6000g2a-firmware
67 iwl6000g2b-firmware
68 iwl6050-firmware
69 iwl7260-firmware
70 kernel-tools
71 libsysfs
72 linux-firmware
73 lshw
74 microcode_ctl
75 postfix
76 sg3_utils
77 sg3_utils-libs
78 dracut-config-generic
79 dracut-fips
80 dracut-fips-aesni
81 dracut-network
82 initial-setup
83 openssh-keycat
84 rdma-core
85 selinux-policy-mls
86 tboot
87 gdb
88 kexec-tools
89 latrace
90 libreport-cli
91 strace
92 systemtap-runtime
93 abrt-addon-ccpp
94 abrt-addon-python
95 abrt-cli
96 crash
97 crash-gcore-command
98 crash-ptdump-command
99 crash-trace-command
100 elfutils
101 kernel-tools
102 libreport-plugin-mailx
103 ltrace
104 memstomp
105 ps_mem
106 trace-cmd
107 valgrind
108 abrt-java-connector
109 gdb-gdbserver
110 glibc-utils
111 memtest86+
112 systemtap-client
113 systemtap-initscrip

Views: 33

CentOS 7原地迁移到版本 8 的 AlmaLinux、Rocky Linux、Oracle Linux

CentOS 官方计划停止维护 CentOS Linux 项目,CentOS 8及 CentOS 7维护情况如下表格。如需了解更多信息,请参见 CentOS 官方公告。

操作系统版本 停止维护时间 使用者影响
CentOS 8 2022年01月01日 停止维护后将无法获得包括问题修复和功能更新在内的任何软件维护和支持。
CentOS 7 2024年06月30日

针对以上情况
Elevate 是一个由 AlmaLinux 团队开发的开源项目,它允许将 CentOS 7 迁移到基于 RHEL 的较新和主要版本的发行版,例如 AlmaLinux 8、Rocky Linux 8、Oracle Linux 8 和 CentOS Stream 8。它结合了 RedHat 的 Leapp 框架带有一个社区开发的库来协助迁移。
本教学指南为您提供了使用 Elevate 将 CentOS 7 升级/迁移到 AlmaLinux 8 的步骤。
注意:Elevate仍处于开发的早期阶段,应该仅用于测试目的。不应该在生产服务器中测试迁移工具。
当前可用的迁移路径:

  • CentOS 7 到 AlmaLinux 8
  • CentOS 7 到 Rocky Linux 8
  • CentOS 7 到 Oracle Linux 8
  • CentOS 7 到 CentOS Stream 8

第 1 步:完全更新系统

首先,更新所有系统包和存储库。
[linuxmi@localhost www.linuxmi.com]$ sudo yum update -y
image.png
然后重启CentOS 7服务器。
[linuxmi@localhost www.linuxmi.com]$ sudo reboot

第 2 步:安装elevate-release包

下一步是安装 elevate-release 包,如下所示。
[linuxmi@localhost www.linuxmi.com]$ sudo yum install -y http://repo.almalinux.org/elevate/elevate-release-latest-el7.noarch.rpm
image.png
安装完成后,现在是时候为要迁移到的首选操作系统安装 Leapp 包和迁移数据了。迁移数据包的可能选项包括:

  • leapp-data-oraclelinux
  • leapp-data-almalinux
  • leapp-data-rocky
  • leapp-data-centos
  • leapp-data-oraclelinux

在我们的例子中,我们正在迁移到 AlmaLinux 8,因此,我们将安装leapp-data-almalinux 包。
[linuxmi@localhost www.linuxmi.com]$ sudo yum install -y leapp-upgrade leapp-data-almalinux
image.png

第 3 步:运行升级前检查

此后,启动升级前检查,如下所示。该命令会运行检查以查看升级是否成功,并提供有关在测试失败时您可以采取的可能补救措施的报告。
[linuxmi@localhost www.linuxmi.com]$ sudo leapp preupgrade
image.png
事实上,测试失败的原因有两到三个,这些原因记录在/var/log/leap /answerfile文件中,带有true/false的问题。有各种各样的建议可以解决无法升级的问题,但是,下面的建议是强制性的。
image.png
image.png
因此, 需要删除多余开发内核依赖以及回答 answerfile 文件
删除其余未使用的内核
image.png
sudo yum remove kernel-devel-3.10.0-1160.88.1.el7 -y
sudo yum remove kernel-devel-3.10.0-1160.108.1.el7 -y
sudo yum remove kernel-devel-3.10.0-1160.118.1.el7 -y
清除并重建缓存
sudo yum clean all
sudo yum makecache

移除已经加载到内核中的模块pata_acpi
[linuxmi@localhost www.linuxmi.com]$ sudo lsmod | grep pata_acpi
[linuxmi@localhost www.linuxmi.com]$ sudo rmmod pata_acpi
image.png
回答 answerfile 文件
[linuxmi@localhost www.linuxmi.com]$ sudo leapp answer --section remove_pam_pkcs11_module_check.confirm=True
image.png

升级还要求 sshd 服务设置允许通过 root 登录, 解决办法
[linuxmi@localhost www.linuxmi.com]$ echo PermitRootLogin yes | sudo tee -a /etc/ssh/sshd_config
PermitRootLogin yes

重启你的系统以确保所有的更改生效:
sudo reboot
再次执行预升级检查
$ sudo leapp preupgrade
image.png

第 4 步:从 CentOS 7 升级到 Almalinux 8

升级前首先备份关键数据
提前创建快照用于回滚

开始升级,请运行以下命令并重新启动系统
[linuxmi@localhost www.linuxmi.com]$ sudo leapp upgrade
image.png
image.png
[linuxmi@localhost www.linuxmi.com]$ sudo reboot
如果使用的时云服务器, 重启后切换 VNC 登录方式
image.png
可以发现升级过程还在继续
image.png
image.png
reboot 重启服务器

在重新启动过程中,将出现一个标有“Elevate-Upgrade-Initramfs”的新引导选项。选择此选项。

问题:

腾讯轻应用服务器, 重启后找不到引导项,尝试在菜单处按 c 进入 grub shell, 查找引导镜像
image.png
但是如何手动加载引导项暂时没找到办法...只能作罢。 并且进入Almalinux系统后输入用户名后直接显示认证失败(无法输入密码). 暂时找不到解决办法, 感觉这种办法不适合在云服务器上使用。最终根据官方指引 https://cloud.tencent.com/document/product/213/70900
从 Centos7 升级到了 TencentOS Server

如果以上没有出现问题, 升级将继续进行,大约需要 25 分钟。
最后,系统将再次重新启动。这次使用 AlmaLinux grub 菜单选项。
登录后,请验证您使用的操作系统版本。
[linuxmi@localhost www.linuxmi.com]$ cat /etc/redhat-release

就我而言,输出确认我已成功从 CentOS 7 升级到 AlmaLinux 8.4。就是这样。我希望本指南可以让你现在可以从 CentOS 7 无缝升级到任何基于 RHEL 8.x 的主要发行版,而不会出现问题。
来自:Linux迷
链接:https://www.linuxmi.com/centos-7update-almalinux-8-rocky-linux-8.html

Views: 46

SpringCloud与微服务-第9章 Spring Cloud Stream 消息驱动的微服务

旧文档(3.0.12.RELEASE):
https://docs.spring.io/spring-cloud-stream/docs/3.0.12.RELEASE/reference/html/

最新文档:
https://docs.spring.io/spring-cloud-stream/docs/current/reference/html/spring-cloud-stream.html#spring-cloud-stream-overview-introducing -->


Spring Cloud Stream 是一个用来为微服务应用构建消息驱动能力的框架。它可以基于 Spring Boot 来创建独立的,可用于生产的 Spring 应用程序。

它通过使用 Spring Integration 来连接消息代理中间件以实现消息事件驱动。Spring Cloud Stream 为一些供应商的消息中间件产品提供了个性化的自动化配置实现。


常见的消息中间件有:

  • RabbitMQ
    一个开源的 AMQP 实现, 服务器端用 Erlang 语言编写,支持多种客户端,如:Python、Ruby、.NET、Java、JMS、C、PHP、ActionScript、XMPP、STOMP 等,支持 AJAX、REST、SOAP 等多种通信协议。
  • Kafka
    一个分布式的基于发布/订阅模式的消息队列,主要应用于大数据实时处理领域。
  • ActiveMQ
    一个完全支持 JMS1.1 和 J2EE 1.4 规范的 JMS Provider 实现,尽管 JMS 规范出台已经是很久的事情了,但是 JMS 在当今的 J2EE 应用中间仍然扮演着特殊的地位。
  • RocketMQ
    一个分布式的消息中间件,具有低延迟、高性能、高可靠、亿级并发的特点,适用于大规模分布式系统的高性能场景。经受了阿里双十一的考验,具有丰富的实战经验。

简单地说,Spring Cloud Stream 本质上就是整合了 SpringBoot 和 Spring Integration, 实现了一套轻量级的消息驱动的微服务框架。

Spring Integration 是一个轻量级的消息代理框架,它的主要目的是为了在应用程序之间提供消息发送和接收的功能。


通过使用 Spring Cloud Stream ,可以有效简化开发人员对消息中间件的使用复杂度,让系统开发人员可以有更多的精力关注于核心业务逻辑的处理。

由于 Spring Cloud Stream 基于 Spring Boot 实现,所以它秉承了 Spring Boot 的优点,自动化配置的功能可帮助我们快速上手使用.


早期的版本比如 Spring Cloud Stream 3.0.12.RELEASE, 只支持 RabbitMQ 和 Kafka 两个著名的消息中间件的自动化配置。

最新版本的 Spring Cloud Stream 已经可以支持几乎所有的主流消息中间件,包括:RabbitMQ、Kafka、ActiveMQ、RocketMQ、Redis、Amazon Kinesis 等。


本章节将主要以 RabbitMQ 为例进行 Spring Cloud Stream 的内容讲解。


目标

在本章中,您将学习:

  • Spring Cloud Stream 快速入门
  • 核心概念
    • 绑定器 Binder
    • 发布-订阅模式 Publis-Subscribe Pattern
    • 消费组 Consumer Group
    • 分区 Partitioning
  • 使用详解
  • 绑定器详解
  • 配置详解

Spring Cloud Stream 快速入门

需求:

  1. 本地安装 RabbitMQ。
  2. 构建一个基于 Spring Boot 的微服务应用,
  3. 这个微服务应用将通过使用消息中间件 RabbitMQ 来接收消息并将消息打印到日志中。

RabbitMQ

消息队列 - MQ(message queue)从字面意思上看,本质是个队列,FIFO 先入先出,只不过队列中存放的是 message 而已,还是一种跨进程的通信机制,用于上下游传递消息。在互联网架构中,MQ 是一种非常常见的上下游"逻辑解耦+物理解耦"的消息通信服务。


通过消息队列可以实现应用程序之间的解耦,提高系统的可扩展性和可维护性。消息队列的应用场景非常多,比如异步处理、应用解耦、流量削锋、日志处理、消息通讯、消息广播等。


RabbitMQ 是 MQ 的实现之一, 具体来说它是一个开源的 AMQP 0-9-1(Advanced Message Queuing Protocol) 的实现。

h:14em


消息代理(Broker)从发布者(publishers)(发布消息的应用程序,也称为生产者)接收消息,并将其路由到消费者(处理这些消息的应用程序)。由于 AMQP 是一个网路协议, 所以发布者、代理和消费者可以在不同的进程中运行, 甚至在不同的主机上运行。

消息被发布到交换机(exchange),交换机可比作邮局或邮箱。然后,交换使用称为绑定(binding)的规则将消息副本分发到队列。然后,代理将消息传递给订阅队列的消费者,或者消费者根据需要从队列中获取/拉取消息。


RabbitMQ 有着运行在所有 Erlang 语言所支持的平台之上的潜力,从嵌入式系统到多核心集群还有基于云端的服务器。为了方便学习和演示,我们将在本地的 windows 系统上配置 RabbitMQ。


RabbitMQ 的安装方式,在其官方网站上有详细的介绍,包括如何下载资源,如何配置,如何启动等。为演示方便我们在 windows 系统上进行 RabbitMQ 安装。

官方文档地址: https://www.rabbitmq.com/install-windows.html 。这里我们下载rabbitmq-server-3.11.15.exe, 安装前请确保已经安装了 erlang。

可以从这里了解到你使用的 rabbitmq 需要安装什么 Erlang 版本.
https://www.rabbitmq.com/which-erlang.html

Erlang 安装: https://erlang.org/download/otp_versions_tree.html


在 Windows 安装 RabbitMQ 具体安装步骤如下:

  1. 下载 otp_win64_25.3.2.exe 并双击安装。

  2. 下载 rabbitmq-server-3.11.15.exe 并双击安装。

  3. 启动 RabbitMQ 服务
    rabbitmq_server-3.11.15\sbin>rabbitmq-server.bat start

  4. 激活监控插件
    rabbitmq_server-3.11.15\sbin>rabbitmq-plugins enable rabbitmq_management


安装成功之后,在浏览器中访问:http://127.0.0.1:15672 (RabbitMQ 的默认 UI 端口号是 15672, 默认本地登录账户 guest, 密码也是 guest):

file


之后点击 login 按钮,进入主界面:

w:35em


稍后,我们就可以在 RabbitMQ 主界面进行一些操作,现在让我们准备开始修改我们之前的两个微服务 orderservice 和 userservice,使用 Spring Cloud Stream 模拟一次简单的异步通讯吧。

异步: 一方发送消息,另一方接收消息,发送方不需要等待接收方的响应,而是继续执行后续的操作。


消息生产者

我们使用 userservice 服务作为消息的生产者,模拟一次消息的生产和发送。


1.引入依赖

引入 Spring Cloud Stream 整合 RabbitMQ 的依赖:

<!--Spring Cloud Stream 的rabbitmq整合依赖-->
<dependency>
    <groupId>org.springframework.cloud</groupId>
    <artifactId>spring-cloud-starter-stream-rabbit</artifactId>
</dependency>

注意: 由于 Spring Cloud Stream 依赖于 Spring Boot,所以我们不需要再引入 Spring Boot 的依赖。但是如果需要创建控制器,则需要引入 Spring Boot 的 web 依赖。


2.修改配置文件

对配置文件进行修改,添加 Spring Cloud Stream 与 RabbitMQ 的相关属性,稍后我们会对这些属性进行解读:


spring:
  cloud:
    stream:
      binders: #在此处配置要绑定的rabbitmq的服务信息
        defaultRabbit: #表示定义绑定名称
          type: rabbit #消息组件类型是RabbitMQ
          environment: #rabbitmq的相关的环境配置
            spring:
              rabbitmq:
                host: localhost
                port: 5672 ## rabbitmq的服务器端口
                username: guest
                password: guest
      bindings: #exchange将路由到队列的规则
        output: #这个名字是一个通道的名称
          destination: niitExchage #表示要使用的交换机的名称定义
          content-type: application/json #设置消息类型
          binder: defaultRabbit #设置要绑定的消息服务的具体设置
          group: niit #消息分组

binding 是一个接口,用于声明输入和输出通道 12。每个通道可以绑定到一个外部的消息代理(如 RabbitMQ 或 Kafka)上,通过 Binder 实现来连接。

binders 是一个抽象,用于实现不同类型的消息代理的连接逻辑。Spring Cloud Stream 提供了一些默认的 Binder 实现,如 RabbitMQ 和 Kafka,也可以自定义 Binder 实现。

一个应用可以使用多个 Binder 实现来连接不同类型的消息代理,但是需要在配置文件中指定每个通道使用哪个 Binder 实现。


3.启动类添加@EnableBinding 注解

创建一个消息生产的接口,定义发送消息的方法签名:

@EnableBinding({Source.class})
@MapperScan("com.niit.user.mapper")
@SpringBootApplication
public class UserApplication {

    public static void main(String[] args) {
        SpringApplication.run(UserApplication.class, args);
    }

}

不是必须


4.实现消息生产服务业务逻辑

创建一个实现类,用来具体实现接口业务逻辑方法:

@EnableBinding(Source.class)
public class UserMessageService {
    @Resource
    private MessageChannel output;

    @Override
    public String send(String msg) {
        Message<String> build = MessageBuilder.withPayload(msg).build();
        boolean sendFlag = output.send(build);
        if (sendFlag) {
            return "消息发送成功: "+msg;
        }
        return "消息发送失败: " + msg;
    }
}

5.创建控制器,对外暴露接口

创建一个 Controller,可以让外部 HTTP 请求访问内部业务,从而向通道中发送消息:

@RestController
@RequestMapping("/message")
public class UserMessageController {
    @Autowired
    private UserMessageService messageService;

    @GetMapping("/rabbitmq/{msg}")
    public String sendMessage(@PathVariable("msg") String msg){
      return   messageService.send(msg);
    }
}

消息消费者

我们使用 orderservice 服务作为消息的消费者,模拟一次消息的接收和消费。


1.引入依赖

引入 Spring Cloud Stream 整合 RabbitMQ 的依赖:

<!--Spring Cloud Stream 的rabbitmq整合依赖-->
<dependency>
    <groupId>org.springframework.cloud</groupId>
    <artifactId>spring-cloud-starter-stream-rabbit</artifactId>
</dependency>

2.修改配置文件

对配置文件进行修改,添加 Spring Cloud Stream 与 RabbitMQ 的相关属性,稍后我们会对这些属性进行解读:


spring:
  cloud:
    stream:
      binders: #在此处配置要绑定的rabbitmq的服务信息
        defaultRabbit: #表示定义的名称,用于与binding整合
          type: rabbit #消息组件类型是RabbitMQ
          environment: #rabbitmq的相关的环境配置
            spring:
              rabbitmq:
                host: localhost
                port: 5672 ## rabbitmq的服务器默认端口
                username: guest
                password: guest
      bindings: #服务的整合处理
        input: #这个名字是一个通道的名称
          destination: niitExchage #表示要使用的交换机的名称定义
          content-type: application/json #设置消息类型
          binder: defaultRabbit #设置要绑定的消息服务的具体设置
          group: niit #消息分组

3.启动类添加@EnableBinding 注解
@EnableBinding({Sink.class})
...
public class OrderApplication {...}

不是必须


4.创建 Java 类,接收并消费消息
@EnableBinding({Sink.class})
public class MessageReceiver {
    @StreamListener(Sink.INPUT) // 注解监听队列, 用于消费者的队列的消息接收
    public void receiveMessage(Message<String> msg){
        System.out.println("order-service接收到的消息是: " + msg.getPayload());
    }
}

核心注解

现在介绍刚才出现的两个 Spring Cloud Stream 核心注解:

  • @EnableBinding: 通过 @EnableBinding (Sink.class) 绑定了 Sink 接口,该接口是 Spring Cloud Stream 中默认实现的对输入消息通道绑定的定义,它的源码如下:

    public interface Sink {
      String INPUT = "input";
    
      @Input("input")
      SubscribableChannel input();
    }

  • @EnableBinding: 通过 @EnableBinding (Source.class) 绑定了 Source 接口,该接口是 Spring Cloud Stream 中默认实现的对输出消息通道绑定的定义,它的源码如下:

    public interface Source {
      String OUTPUT = "output";
    
      @Output("output")
      MessageChannel output();
    }

  • @EnableBinding: 通过 @EnableBinding (Processor.class) 绑定了 Processor接口, 它是一个消息通道的定义,它继承了 Sink 接口,同时还继承了 Source 接口, 作用是既可以作为消息的生产者,也可以作为消息的消费者,它的源码如下:

    public interface Processor extends Source, Sink {}

通过 @Input 注解绑定了一个名为 input 的通道。除了 Sink 之外, Spring Cloud Stream 还默认实现了绑定 output 通道的 Source 接口,还有结合了 Sink 和 Source 的 Processor 接口,实际使用时也可以自己通过 @Input 和 @Output 注解来定义绑定消息通道的接口。

在 Spring Cloud Stream 3.1 版本中,@EnableBinding 注解已被弃用,取而代之的是功能编程模型。


  • @StreamListener :主要定义在方法上,作用是将被修饰的方法注册为消息中间件上数据流的事件监听器,注解中的属性值对应了监听的消息通道名。在上面的例子中,通过@StreamListener (Sink.INPUT) 注解将 receiveMessage 方法注册为 input 消息通道的监听处理器,所以在 RabbitMQ 的控制页面中发布消息的时候, receiveMessage 方法会做出对应的响应动作。

测试

同时启动生产者和消费者

启动 Nacos 服务,然后分别启动 userservice 服务和 orderservice 服务。现在,我们访问: http://localhost:8081/message/rabbitmq/helloworld ,可以在 web 浏览器页面看到消息是否发送成功:

image-20220804095142675


多发送几次请求之后,我们可以在 RabbitMQ 的控制台页面看到消息流量的波动图:

h:14em


此时,我们在 orderservice 服务的控制台上,也可以观察到我们接收的消息:

image-20220804095507894


先启动生产者后启动消费者

上面的测试步骤,还不能直观的体现出异步消息的特性,我们现在模拟现实中的异步情境。我们可以先把消息消费者 orderservice 服务停掉,然后让消息生产者 userservice 服务多生产一些消息,发送到消息队列中,随后再启动消息消费者 orderservice 服务,查看 orderservice 服务的控制台输出日志:


image-20220804100314693


由此我们可以发现,当消费者服务重启之后,可以重新消费之前没有消费的消息。

需要指定消息分组(group)才有消息持久化的效果


通过 Spring Cloud Stream 框架驱动 RabbitMQ 消息中间件,解除了 userservice 和 orderservice 之间的耦合性,并且屏蔽了 RabbitMQ 底层复杂 API 的使用。

通过上一节介绍的快速入门示例,相信同学们对 Spring Cloud Stream 的工作模式己经有了一些基础概念,比如输入、输出通道的绑定,通道消息事件的监听等。在本节中,将详细介绍在 Spring Cloud Stream 中是如何通过定义一些基础概念来对各种不同的消息中间件做抽象的。


官网提供的 Spring Cloud Stream 的模型结构图。从中可以看到, Spring Cloud Stream 构建的应用程序与消息中间件之间是通过绑定器 Binder 相关联的,绑定器对于应用程序而言起到了隔离作用,它使得不同消息中间件的实现细节对应用程序来说是透明的。

bg right fit


所以对于每一个 Spring Cloud Stream 的应用程序来说,它不需要知晓消息中间件的通信细节,它只需知道 Binder 对应程序提供的抽象概念来使用消息中间件来实现业务逻辑即可,而这个抽象概念就是在快速入门中提到的消息通道: Channel 。如图所示,在应用程序和 Binder 之间定义了两条输入通道和三条输出通道来传递消息,而绑定器则是作为这些通道和消息中间件之间的桥梁进行通信。


绑定器

Binder 绑定器是 Spring Cloud Stream 中一个非常重要的概念。

绑定器的作用: 屏蔽底层消息中间件的差异, 降低消息中间件切换成本,统一消息的编程模型。


在没有绑定器这个概念的情况下,Spring Boot 应用要直接在业务代码中与消息中间件进行信息交互的时候,由于各消息中间件构建的初衷不同,所以它们在实现细节上会有较大的差异,这使得实现的消息交互逻辑就会非常笨重。

因为对具体的中间件实现细节有太重的依赖,当中间件有较大的变动升级或是更换中间件的时候,就需要付出非常大的代价来实施, 比如修改业务中的相关代码。


通过定义绑定器作为中间层,完美地实现了应用程序与消息中间件细节之间的隔离。通过向应用程序暴露统一的 Channel 通道,使得应用程序不需要再考虑各种不同的消息中间件的实现。当需要升级消息中间件,或是更换其他消息中间件产品时,要做的就是更换它们对应的 Binder 绑定器而不需要修改任何 Spring Boot 的应用逻辑。


早期 3.x 版本的 Spring Cloud Stream 只为 RabbitMQ 和 Kafka 提供了默认的 Binder 实现,在快速入门的例子中,就使用了 RabbitMQ 的 Binder。

bg right fit


另外, Spring Cloud Stream 还实现了一个专门用于单元测试的 TestSupportBinder, 开发者可以直接使用它来对通道的接收内容进行可靠的测试断言。

如果要使用除了 RabbitMQ 和 Kafka 以外的消息中间件的话,也可以过使用它所提供的扩展 API 来自行实现其他中间件的 Binder.

或者也可以使用较新版本的 Spring Cloud Stream 从而自动实现对其他消息中间件的支持。


相信大家己经发现,在快速入门示例中,使用 application.yml 做了一些属性设置。

当然,我们还可以通过 Spring Boot 应用支持的任何方式来修改这些配置,比如,通过应用程序参数、环境变量、 application.properties 或是 application.yaml 配置文件等。


发布-订阅模式

Spring Cloud Stream 中的消息通信方式遵循了发布 - 订阅模式,当一条消息被投递到消息中间件之后,它会通过共享的 Topic 主题进行广播,消息消费者在订阅的主题中收到它并触发自身的业务逻辑处理。这里所提到的 Topic 主题是 Spring Cloud Stream 中的一个抽象概念,用来代表发布共享消息给消费者的地方。


在不同的消息中间件中, Topic 可能对应不同的概念,比如,在 RabbitMQ 中,它对应 Exchange ,而在 Kakfa 中则对应 Kafka 中 的 Topic 。

在快速入门的示例中,通过 RabbitMQ 的 Channel 发布消息给我们编写的应用程序消费,而实际上 Spring Cloud Stream 应用启动的时候,在 RabbitMQ 的 Exchange 中也创建了一个名为 niitExchange 的交换器,

由于 Binder 的隔离作用,应用程序并无法感知它的存在,应用程序只知道自己指向 Binder 的输入或是输出通道。


在 Exchanges 选项卡中,还能找到名为 niitExchange 的交换器,单击进入可以看到如下图所示的详情页面。


bg fit


可以通过 Exchange 页面的 Publish Message 来发布消息:


bg fit


此时可以发现消费者服务接收到了消息:


bg fit


发布 - 订阅模式

  • 相对于点对点队列实现的消息通信来说, Spring Cloud Stream 采用的发布 - 订阅模式可以有效降低消息生产者与消费者之间的耦合。当需要对同一类消息增加一种处理方式时,只需要增加一个应用程序并将输入通道绑定到既有的 Topic 中就可以实现功能的扩展,而不需要改变原来己经实现的任何内容。

bg right fit


消费组

虽然 Spring Cloud Stream 通过发布 - 订阅模式将消息生产者与消费者做了很好的解耦,基于相同主题的消费者可以轻松地进行扩展,但是这些扩展都是针对不同的应用实例而言的。在现实的微服务架构中,每一个微服务应用为了实现高可用和负载均衡,实际上都会部署多个实例。在很多情况下,消息生产者发送消息给某个具体微服务时,只希望被消费一次。但是同一个应用的多个实例都会接收到消息,这个消息将会被重复消费。


为了解决这个问题,在 Spring Cloud Stream 中提供了消费组的概念。如果在同一个主题上的应用需要启动多个消费者实例的时候,可以通过 spring.cloud.stream.bindings.input.group 属性为应用指定一个组名,这样这个应用的多个实例在接收到消息的时候,只会有一个成员真正收到消息并进行处理。


如下图所示,为 Service-A 和 Service-B 分别启动了两个实例,并且根据服务名进行了分组,这样当消息进入主题之后, Group-A 和 Group-B 都会收到消息的副本,但是在两个组中都只会有一个实例对其进行消费。

bg right fit


默认情况下,当没有为应用指定消费组的时候, Spring Cloud Stream 会为其分配一个独立的匿名消费组。所以,如果同一主题下的所有应用都没有被指定消费组的时候,当有消息发布之后,所有的应用都会对其进行消费,因为它们各自都属于一个独立的组。如果这个消息根据业务要求只需要被消费一次,那么就会出现重复消费的问题。

匿名消费者组订阅的消息是不会持久化的, 因此是不可靠的.


也可以为多个消费者指定相同的组, 可以逻辑上看做是一个消费者. 这样可以实现负载均衡和故障转移的效果, 因为每个消息只会被消费者组中的成员消费一次. 具体由哪个消费者消费, 由消息中间件决定.

指定的消费者组订阅的消息会持久化, 因此更加可靠.


消息分区

通过引入消费组的概念,己经能够在多实例的情况下,保障每个消息只被组内的一个实例消费。但是消费组无法控制消息具体被哪个实例消费。也就是说,对于同一条消息,它多次到达之后可能是由不同的实例进行消费的。


但是对于一些业务场景,需要对一些具有相同特征的消息设置每次都被同一个消费实例处理,比如,一些用于监控服务,为了统计某段时间内消息生产者发送的报告内容,监控服务需要在自身聚合这些数据,


那么消息生产者可以为消息在消息的 header 中增加一个固有的特征 ID 来进行分区,使得拥有这些 ID 的消息每次都能被发送到一个特定的实例上实现累计统计的效果,否则这些数据就会分散到各个不同的节点导致监控结果不一致的情况。


而分区概念的引入就是为了解决这样的问题:当生产者将消息数据发送给多个消费者实例时,保证拥有共同特征的消息数据始终是由同一个消费者实例接收和处理。


Spring Cloud Stream 为分区提供了通用的抽象实现,用来在消息中间件的上层实现分区处理,所以它对于消息中间件自身是否实现了消息分区并不关心,这使得 Spring Cloud Stream 为不具备分区功能的消息中间件也增加了分区功能扩展。


小问题 :
什么是 Spring Cloud Stream 中一个非常重要的概念?

  1. Spring 绑定器
  2. Prop 绑定器
  3. Binder 绑定器
    查看答案

    正确答案:3


使用详解

在介绍了 Spring Cloud Steam 的基础结构和核心概念之后,我们来详细地学习一下它所提供的一些核心注解的具体使用方法。


开启绑定功能

在 Spring Cloud Stream 中,需要通过 @EnableBinding 注解来为应用启动消息驱动的功能,该注解在快速入门中己经有了基本的介绍,下面来详细看看它的定义:


@Target({ElementType.TYPE, ElementType.ANNOTATION_TYPE})
@Retention(RetentionPolicy.RUNTIME)
@Documented
@Inherited
@Configuration
@Import({BindingBeansRegistrar.class, BinderFactoryAutoConfiguration.class})
@EnableIntegration
public @interface EnableBinding {
    Class<?>[] value() default {};
}

从该注解的定义中可以看到,它自身包含了 @Configuration 注解,所以用它注解的类也会成为 Spring 的基本配置类。另外该注解还通过 @Import 加载了 Spring Cloud Stream 运行需要的几个基础配置类。

  • BindingBeansRegistrar :该类是 ImportBeanDefinitionRegistrar 接口的实现,主要是在 Spring 加载 Bean 的时候被调用,用来实现加载更多的 Bean 。由于 BindingBeansRegistrar 被 @EnableBinding 注解的 @Import 所引用,所以在其他配置加载完后,它的实现会被回调来创建其他的 Bean, 而这些 Bean 则从 @EnableBinding 注解的 value 属性定义的类中获取。就如入门实例中定义的 @EnableBinding (Sink.class) ,它在加载用于消息驱动的基础 Bean 之后, 会继续加载 Sink 中定义的具体消息通道绑定。

  • BinderFactoryConfiguration : Binder 工厂的配置,主要用来加载与消息中间件相关的配置信息,比如,它会从应用工程的 META-INF/spring.binders 中 加载针对具体消息中间件相关的配置文件等。


@EnableBinding 注解只有一个唯一的属性: value 。上面己经介绍过,由于该注解 @Import 了 BindingBeansRegistrar 实现,所以在加载了基础配置内容之后,它会回调来读取 value 中的类,以创建消息通道的绑定。另外,由于 value 是一个 Class 类型的数组,所以可以通过 value 属性一次性指定多个关于消息通道的配置。


绑定消息通道

在 Spring Cloud Steam 中,可以在接口中通过 @Input 和 @Output 注解来定义消息通道,而用于定义绑定消息通道的接口则可以被 @EnableBinding 注解的 value 参数来指定,从而在应用启动的时候实现对定义消息通道的绑定。


bg left w:36em

w:18em


在快速入门的示例中,演示了使用 Sink 接口绑定的消息通道。 Sink 接口是 Spring Cloud Steam 提供的一个默认实现,除此之外还有 Source 和 Processor ,可从它们的源码中学习它们的定义方式:


//Sink
public interface Sink {
    String INPUT = "input";

    @Input("input")
    SubscribableChannel input();
}

//Source
public interface Source {
    String OUTPUT = "output";

    @Output("output")
    MessageChannel output();
}

//Processor
public interface Processor extends Source, Sink {
}

从上面的源码中,可以看到, Sink 和 Source 中分别通过 @Input 和 @Output 注解定义了输入通道和输出通道,而 Processor 通过继承 Source 和 Sink 的方式同时定义了一个输入通道和一个输出通道。
bg right fit


另外, @Input 和 @Output 注解都还有一个 value 属性,该属性可以用来设置消息通道的名称,这里 Sink 和 Source 中指定的消息通道名称分别为 input 和 output。如果直接使用这两个注解而没有指定具体的 value 值,将默认使用方法名作为消息通道的名称。


最后,需要注意一点,

  • 当定义输出通道的时候,需要返回 MessageChannel 接口对象,该接口定义了向消息通道发送消息的方法;

  • 而定义输入通道时,需要返回 SubscribableChannel 接口对象,

    • 该接口继承自 MessageChannel 接口,它定义了维护消息通道订阅者的方法。

注入绑定接口

在完成了消息通道绑定的定义之后, Spring Cloud Stream 会为其创建具体的实例,而开发者只需要通过注入的方式来获取这些实例并直接使用即可。举个简单的例子,在快速入门示例中 orderservice 服务己经为 Sink 接口绑定的 input 消息通道实现了具体的消息消费者,下面可以通过注入的方式实现一个消息生成者,向 input 消息通道发送数据。


  1. 创建一个将消息通道"input"作为输出通道的发送消息的接口,具体如下:

    public interface SinkSender {
       @Output(Sink.INPUT)
       MessageChannel output();
    }

  1. 对 orderservice 中 定义的 ReceiveMessage 做一些修改:在@Enablebinding 注解中增加对 SinkSender 接口的指定,使 Spring Cloud Stream 能创建出对应的 Java Bean 的实例。

    @EnableBinding({Sink.class,SinkSender.class})
    public class MessageReceiver {
    
     @StreamListener(Sink.INPUT)
       public void receiveMessage(Message<String> msg){
           System.out.println("order-service接收到的消息是: "+ msg.getPayload());
       }
    
    }

  • 创建 OrderMessageController 控制器类, 通过 @Autowired 注解注入 SinkSender 的实例,并在控制器方法中调用它的发送消息方法。

    @RestController
    @RequestMapping("/message")
    public class OrderMessageController {
    
      @Autowired
      private SinkSender sinkSender;
    
      @RequestMapping("/sinkSender/{msg}")
      public String contextLoads(@PathVariable("msg")String msg){
          boolean send = sinkSender.output().send(MessageBuilder.withPayload(msg).build());
          if (send) {
              return "成功发送消息: " + msg;
          }
          return "消息发送失败: " + msg;
      }
    }
    

  1. 重启 orderservice 服务,在浏览器中访问:http://localhost:8080/message/sinkSender/hello 。如果可以在控制台中找到如下输出内容,表明试验己经成功了,消息被正确地发送回 orderservice 服务的 input 通道中,并被消息消费者输出(等于是自己把消息发送给了自己)。

image-20220804222040951


注入消息通道进行消息发送

由于 Spring Cloud Stream 会根据绑定接口中的 @Input 和 @Output 注解来创建消息通道实例,所以也可以通过直接注入的方式来使用消息通道对象。比如,可以通过下面的示例,注入上面例子中 SinkSender 接口中定义的名为 input 的消息输入通道。


@RestController
@RequestMapping("/message")
@EnableBinding({Sink.class, SinkSender.class})
public class OrderMessageController {
    @Autowired
    private MessageChannel input;

    @RequestMapping("/inputChannel/{msg}")
    public String channelMessage(@PathVariable("msg")String msg){
        boolean send = input.send(MessageBuilder.withPayload(msg).build());
        if (send) {
            return "成功发送消息: " + msg;
        }
        return "消息发送失败: " + msg;
    }

}

上面定义的内容,完成了与之前通过注入绑定接口 SinkSender 方式实现的测试用例相同的操作。因为在通过注入绑定接口实现时, sinkSender.output() 方法实际获得的就是 SinkSender 接口中定义的 MessageChannel 实例,只是在这里直接通过注入的方式来实现了而己。

注入绑定接口 SinkSender这种用法虽然很直接,但是也容易犯错,很多时候在一个微服务应用中可能会创建多个不同名的 MessageChannel 实例,这样通过 @Autowired 注入时,要注意参数命名需要与通道同名才能被正确注入,或者也可以使@Qualifier 注解来特别指定具体实例的名称,该名称需要与定义 MessageChannel 的 @Output 中的 value 参数一致,这样才能被正确注入。


比如下面的例子,在一个接口中定义了两个输出通道,分别命名为 output1 和 output2, 当要使用 通道output1 的时候,可以通过
@Qualifier( "Output1") 来指定这个具体的实例来注入使用。


  • MySource 接口

    //定义通道
    public interface MySource {
      @Output("output1")
      MessageChannel output1();
    
      @Output("output2")
      MessageChannel output2();
    }
  • OutputSender 类

    //注入
    @Component
    public class OutputSender {
      @Autowired
      @Qualifier("output1")
      MessageChannel output;
    ...
    }
    

提示:
@EnableBinding 注解也可以只写在启动类上,Binding 组件及其通道只需要注册一次


注解 @StreamListener 详解

通过入门示例,对于 @StreamListener 注解,应该都己经有了一些基本的认识,通过该注解修饰的方法, Spring Cloud Steam 会将其注册为输入消息通道的监听器。当输入 消息通道中有消息到达的时候,会立即触发该注解修饰方法的处理逻辑对消息进行消费。
@SteamListener 注解都实现了对输入消息通道的监听,并且内置了一系列的消息转换功能,这使得基于 @SteamListener 注解实现的消息处理模型更为简单。


消息转换

大部分情况下,通过消息来对接服务或系统时,消息生产者都会以结构化的字符串形式来发送,比如 JSON 或 XML 。当消息到达的时候,输入通道的监听器需要对该字符串做一定的转化,将 JSON 或 XML 转换成具体的对象,然后再做后续的处理。


可以这样为通道设置绑定消息的类型, 这样传输消息时会对格式进行校验如:

  • cloud.stream.bindings.input.content-type=application/json
  • cloud.stream.bindings.output.content-type=application/json

假设,需要传输一个 User 对象,该对象有 username 和 address 两个字段,这时,如果使用 @SteamListener 注解的话,代码实现将变得非常简单优雅:

@EnableBinding(MySource.class)
public class MessageReceiver {

    /**
     * 监听output2通道,完成自动类型转换
     * @param user
     */
    @StreamListener("output2")
    public void receive(User user) {
        System.out.println("Received: " + user);
    }
}

我们可以在控制器中增加一个方法,使用我们刚刚定义的 output2 通道进行信息发送:

@RestController
@RequestMapping("/message")
public class OrderMessageController {

    @Autowired
    private MySource mySource;

    @RequestMapping("/userInfo/")
    public String output2ChannelMessage(User user){
        boolean send = mySource.output2().send(MessageBuilder.withPayload(JSONUtil.toJsonStr(user)).build());
        if (send) {
            return "成功发送消息: " + JSONUtil.toJsonStr(user);
        }
        return "消息发送失败: " + JSONUtil.toJsonStr(user);
    }
}

我们可以通过之前介绍的小插件: FastRequest,模拟一个 User 信息,然后发送:

image-20220804235921454


控制台上将会显示:

image-20220805000138729


消息反馈

很多时候在处理完输入消息之后,需要反馈一个消息给对方,这时候可以通过 @SendTo 注解来指定返回内容的输出通道。我们对 orderservice 中消息接收类 MessageReceiver 进行修改,在其中一个类上添加@SendTo 注解:

@EnableBinding({Processor.class, SinkSender.class,MySource.class})
public class MessageReceiver {

    @StreamListener(Sink.INPUT)
    @SendTo(Processor.OUTPUT)
    public String receiveMessage(Message<String> msg){
        System.out.println("order-service接收到的消息是: "+msg.getPayload());
        return msg.getPayload();
    }
}

在 userservice 中同样创建一个消息接收类,用来监听来自 orderservice 的消息回执:

@Component
public class UserMessageReceiver {

    @StreamListener(Processor.INPUT)
    public void receiveReturnMessage(Message<String> msg) {
        System.out.println("receiveReturnMessage 接收到的消息回执是: " + msg.getPayload());
    }
}

userservice 是 orderservice 应用中 input 通道的生产者以及 output 通道的消费者。可以在配置文件中将两个应用的通道绑定反向地做一些配置。因为对于 userservice 来说,它的 input 绑定通道实际上是对 output 主题的消费者,而 output 绑定通道实际上是对 input 主题的生产者,


做如下具体配置。指定通道的 destination 来实现两个应用的消息交互:

orderservice 的配置

spring:
  cloud:
    stream:
      binders: #在此处配置要绑定的rabbitmq的服务信息
        defaultRabbit: #表示定义的名称,用于与binding整合
          type: rabbit #消息组件类型是RabbitMQ
          environment: #rabbitmq的相关的环境配置
            spring:
              rabbitmq:
                host: localhost
                port: 5672
                username: guest
                password: guest
      bindings: #服务的整合处理
        input: #这个名字是一个通道的名称
          destination: niitExchange #表示要使用的交换机的名称定义
          content-type: application/json #设置消息类型
          binder: defaultRabbit #设置要绑定的消息服务的具体设置
          group: niit
        output:
          destination: commonExchange
          content-type: application/json
          binder: defaultRabbit
          group: niit

userservice 的配置

spring:
  cloud:
    stream:
      binders:  #在此处配置要绑定的rabbitmq的服务信息
        defaultRabbit:  #表示定义的名称,用于与binding整合
          type: rabbit  #消息组件类型是RabbitMQ
          environment:  #rabbitmq的相关的环境配置
            spring:
              rabbitmq:
                host: localhost
                port: 5672
                username: guest
                password: guest
      bindings: #服务的整合处理
        output:
          destination: niitExchange
          content-type: application/json
          binder: defaultRabbit
          group: niit
        input:
          destination: commonExchange
          content-type: application/json
          binder: defaultRabbit
          group: niit

当 receiveMessage 方法处理消息之后, return 的值将会最终被 UserMessageReceiver 类的 receiveReturnMessage 方法消费。


消费组与消息分区

在“核心概念” 一节中,对消费组和消息分区己经进行了基本的介绍,在这里来详细介绍一下这两个概念的使用方法。


消费组

通常每个服务都不会以单节点的方式运行在生产环境中,当同一个服务启动多个实例的时候,这些实例会绑定到同一个消息通道的目标主题上。默认情况下,当生产者发出一 条消息到绑定通道上,这条消息会产生多个副本被每个消费者实例接收和处理。但是在有些业务场景之下,希望生产者产生的消息只被其中一个实例消费,这个时候就需要为这些消费者设置消费组来实现这样的功能。实现的方式非常简单,只需在服务消费者端设置 spring.cloud.stream.bindings.input.group 属性即可。在上面的消息回执例子中,我们正是基于此种原因,所以设置了 group 的属性。


分别运行 orderservice 和 userservice,其中 userservice 启动多个实例。以消息回执案例进行测试,可以发现,消息的回执会被启动的多个 userservice 实例以轮询的方式进行接收和输出。

image-20220805172532126


消息分区

通过消费组的设置,虽然己经能够在多实例环境下,保证同一消息只被一个消费者实例进行接收和处理,但是,对于一些特殊场景,除了要保证单一实例消费之外,还希望那些具备相同特征的消息都能够被同一个实例进行消费。这时候我们就需要对消息进行分区处理。


在 Spring Cloud Stream 中实现消息分区非常简单,对消费组示例做一些配置修改就能实现,具体如下所示。


  1. 以 orderservice 作为消息生产者,修改其配置文件:

    spring:
     cloud:
       stream:
         binders: #在此处配置要绑定的rabbitmq的服务信息
           defaultRabbit: #表示定义的名称,用于与binding整合
             type: rabbit #消息组件类型是RabbitMQ
             environment: #rabbitmq的相关的环境配置
               spring:
                 rabbitmq:
                   host: localhost
                   port: 5672
                   username: guest
                   password: guest
         bindings: #服务的整合处理
           input: #这个名字是一个通道的名称
             destination: niitExchange #表示要使用的交换机的名称定义
             content-type: application/json #设置消息类型
             binder: defaultRabbit #设置要绑定的消息服务的具体设置
             group: niit
           output:
             destination: commonExchange
             content-type: application/json
             binder: defaultRabbit
             group: niit
             producer:
               partitionKeyExpression: payload
               partitionCount: 3

从上面的配置中,我们可以看到增加了下面这两个参数。

  • spring.cloud.stream.bindings.output.producer.partitionKeyExpression :通过该参数指定了分区键的表达式规则,可以根据实际的输出消息规则配置 SpEL 来生成合适的分区键。
    • payload : 根据类型自动分区
    • header: 设定消息头 ,根据消息头手动设置分区, 如 partitionKeyExpression: headers['partitionKey']表示根据消息头中的 partitionKey 的值来分区
  • spring.cloud.stream.bindings.output.producer.partitionCount :该参数指定了消息分区的数量。我们准备启动三个消费者实例,所以设置为 3 个分区。

  1. 以 userservice 作为消息消费者,修改其配置文件:

    spring:
     cloud:
       stream:
         binders:  #在此处配置要绑定的rabbitmq的服务信息
           defaultRabbit:  #表示定义的名称,用于与binding整合
             type: rabbit  #消息组件类型是RabbitMQ
             environment:  #rabbitmq的相关的环境配置
               spring:
                 rabbitmq:
                   host: localhost
                   port: 5672
                   username: guest
                   password: guest
         bindings: #服务的整合处理
         output:
             destination: niitExchange
             content-type: application/json
             binder: defaultRabbit
             group: niit
           input:
             destination: commonExchange
             content-type: application/json
             binder: defaultRabbit
             group: niit
             partitioned: true
         instance-count: 3
         instance-index: 0

从上面的配置中,可以看到增加了下面这三个参数。

  • spring.cloud.stream.bindings.input.consumer.partitioned :通过该参数开启消费者分区功能。
  • spring.cloud.stream.instanceCount :该参数指定了当前消费者的总实例数量。
  • spring.cloud.stream.instanceIndex :该参数设置当前实例的索引号,从 0 开始 。试验的时候需要启动多个实例,可以通过运行参数来为不同实例设置不同的索引值。

到这里消息分区配置就完成了,可以再次启动这两个应用,同时启动多个消费者。但需要注意的是,要为消费者指定不同的实例索引号,这样当同一个消息被发送给消费组时,可以发现只有一个消费实例在接收和处理这些相同的消息。


消息类型

Spring Cloud Stream 为了让开发者能够在消息中声明它的内容类型,在输出消息中定义了一个默认的头信息:contentType 。对于那些不直接支持头信息的消息中间件, Spring Cloud Stream 提供了自己的实现机制,它会在消息发出前自动将消息包装进它自定义的消息封装格式中,并加入头信息。而对于那些自身就支持头信息的消息中间件, Spring Cloud Stream 构建的服务可以接收并处理来自非 Spring Cloud Stream 构建但包含符合规范头信息的应用程序发出的消息。


Spring Cloud Stream 允许使用 spring.cloud.stream.bindngs.<channelName>.content-type 属性以声明式的配置方式为绑定的输入和输出通道设置消息内容的类型。 此外,原生的消息类型转换器依然可以轻松地用于我们的应用程序。


目前, Spring Cloud Stream 中自带支持了以下几种常用的消息类型转换。

  • JSON 与 POJO 的互相转换。

  • JSON 与 org.springframework.tuple.Tuple 的互相转换。

  • Object 与 byte[] 的互相转换。为了实现远程传输序列化的原始字节,应用程序需要发送 byte 类型的数据,或是通过实现 Java 的序列化接口来转换为字节( Object 对象必须可序列化)。

  • String 与 byte[] 的互相转换。

  • Object 向纯文本的转换: Object 需要实现 toString()方法。


上面所指的 JSON 类型可以表现为一个 byte 类型的数组,也可以是一个包含有效 JSON 内容的字符串。另外, Object 对象可以由 JSON 、 byte 数组或者字符串转换而来,但是在转换为 JSON 的时候总是以字符串的形式返回。


MIME 类型

在 Spring Cloud Stream 中定义的 content-type 属性采用了 Media Type ,即 Internet Media Type (互联网媒体类型),也被称为 MIME 类型,常见的有 application/json 、text/plain;charset=UTF-8, 相信接触过 HTTP 的工程师们对这些类型都不会感到陌生。


MIME 类型对于标示如何转换为 String 或 byte [] 非常有用。并且,还可以使用 MIME 类型格式来表示 Java 类型,只需要使用带有类型参数的一般类型: application/x-java-object 。

比如,我们可以使用 application/x-java-object;type=java.util.Map 来表示传输的是一个 java.util.Map 对象,或是使用 application/x-java-object;type=com.niit.pojo.User 来表示传输的是一个 com.niit.pojo.User 对象;除此之外,更重要的是,它还提供了自定义的 MIME 类型, 比如通过 application/x-spring-tuple 来指定 Spring 的 Tuple 类型。


在 Spring Cloud Stream 中默认提供了一些可以开箱即用的类型转换器,具体如下表所示。


h:18em


消息类型的转换行为只会在需要进行转换时才被执行,比如,当服务模块产生了一个头信息为 application/json 的 XML 字符串消息, Spring Cloud Stream 是不会将该 XML 字符串转换为 JSON 的,这是因为该模块的输出内容己经是一个字符串类型了,所以它并不会将其做进一步的转换。


另外需要注意的是, Spring Cloud Stream 虽然同时支持输入通道和输出通道的消息类 型转换,但还是推荐开发者尽量在输出通道中做消息转换。因为对于输入通道的消费者来说,当目标是一个 POJO 的时候,使用 @StreamListener 注解是能够支持自动对其进行转换的。


Spring Cloud Stream 除了提供上面这些开箱即用的转换器之外,还支持开发者自定义的消息转换器。这使得用户可以使用任意格式(包括二进制)的数据进行发送和接收,并且将这些数据与特定的 contentType 相关联。在应用启用的时候,Spring Cloud Stream 会将所有 org.springframework.messaging.converter.MessageConverter 接口实现的自定义转换器以及默认实现的那些转换器都加载到消息转换工厂中,以提供给消息处理时使用。


绑定器详解

在“核心概念”一节中,己经简单介绍过 Binder 绑定器的基本概念和作用:它是定义在应用程序与消息中间件之间的抽象层,用来屏蔽消息中间件对应用的复杂性,并提供简单而统一的操作接口给应用程序使用。在本节中,将详细介绍绑定器背后的细节和行为。


绑定器 SPI (了解)

绑定器 SPI 涵盖了一套可插拔的用于连接外部中间件的实现机制,其中包含了许多接口、开箱即用的实现类以及发现策略等内容。其中,最为关键的就是 Binder 接口,它是用来将输入和输出连接到外部中间件的抽象:

package org.springframework.cloud.stream.binder;

public interface Binder<T, C extends ConsumerProperties, P extends ProducerProperties> {
    Binding<T> bindConsumer(String name, String group, T inboundBindTarget, C consumerProperties);

    Binding<T> bindProducer(String name, T outboundBindTarget, P producerProperties);
}

当应用程序对输入和输出通道进行绑定的时候,实际上就是通过该接口的实现来完成的。

  • 向消息通道发送数据的生产者调用 bindProducer 方法来绑定输出通道时,第一个参数代表了发往消息中间件的目标名称,第二个参数代表了发送消息的本地通道实例,第三个参数是用来创建通道时使用的属性配置 ( 比如分区键的表达式等 ) 。

  • 从消息通道接收数据的消费者调用 bindConsumer 方法来绑定输入通道时,第一个参数代表了接收消息中间件的目标名称,第二个参数代表了消费组的名称 ( 如果多个消费者实例使用相同的组名,则消息将对这些消费者实例实现负载均衡,每个生产者发出的消息只会被组内一个消费者实例接收和处理 ) ,第三个参数代表了接收消息的本地通道实例,第四个参数是用来创建通道时使用的属性配置。


另外,从 Binder 的定义中,还可以知道 Binder 是一个参数化并且可扩展的接口。

  • 对于输入与输出的绑定类型,在 1.0 版本中仅支持 MessageChannel ,但是在接口中通过泛型定义,所以在未来可以对其进行扩展。
  • 对于属性配置也提供了可扩展的定义,可以为特定的 Binder 以类型安全的方式来补充一些特有的属性。

一个典型的 Binder 绑定器实现一般包含以下内容。

  • 一个实现 Binder 接口的类。
  • 一个 Spring 配置加载类,用来创建连接消息中间件的基础结构使用的实例。
  • 一个或多个能够在 classpath 下的 META-INF/spring.binders 路径找到的绑定器定义文件。比如用户可以在 spring-cloud-starter-stream-rabbit 中找到该文件,该文件中存储了当前绑定器要使用的自动化配置类的路径:

image-20220805203723402


绑定器的自动化配置

Spring Cloud Stream 通过绑定器 SPI 的实现将应用程序逻辑上的输入输出通道连接到物理上的消息中间件。消息中间件之间通常都会有或多或少的差异性,所以为了适配不同的消息中间件,需要为它们实现各自独有的绑定器。目前, Spring Cloud Stream 中默认实现了对 RabbitMQ 、 Kafka 的绑定器,在上面的示例中引入的 spring-cloud- starter- stream-rabbit 依赖中就包含了 RabbitMQ 的绑定器 spring-cloud-stream-binder- rabbit 。


默认情况下, Spring Cloud Stream 也遵循 Spring Boot 自动化配置的特性。如果在 classpath 中能够找到单个绑定器的实现,那么 Spring Cloud Stream 会自动加载它。而在 classpath 中引入绑定器的方法也非常简单,只需要在 pom.xml 中增加对应消息中间件的绑定器依赖即可,比如增加 RabbitMQ 的绑定器依赖:

<dependency>
    <groupId>org.springframework.cloud</groupId>
    <artifactId>spring-cloud-stream-binder-rabbit</artifactId>
</dependency>

如果使用 Kafka, 则引入:

<dependency>
    <groupId>org.springframework.cloud</groupId>
    <artifactId>spring-cloud-stream-binder-kafka</artifactId>
</dependency>

多绑定器配置

当应用程序的 classpath 下存在多个绑定器时, Spring Cloud Stream 在为消息通道做绑定操作时,无法判断应该使用哪个具体的绑定器,所以需要为每个输入或输出通道指定具体的绑定器。在一个应用程序中使用多个绑定器时,往往其中一个绑定器会是主要使用的,而第二个可能是为了适应一些特殊要求(比如性能等原因)。可以先通过设置默认绑定器 来为大部分的通道设置绑定器。比如,使用 RabbitMQ 设置默认绑定器:

spring:
  cloud:
    stream:
      defaultBinder: rabbit

在设置了默认绑定器之后,再为其他一些少数的消息通道单独设置绑定器,比如:

spring:
  cloud:
    stream:
      bindings:
        input:
          binder: kafka

需要注意的是,上面设置参数时用来指定具体绑定器的值并不是消息中间件的名称,而是在每个绑定器实现的

META-INF/spring.binders 文件中定义的标识(一个绑定器实现的标识可以定义多个,以逗号分隔),所以上面配置的 rabbit 和 kafka 分别来自于各自的配置定义,它们的具体内容如下所示:

rabbit:\
org.springframework.cloud.stream.binder.rabbit.config.RabbitServiceAutoConfiguration;
kafka:\
org.springframework.cloud.stream.binder.kafka.config.KafkaBinderConfiguration

另外,当需要在一个应用程序中使用同一类型不同环境的绑定器时,也可以通过配置轻松实现通道绑定。比如,当需要连接两个不同的 RabbitMQ 实例的时候,可以参照如下配置(application.properties 格式):


spring.cloud.stream.bindings.input.binder=rabbit1
spring.cloud.stream.bindings.output.binder=rabbit2

spring.cloud.stream.binders.rabbitl.type=rabbit
spring.cloud.stream.binders.rabbitl.environment.spring.rabbitmq. hos t=l92.168.0.101
spring.cloud.stream.binders.rabbitl.environment.spring.rabbitmq.port=5672
spring.cloud.stream.binders.rabbitl.environment.spring.rabbitmq.username=springcloud
spring.cloud.stream.binders.rabbitl.environment.spring.rabbitmq.password=123456

spring.cloud.stream.binders.rabbit2.type=rabbit
spring.cloud.stream.binders.rabbit2.environment.spring.rabbitmq-host=l92.168.0,102
spring.cloud.stream.binders.rabbit2.environment.spring.rabbitmq.port=5672
spring.cloud.stream.binders.rabbit2.environment.spring.rabbitmq.username=springcloud
spring.cloud.stream.binders.rabbit2.environment.spring.rabbitmq.password=123456

从上面的配置中,可以看到对于输入输出通道指定的绑定器采用了显式别名的配置方式,其中 input 通道的绑定指定了名为 rabbit1 的配置,而 output 通道的绑定指定了名为 rabbit2 的配置。当采用显式配置方式时会自动禁用默认的绑定器配置,所以当定义了显式配置别名后,对于这些绑定器的配置需要通过 spring.cloud.stream.binders.< configurationName> 属性来进行设置。


对于绑定器的配置主要有下面 4 个参数。

  • Spring.cloud.stream.binders.<configurationName>.type 指定了绑定器的类型,可以是 rabbit 、 kafak 或者其他自定义绑定器的标识名,绑定器标识名的定义位于绑定器的 META-INF/spring.binders 文件中。
  • Spring.cloud.stream.binders.<configurationName>.environment 参数可以直接用来设置各绑定器的属性,默认为空。
  • spring.cloud.stream.binders.<configurationName>.inheritEnvironment 参数用来配置当前绑定器是否继承应用程序自身的环境配置,默认为 true 。
  • spring.cloud.stream.binders.<configurationName>.defaultCandidate 参数用来设置当前绑定器配置是否被视为默认绑定器的候选项,默认为 true 。 当需要让当前配置不影响默认配置时,可以将该属性设置为 false 。

RabbitMQ 与 Kafka 绑定器

在之前的章节中多次提到, Spring Cloud Stream 自身就提供了对 RabbitMQ 和 Kafka 的绑定器实现。由于 RabbitMQ 和 Kafka 自身的实现结构有所不同,理解绑定器实现与消息中间件自有概念之间的对应关系,对于正确使用绑定器和消息中间件会有非常大的帮助。下面就来分别说说 RabbitMQ 与 Kafka 的绑定器是如何使用消息中间件中不同概念来实现消息的生产与消费的。


  • RabbitMQ 绑定器 : 在 RabbitMQ 中,通过 Exchange 交换器来实现 Spring Cloud Stream 的主题概念,所以消息通道的输入输出目标映射了一个具体的 Exchange 交换器。而对于每个消费组,则会为对应的 Exchange 交换器绑定一个 Queue 队列进行消息收发。
  • Kafka 绑定器:由于 Kafka 自身就有 Topic 概念,所以 Spring Cloud Stream 的主题直接釆用了 Kafka 的 Topic 主题概念,每个消费组的通道目标都会直接连接 Kaflca 的主题进行消息收发。

配置详解

在 Spring Cloud Stream 中对绑定通道和绑定器提供了通用的属性配置项,一些绑定器还允许使用附加属性来对消息中间件的一些独有特性进行配置。这些属性的配置可以通过 Spring Boot 支持的任何配置方式来进行,包括使用环境变量、 YAML 或者 properties 配置文件等。


基础配置

下表是 Spring Cloud Stream 应用级别的通用基础属性,这些属性都以 <code>spring.cloud.stream.为前缀。
参数名 说明 默认值
instanceCount 应用程序部署的实例数量,当使用 Kafka 的时候需要设置分区 1
instanceIndex 应用程序实例的索引,该值从 0 开始,最大值设置为-1。当使用分区和 Kafka 的时候使用
dynamicDestinations 动态绑定的目标列表,该列表默认为空,当设置了具体列表之后,只有列表中的目标才能被发现 空
defaultBinder 默认绑定器配置,在应用程序中有多个绑定器时使用 空
overrideCloudConnectors 该属性只适用于激活 cloud 配置并且提供了 Spring Cloud Connectors 的应用。当使用默认属性 false 时,绑定器会自动检测合适的服务来绑定(比如,在 Cloud Foundry 中绑定的 RabbitMQ 服务)。当设置为 true 时,绑定器将忽略绑定的服务,而是依赖应用程序中的设置属性来进行绑定和连接

绑定通道配置

对于绑定通道的属性配置,在之前的示例中已经有过一些介绍,这些配置可以在属性文件中通过 spring.cloud.stream.bindings.<channelName>.<property>=<value> 格式的参数来进行设置。其中 <channelName> 代表在绑定接口中定义的通道名称,比如, Sink 中的 input 、Source 中的 output。

由于绑定通道分为输入通道和输出通道,所以在绑定通道的配置中包含了三类面向不同通道类型的配置:通用配置、消费者配置、生产者配置。在下面介绍各具体配置属性时将省略 spring.cloud.stream.bindings.<channelName>. 前缀,但在实际使用的时候记得使用完整的参数名称进行配置。


通用配置
对于绑定通道的通用配置,它们既适用于输入通道,也适用于输出通道,它们通过 spring.cloud.stream.bindings.<channelName>. 前缀来进行设置,具体可配置的属性如下表所示。
参数名 说明 默认值
destination 该参数用来配置消息通道绑定在消息中间件中的目标名称,比如 RabbitMQ 的 Exchange 或 Kafka 的 Topic。如果配置的绑定通道是一个消费者(输入通道),那么它可以绑定多个目标,这些目标名称通过逗号分隔。如果没有设置该属性,将使用通道名
group 该参数用来设置绑定通道的消费组,该参数主要作用于输入通道,以保证同一消息组中的消息只会有一个消费实例接收和处理 null
contentType 该参数用来设置绑定通道的消息类型 null
binder 当存在多个绑定器时使用该参数来指定当前通道使用哪个具体的绑定器 null

消费者配置
下面这些配置仅对输入通道的绑定有效,它们以 spring.cloud.stream.bindings.<channelName>.consumer. 格式作为前缀。
参数名 说明 默认值
concurrency 输入通道消费者的并发数 1
partitioned 来自消息生产者的数据是否采用了分区 false
headerMode 当设置为 raw 的时候将禁用对消息头的解析,该属性只有在使用不支持消息头功能的中间件时有效,因为 Spring Cloud Stream 默认会解析嵌入的头部信息 embeddedheaders
maxAttempts 对输入通道消息处理的最大重试次数 3
backOffInitialInterval 重试消息处理的初始间隔时间 1000
backOffMaxInterval 重试消息处理的最大间隔时间 10000
backOffMultiplier 重试消息处理时间间隔的递增乘数 2.0

生产者配置
参数名 说明 默认值
partitionKeyExpression 该参数用来配置输出通道数据分区键的 SpEl 表达式,当设置该属性之后,将对当前绑定通道的输出数据进行分区处理。同时,partitionCount 的参数必须大于 1 才能生效。该参数与 partitionKeyExtractorClass 参数互斥,不能同时设置 null
partitionKeyExtractorClass 该参数用来配置分区键提取策略接口 PartitionKeyExtractorStrategy 的实现。当设置该属性之后,将对当前绑定通道的输出数据进行分区处理,同时,partitionCount 的参数必须大于 1 才能生效。该参数与 partitionKeyExpression 参数互斥,不能同时设置 null

参数名 说明 默认值
partitionSelectorClass 该参数用来指定分区选择器接口 PartitionSelectorStrategy 的实现。它与 partitionSelectorExpression 参数互斥,不能同时设置。如果两者都不设置,那么分区选择计算规则为 hashCode(key)%partitionCount,这里的 key 根据 partitionKeyExpression 或 partitionKeyExtractorClass 的配置计算得到 null
partitionSelectorExpression 该参数用来配置自定义分区选择器的 SpEL 表达式。它与 partitionSelectorClass 参数互斥,不能同时设置。如果两者都不设置,那么分区选择计算规则为 hashCode(key)%partitionCount,这里的 key 根据 partitionKeyExpression 或 partitionKeyExtractorClass 的配置计算得到 null
partitionCount 当分区功能开启时,使用该参数来配置消息数据的分区数。如果消息生产者已经配置了分区键的生成策略,那么它的值必须大于 1 1
headerMode 当设置为 raw 的时候将禁用对消息头的解析,该属性只有在使用不支持消息头功能的中间件时有效,因为 Spring Cloud Stream 默认会解析嵌入的头部信息 embeddedheaders

绑定器的配置属性

由于 Spring Cloud Stream 目前只实现了对 RabbitMQ 和 Kafka 的绑定器,所以对于绑定器的属性配置主要是针对这两个中间件。同时因为这两个中间件自身结构的不同,所以会有不同的附加属性来配置各自的一些特殊功能和基础设施。


RabbitMQ 配置

RabbitMQ 绑定器的配置同绑定通道的配置一样,分为三种不同类型:通用配置、消费者配置以及生产者配置。


通用配置

由于 RabbitMQ 绑定器默认使用了 Spring Boot 的 ConnectionFactory, 所以 RabbitMQ 绑定器支持在 Spring Boot 中的配置选项,它们以 spring.rabbitmq. 为前缀。 在之前的示例中,使用了这些配置来指定具体的 RabbitMQ 地址、端口、用户信息等,更多配置可见 Spring Boot 文档中对 RabbitMQ 支持的章节内容,或是通过 spring-boot-starter-amqp 模块中 RabbitProperties 类的源码来查看。


在 Spring Cloud Stream 对 RabbitMQ 实现的绑定器中主要有下面几个属性,它们都以 spring.cloud.stream.rabbit.binder. 为前缀。这些属性可以在 org.springframework.cloud.stream.binder.rabbit.config.RabbitBinderConfigurationProperties 中找到它们。


参数名 说明 默认值
adminAddresses 该参数用来配置 RabbitMQ 管理插件的 URL,当需要配置多个时用逗号分隔。该参数只有在 nodes 参数包含多个时使用,并且这里配置的内容必须在 spring.rabbitmq.addresses 中存在
nodes 该参数用来配置 RabbitMQ 的节点名称,当需要配置多个时用逗号分隔。在配置多个的情况下,可以用来定位队列所在的服务器的地址。这里配置的内容必须在 spring.rabbitmq.addresses 中存在
compressionLevel 绑定通道的压缩级别,它的具体可选值及含义可见 java.util.zip.Deflater 中的定义 1

消费者配置

下面这些配置仅对 RabbitMQ 输入通道的绑定有效,它们以 spring.cloud.stream.rabbit.bindings.<channelName>.consumer. 格式作为前缀。


参数名 说明 默认值
acknowledgeMode 用来设置消息的确认模式,可选配置包含:NONE,MANUAL,AUTO AUTO
autoBindDlq 用来设置是否自动声明 DLQ(Dead-Letter-Queue),并绑定到 DLX(Dead-Letter-Exchange)上 false
durableSubscription 用来设置订阅是否被持久化,该参数仅在 group 被设置的时候有效 true
maxConcurrency 用来设置消费者的最大并发数 1
prefetch 用来设置预取数量,它表示在一次会话中从消息中间件获取的消息数量,该值越大消息处理越快,但是会导致非顺序处理的风险 1
prefix 用来设置统一的目标和队列名称前缀
recoveryInterval 用来设置恢复连接的尝试时间间隔,以毫秒为单位 5000

参数名 说明 默认值
requestRejected 用来设置消息传递失败时候重传 true
requestHeaderPatterns 用来设置需要被传递的请求头信息 [STANDARD_REQUEST_HEADERS,'*']
replyHeaderPatterns 用来设置需要被传递的响应头信息 [STANDARD_REPLY_HEADERS,'*']
republishToDlq 默认情况下,消息在重试也失败之后会被拒绝。如果 DLQ 被配置的时候,RabbitMQ 会将失败的消息路由到 DLQ 中。如果该参数设置为 true,总线会将失败的消息附加一些头信息(包括异常消息,引起失败的跟踪堆栈)之后重新发布到 DLQ 中
transacted 用来设置是否启用 channel-transacted,即是否在消息中使用事务 false
txSize 用来设置 transaction-size 的数量,当 acknowledgeMode 被设置为 AUTO 时,容器会在处理 txSize 个消息之后才开始英达 1

生产者配置
下面这些配置仅对 RabbitMQ 输出通道的绑定有效,它们以 spring.cloud.stream.rabbit.bindings.<channelName>.producer. 格式作为前缀。
参数名 说明 默认值
autoBindQlp 用来设置是否自动声明 DLQ(Dead-Letter-Queue),并绑定到 DLX(Dead-Letter-Exchange)上 false
batchingEnable 是否启用消息批处理 false
batchSize 当批处理开启时,用来设置缓存的批处理消息数量 100
batchBufferLimit 批处理缓存限制 10000
batchTimeout 批处理超时时间 5000
compress 消息发送时是否启用压缩 false
deliveryMode 消息发送模式 PERSISTENT
prefix 用来设置统一的目标前缀
requestHeaderPatterns 用来设置需要被传递的请求头信息 [STANDARD_REQUEST_HEADERS,'*']
replyHeaderPatterns 用来设置需要被传递的响应头信息 [STANDARD_REPLY_HEADERS,'*']

Kafka 配置

Kafka 绑定器的配置在类别上与 RabbitMQ —样,分为三种不同类型:通用配置、消费者配置以及生产者配置。但是由于 RabbitMQ 与 Kafka 自身有一些差异,所以它们的配置也不一样。


通用配置
Spring Cloud Stream 实现的 Kafka 绑定器包含下面这些通用配置,它们都以 spring.cloud.stream.kafka.binder. 为前缀。
参数名 说明 默认值
brokers Kafka 绑定器链接的消息中间件列表。需要配置多个时用逗号分隔,每个地址可以是单独的 host,也可以是 host:port 的形式 localhost
defaultBrokerPort 用来设置默认的消息中间件端口号。当 brokers 中的配置地址没有包含端口信息时,将使用该参数配置的默认端口进行连接 9092
zkNodes kafka 绑定器使用的 Zookeeper 节点列表。需要配置多个时用逗号分隔,每个地址可以是单独的 host,也可以是 host:port 的形式 localhost
defaultZkPort 用来设置默认的 Zookeeper 端口号。当 zkNodes 中的配置地址没有包含端口信息时,将使用该参数配置的默认端口进行连接 2181
headers 用来设置会被传输的自定义头信息
offsetUpdateTimeWindow 用来设置 offset 的更新频率,以毫秒为单位,如果设置为 0 则忽略 10000
offsetUpdateCount 用来设置 offset 以次数表示的更新频率,如果为 0 则忽略,该参数与 offsetUpdateTimeWindow 互斥 0

参数名 说明 默认值
minPartitionCount 该参数仅在设置了 autoCreateTopics 和 autoAdd-Partitions 时生效,用来设置该绑定器所使用主题的全局分区最小数量。如果当生产者的 partitionCount 参数或 instanceCount*concurrency 的值大于该参数配置时,该参数值将被覆盖 1
requiredAcks 用来设置确认消息的数量 1
replicationFactor 当 autoCreateTopics 参数为 true 时候,用来配置自动创建主题的副本数量 1
autoCreateTopics 该参数默认为 true,绑定器会自动地创建新主题。如果设置为 false,那么绑定器使用已经配置的主题,但是在这种情况下,如果需要使用的主题不存在,绑定器会启动失败 true
autoAddPartitions 该参数默认为 false,绑定器会根据已经配置的主题分区来实现,如果目标主题的分区数小于预期值,那么绑定器会启动失败。如果该参数设置为 true,绑定器将在需要的时候自动创建新的分区 false
socketBufferSize 该参数用来设置 Kafka 的 Socket 缓存大小 2097152

消费者配置
下面这些配置仅对 Kafka 输入通道的绑定有效,它们以 spring.Cloud.Stream.kafka.bindings.<channelName>.consumer. 格式作为前缀。
参数名 说明 默认值
bufferSize kafka 批量发送前的缓存数据上限,以字节为单位 16384
sync 该参数用来设置 Kafka 消息生产者的发送模式,默认为 false,即采用 async 配置,允许批量发送数据。当设置为 true 时,将采用 sync 配置,消息将不会被批量发送,而是一条一条地发送 false
batchTimeout 消息生产者批量发送时,为了积累更多发送数据而设置的等待时间。通常情况下,生产者基本不会等待,而是直接发送所有在前一批次发送时积累的消息数据。当我们设置一个非 0 值时,可以以延迟为代价来增加系统的吞吐量 0

活动 9.1

Spring Cloud Stream 的使用


练习问题


  1. 说明以下语句是正确还是错误。
    通过定义绑定器作为中间层,完美地实现了应用程序与消息中间件细节之间的隔离 。
    a. 正确
    b. 错误

    答案

    a 正确


  1. Spring Cloud Stream 中的消息通信方式遵循了什么模式?
    a. 发布-订阅
    b. 发布
    c. 订阅
    d. 订阅-发布

    答案

    a 正确


  1. 在 Spring Cloud Stream 中,需要通什么注解来为应用启动消息驱动的功能?
    a. @EnaBinding
    b. @EnableBinding
    c. @Enableding
    d. @EnableBind

    答案

    b 正确


  1. 应用 Spring Cloud Stream 是基于什么构建起的
    a. Spring Integration
    b. Integration
    c. Spring
    d. Spring Integrat

    答案

    a 正确


小结


在本章中,您学习了:

  • 本地配置安装 RabbitMQ
  • Spring Cloud Stream 的核心概论
    • 绑定器
    • 发布-订阅模式
    • 消费组
  • 对于一些核心注解的使用

    • 开启绑定功能

    • 绑定消息通道

    • 消息生产与消费

    • 消费组与消息分区

    • 消息类型


  • 绑定器背后的细节和行为

    • 绑定器 SPI
    • 自动化配置
    • RabbitMQ 与 Kafka 绑定器
  • 了解 Spring Cloud Stream 中相应的配置

    • 基础配置
    • 绑定通道配置
    • 绑定器配置

谢谢

Views: 33

SpringCloud与微服务-第8章-Nacos分布式配置中心

在微服务架构中,当系统从一个单体应用,被拆分成分布式系统上一个个服务节点后,配置文件也必须跟着迁移(分割)。在系统架构中,配置中心是整个微服务基础架构体系中的一个组件,它的功能看上去并不起眼,无非就是配置的管理和存取,但它是整个微服务架构中不可或缺的一环。


应用程序在启动和运行的时候往往需要读取一些配置信息,配置基本上伴随着应用程序的整个生命周期,应用在启动时通过读取配置来初始化,在运行时根据配置调整行为。同一份程序在不同的环境(开发、测试、生产)、不同的集群(如不同的数据中心)经常需有不同的配置,所以需要有完善的环境、集群配置管理。


Spring Cloud Alibaba Nacos 的一大优势是整合了注册中心、配置中心功能,部署和操作更加直观简单,它简化了架构复杂度,并减轻运维及部署工作。在之前的章节中,我们已经使用 Nacos 作为注册中心,本章节我们将详细介绍 Nacos 配置中心的功能。


目标

在本章中,您将学习:

  • Nacos 配置管理
  • 配置拉取
  • 配置热更新
  • 多环境配置共享
  • Nacos 集群搭建

快速入门

当微服务部署的实例越来越多,达到数十、数百时,逐个修改微服务配置是一件效率非常低的事情,而且很容易出错。我们需要一种统一配置管理方案,可以集中管理所有实例的配置。下面我们就通过一个小案例来感受一下 Spring Cloud Alibaba Nacos 配置中心的强大功能吧。


在 Nacos 中创建配置

我们以 userservice 服务为例,在 Nacos 配置中心创建一份配置文件。进入 Nacos 的管理界面,然后选择配置管理菜单,点击配置列表选项,点击右侧"+"符号,如下图所示:


bg fit


点击之后,填写弹出的表单页,最后点击发布,我们就创建好一份 userservice 服务的配置,如下图所示:


image-20220731211155187


发布成功之后,我们在控制台的配置列表中可以看到新增配置,并且可以进行再次编辑,查看详情,删除等操作。


image-20220731211354245


此时,在 Nacos 配置中心的工作暂时告一段落,我们已经创建好一份 userservice 服务的配置文件,那么怎么样才能让 userservice 服务读取到这份配置呢?


读取 Nacos 中配置

userservice 服务要读取 Nacos 中管理的配置,并且与本地的 application.yml 配置合并,才能完成项目启动。但是此时有一个问题必须解决: 尚未读取 application.yml,又如何得知 nacos 地址呢?


为解决此问题,Spring 引入了一种新的配置文件:bootstrap.yaml 文件,该配置文件的优先级很高,会在 application.yml 之前被读取,流程如下:


image-20220731213028366


首先,我们要给 userservice 引入 Nacos 配置中心的 Maven 依赖:

<!--nacos配置管理依赖-->
<dependency>
    <groupId>com.alibaba.cloud</groupId>
    <artifactId>spring-cloud-starter-alibaba-nacos-config</artifactId>
</dependency>

其次,在 resources 目录下创建 bootstrap.yml 文件,并进行配置:

spring:
  application:
    name: userservice # 服务名称
  profiles:
    active: dev #开发环境,这里是dev
  cloud:
    nacos:
      server-addr: localhost:8848 # Nacos地址
      config:
        file-extension: yaml # 文件后缀名

此时,同学们会发现,在 bootstrap.yml 中配置的这些信息,正好可以满足我们的需求: 去 nacos 中查找一个名字为 userservice-dev.yaml 文件。这些信息正是${spring.application.name}-${spring.profiles.active}.${spring.cloud.nacos.config.file-extension}的值拼接而成。


image-20220731213648989


这也是我们在 nacos 中创建配置文件时的命名规则的含义所在。


测试验证

我们在 userservice 中新增一个控制器类,进行测试:

@Slf4j
@RestController
@RequestMapping("/config")
public class NacosConfigTestController {
    @Value("${myconfig.msg}")
    private String msg;

    @RequestMapping("/msg")
    public String getMsg(){
        return this.msg;
    }
}

修改完成之后,我们重启 userservice 服务,访问: http://localhost:8081/config/msg ,结果如图所示:

image-20220731214102597

我们通过简单的几个步骤,成功的使用 Nacos 配置中心的功能,实现了在 Nacos 控制台对配置文件的简单管理。


现在如果尝试修改 Nacos 配置中心的config.msg的值,然后再次访问 http://localhost:8081/config/msg ,会发现返回的结果并没有发生变化,
要想使新的配置生效仍需重启服务.


配置热更新

如果每次线上修改修改了配置文件,都需要重启服务才能生效,这样的话,就会造成服务的中断,这是我们不希望看到的。那么,有没有一种方式,可以实现在线上修改配置文件,而不需要重启服务呢?

答案是可以的, 我们可以通过为 Nacos 的配置热更新来实现此类需求。

所谓的热更新是指,在 Nacos 中的配置文件变更后,对应的微服务无需重启就可以感知到最新的配置。


所谓的 Nacos 配置热更新,就是指 Nacos 中的线上配置变更后, 无需重启服务就可以感知到最新的配置变化。

在微服务中,可以通过 SpringBoot 的方式进行配置的注入 :

  • 通过 @Value 注解进行注入单个配置
  • 通过 @AutoWired 注解进行注入包含多个配置属性的对象
  • 在需要注入配置的类上添加 @RefreshScope 注解,使得注入的配置可以和线上配置保持同步

使用@Value 注解进行注入

在@Value 注入的变量所在类上添加注解@RefreshScope,:

@Slf4j
@RestController
@RequestMapping("/config")
@RefreshScope // 添加此注解, 使得配置文件可以热更新
public class NacosConfigTestController {
    @Value("${myconfig.msg}") // 通过@Value 注解进行注入响应的配置
    private String msg;

    @RequestMapping("/msg")
    public String getMsg(){
        return this.msg;
    }
}

通过查看@RefreshScope 的注解发现,其底层是重新创建了实例,并进行了依赖注入。

/**
 * 添加了@RefreshScope的Bean 可以在运行时刷新
 * 任何使用它们的组件都将在下一次方法调用时获得一个新实例,完全初始化并注入所有依赖项。
 * @author Dave Syer
 */
@Target({ ElementType.TYPE, ElementType.METHOD })
@Retention(RetentionPolicy.RUNTIME)
@Scope("refresh")
@Documented
public @interface RefreshScope {
   ScopedProxyMode proxyMode() default ScopedProxyMode.TARGET_CLASS;
}

使用@Autowire 注解进行注入

用@Value 可以进行简单值的属性注入,如果我们想要在 Bean 中注入一个对象,往往使用@Autowire 注解进行注入。针对这种注入方式,我们可以使用@ConfigurationProperties 注解进行配置读取。


我们仍然以myconfig.msg配置为例,使用@ConfigurationProperties 进行配置读取。​ 首先,编写一个 Java Bean:

@Data
@Component
@ConfigurationProperties(prefix = "myconfig")
public class UserNacosConfig {
    private String msg;
    private String chapter;
}

在编写 JavaBean 时,@ConfigurationProperties 注解的 prefix 属性和 Nacos 中配置的属性前缀保持一致,成员变量的名字和 Nacos 中配置的属性的 Key 保持一致。 Nacos 中配置我们修改如下:

w:35em

修改完 Nacos 中配置之后,一定要记得点击发布按钮进行发布。


最后修改 NacosConfigTestController,注入我们编写的 JavaBean:

@Slf4j
@RestController
@RequestMapping("/config")
@RefreshScope
public class NacosConfigTestController {
    @Value("${myconfig.msg}")
    private String msg;

    @Autowired
    private UserNacosConfig userNacosConfig;

    @RequestMapping("/msg")
    public String getMsg() {
        return this.msg;
    }

    @RequestMapping("/message")
    public String getMessage() {
        return userNacosConfig.getMsg() + userNacosConfig.getChapter();
    }
}

重启 userservice 应用,访问:http://localhost:8081/config/message ,结果如图所示:

image-20220801112128380

此时是因为我们重启了服务,所以能够读取到最新的配置,这是可以预知的结果。


现在,我们继续修改配置,而不重启服务,观察是否能够得到最新的配置内容。

最新配置更改如下:

w:35em


点击发布之后,我们不用重启 userservice 服务,直接再次访问:http://localhost:8081/config/message , 结果如图所示:

w:35em

可以看到,我们已经可以实现配置的热更新。


加载指定的 Nacos 配置

配置默认的 GroupID 为 DEFAULT_GROUP, 也可以创建自定义的 GroupID, 但是在加载配置时,需要指定 GroupID。


在 Nacos 配置中心, 创建配置如下:

w:24em


pom.xml 加入依赖

  <dependency>
    <groupId>com.alibaba.boot</groupId>
    <artifactId>nacos-config-spring-boot-starter</artifactId>
    <version>0.2.12</version>
  </dependency>

启动类加入注解

// 读取 Nacos 配置中心
@NacosPropertySource(dataId = "retryTimes", groupId = "RiskModule",autoRefreshed = true)
@MapperScan("com.niit.user.mapper")
@SpringBootApplication
public class UserApplication {...}

使用时通过@NacosValue 注解进行注入

    @NacosValue(value = "${readRetryTimes}", autoRefreshed = true)
    private String readRetryTimes;
    @NacosValue(value = "${writeRetryTimes}", autoRefreshed = true)
    private String writeRetryTimes;

此外, 也可以通过 Nacos 的 API 进行配置的设置和读取

发布配置
curl -X POST "http://127.0.0.1:8848/nacos/v1/cs/configs?dataId=nacos.cfg.dataId&group=test&content=HelloWorld"

获取配置
curl -X GET "http://127.0.0.1:8848/nacos/v1/cs/configs?dataId=nacos.cfg.dataId&group=test"


注意:

  1. 不是所有的配置都适合放到配置中心,维护起来比较麻烦。
  2. 建议将一些关键参数,需要在运行时进行调整的参数放到 Nacos 配置中心,这些参数一般都是自定义配置信息。

多环境配置

微服务系统中, 及存在共通的配置比如应用的名称。而有些配置在不同的环境下,比如开发环境和生产环境下的微服务的配置是不尽相同的, 比如端口号, 数据库的链接,Nacos 的地址。

此时,我们想尽量减少重复性工作,提升开发和维护的效率,那就需要将这些跨环境的那些相同配置提取出来。我们将共同配置提取到一个地方,就可以做到只用设置一次, 所有环境(比如开发.测试以及生产环境)都会读取该配置,这样就达到了配置共享的目的。


配置读取

在微服务启动的时候,会去 Nacos 中读取多个配置:

  • bootstrap.yml 启动配置:
    会优先加载,用于加载一些系统级别的配置,比如连接到配置中心的配置。
  • ${spring.application.name}.yaml
    例如 userservice.yaml。无论微服务的环境是开发环境也好, 是生产环境也好,${spring.application.name}.yaml这个配置文件都会userservice加载。

在 Nacos 中添加共享配置

我们在 nacos 中添加一个 userservice.yaml 文件:

w:22em

注意: 命名规则是 ${spring.application.name}.yaml ,该共享配置与环境无关。


读取新增配置

在 userservice 服务中,修改配置类,读取新增的author属性:

@Data
@Component
@ConfigurationProperties(prefix = "myconfig")
public class UserNacosConfig {
    private String chapter;
    private String msg;
    private String author;
}

新增一个控制器方法,以便于进行测试验证:

@Slf4j
@RestController
@RequestMapping("/config")
@RefreshScope
public class NacosConfigTestController {

    @Autowired
    private UserNacosConfig userNacosConfig;
    ...
    @RequestMapping("/author")
    public String getAuthor() {
        return userNacosConfig.getAuthor();
    }
}

多环境测试

使用 dev 环境启动一个 userservice 实例,端口号 8081,再用 test 环境启动另一个 userservice 实例,端口号 8082,然后我们分别访问: http://localhost:8081/config/author 和 http://localhost:8082/config/author:


结果如图所示:

w:28em


我们知道,在 bootstrap.yml 文件中指定了应用读取 Nacos 中的 userservice-dev.yaml 文件,并且 Nacos 中并没有 userservice-test.yaml 配置文件,但是我们访问 8082 端口的 userservice 时,仍然能够读取myconfig.author的值,这说明userservice.yaml文件确实是多环境共享的一个配置文件。


配置优先级

我们现在考虑一个问题,如果多个配置文件中都对某个属性进行了定义,那么微服务究竟应该读取那个配置来使用呢?我们不妨来做一个小小的测验。

在每一个配置文件中都定义一个相同的属性,比如myconfig.bookName属性,如下所示:

  • 在 userservice-dev.yaml 中,bookName 的值是 Spring Cloud Dev。
  • 在 userservice.yaml 中,bookName 的值是 Spring Cloud Default。
  • 在本地 application-dev.yaml 中,bookName 的值是 Spring Cloud Local。

修改配置类,读取 bookName 属性:

@Data
@Component
@ConfigurationProperties(prefix = "myconfig")
public class UserNacosConfig {
    private String chapter;
    private String msg;
    private String author;
    private String bookName;
}

新增控制器方法,读取myconfig.bookName属性的值:

@Slf4j
@RestController
@RequestMapping("/config")
@RefreshScope
public class NacosConfigTestController {
    @Autowired
    private UserNacosConfig userNacosConfig;
    ...
    @RequestMapping("/bookName")
    public String getBookName() {
        return userNacosConfig.getBookName();
    }
}

在 dev 环境下,重启 userservice,端口号 8081,现在我们访问: http://localhost:8081/config/bookName ,结果如图所示:

w:32em


这说明,Nacos 线上配置的带有环境变量的配置文件的优先级是最高的。那么线上的共享配置和本地配置,哪个优先级更高呢?我们暂时删除 Nacos 线上配置的带有环境变量的配置文件中的属性,重新访问: http://localhost:8081/config/bookName ,结果如图所示:

image-20220801174120879


通过这个小小的测验,对于配置优先级问题,我们总结出如下规律:

  • 线上配置优先于本地配置
  • 线上自定义环境配置优先于线上共享配置

bg right fit


这样的优先级规律,也恰恰符合配置中心的设计原则,即线上可配置优先于本地预配置,自定义个性化配置优先于多环境共享配置。详见 Nacos-config 参考文档


搭建 Nacos 高可用集群


在之前的内容中,我们已经学习了 Nacos 的基本用法。不过我们一直是使用的 nacos 单节点服务,Nacos 单节点模式只适用于线下环境,在企业的生产环境中,我们是不能使用 Nacos 单节点进行业务架构部署的。


单节点对于高可用设计来说是远远不够的,因为单节点一旦出先故障,那么所有依赖于该节点的其他微服务,都会出现问题,这样会导致整个业务系统的大崩溃。所以,我们接下来就开始学习如何搭建 Nacos 集群,构建一个高可用的服务治理体系。


Nacos 集群

一个 Nacos 集群,至少要有三个节点。下方是 Spring Cloud Alibaba Nacos 官方给出的最简易的 Nacos 集群架构图:


image-20220802110040457

其中包含 3 个 nacos 节点,然后一个负载均衡器代理 3 个 Nacos。这里负载均衡器可以使用 nginx。接下来,我们以 windows 系统为例,搭建一个最简单的 Nacos 集群进行学习。


注意:

在实际生产环境中,需要给做反向代理的 nginx 服务器设置一个域名,这样后续如果有服务器迁移,nacos 的客户端也无需更改配置。Nacos 的各个节点应该部署到多个不同服务器,做好容灾和隔离。Mysql 应该搭建一个主从高可用的集群。 本案例以单机 windows 系统为例,模拟三个 Nacos 节点,并且使用单机的 Mysql8.0 以上的版本进行集群搭建的演示。


Nacos 集群搭建步骤如下:

  1. 搭建数据库,初始化数据库表结构
  2. 下载 nacos 安装包
  3. 配置 nacos 节点
  4. 启动 nacos 集群
  5. nginx 反向代理
  6. 修改微服务 Nacos 地址

初始化数据库

Nacos 默认数据存储在内嵌数据库 Derby 中,不属于生产可用的数据库。官方推荐的最佳实践是使用带有主从的高可用数据库集群。这里我们以单点的数据库为例来演示。首先新建一个数据库,命名为 nacos,而后导入 Nacos 的 MySQL 数据库的初始化 SQL 文件.


使用 Nacos 项目官方提供的初始化 sql 文件进行初始化操作, 具体位置在 conf 目录下的mysql-schema.sql文件


h:13em

Nacos 的 MySQL 初始化脚本在线地址
https://github.com/alibaba/nacos/tree/develop/config/src/main/resources/META-INF


在本地 Mysql 中执行成功后,nacos 数据库中将会出现下面的表:

image-20220802111857744


下载 Nacos 安装包

本知识点在本书第三章节已经介绍过,不再赘述,我们以 nacos2.1.0 版本为例,进行搭建。

image-20220802112128986


配置 Nacos 节点

我们在本地 windows 系统上配置三个 Nacos 节点,各个节点的 IP 和端口如下表所示:

节点 IP port
nacos1 127.0.0.1 8858
nacos2 127.0.0.1 8868
nacos3 127.0.0.1 8878

注意: 本地 windows 系统上,如果端口被占用,可以修改端口号,但是要保证三个节点的端口号不一样, 且不要相邻(容易冲突)。


我们先配置好一个 nacos 节点,然后再复制成三个节点。

下面以一个节点操作为例:


1.解压缩

将下载的安装包,解压缩至一个没有中文目录的文件夹,命名为 nacos1,解压完成之后,Nacos 目录如图所示:

image-20220802114728043


-目录说明:

  • bin:启动脚本
  • conf:配置文件

2.修改配置

进入 nacos 的 conf 目录,修改配置文件 cluster.conf.example,重命名为 cluster.conf:

image-20220802114924887


编辑此文件,添加 nacos 集群配置:

192.168.65.1:8858
192.168.65.1:8868
192.168.65.1:8878

192.168.65.1 为本机真实的 IP 地址
真实 ip 可以在启动单机版 nacos 后从日志提示中查看


然后修改 application.properties 文件,进行数据库配置:

image-20220802115229225


属性说明:

  • spring.datasource.platform=mysql: 指定数据库平台为 mysql,Nacos 目前只支持 Mysql。
  • db.num=1 : Mysql 的实例个数,本案例中为 1 个 Mysql 实例。
  • db.url.0: 数据库链接信息,本案例中我们使用的是 Mysql8.0 以上的版本,并且创建了 nacos 数据库,所以链接如上所示。
  • db.user.0: 数据库用户名
  • db.password.0: 数据库密码

接着修改应用的端口号,第一个节点是: 8858,如下图所示:

server.port=8858

3.复制节点

修改完上述配置并保存后,我们将该节点复制两份,命名为 nacos2 和 nacos3,并将 nacos1的端口修改为8858, nacos2 的端口修改为 8868,将 nacos3 的端口修改为 8878。完成之后,nacos 三个节点的集群整体如下所示:


启动 Nacos 集群

在每个 Nacos 节点的 bin 目录下,使用 cmd 命令窗口执行: startup.cmd。 Nacos 的该命令,默认就是以集群方式进行启动。启动成功,最后会打印出日志信息:

image-20220802150700225


Nginx 反向代理

经过以上步骤,我们已经成功搭建好具有三个节点的 Nacos 集群,现在我们使用 Nginx 给该集群配置负载均衡和反向代理。


1.下载 Nginx

下载 Nginx,下载地址:https://nginx.org/en/download.html 我们选择最新的稳定版本。

h:12em


2.解压缩

将安装包解压缩到任意的非中文目录下:

image-20220802153243148


3.添加代理配置

修改 conf/nginx.conf 文件,将下面的配置放在配置文件的 http 节点之内:

upstream nacos-cluster {
  server 127.0.0.1:8858;
    server 127.0.0.1:8868;
    server 127.0.0.1:8878;
}

server {
    listen       8848;
    server_name  localhost;

    location /nacos {
        proxy_pass http://nacos-cluster;
    }
}

该代理配置监听 8848 端口,虚拟路径为nacos。


4.启动 Nginx

修改完配置并保存之后,回到 Nginx 的安装目录,双击nginx.exe即可启动。此时,我们可以访问: http://localhost/nacos

bg right fit


通过固定访问方式 http://localhost:8848/nacos 访问集群管理 - 节点列表, 如果看到集群中的三个节点都已经启动成功, 且为 UP 状态, 则说明 Nginx 反向代理配置成功。


微服务配置 Nacos 地址

由于我们现在已经搭建了 Nacos 集群,并且使用 Nginx 做了负载均衡和反向代理,所以各个微服务之前配置的 Nacos 地址肯定无法再使用,微服务中的 Nacos 地址配置要与 Nginx 中保持一致。

以 userservice 服务为例,修改其 bootstrap.yaml 文件即可:


spring:
  application:
    name: userservice
  profiles:
    active: dev #环境
  cloud:
    nacos:
      server-addr: localhost:8848 #nacos集群的NGINX反向代理地址
      config:
        file-extension: yaml #配置文件后缀名

现在我们启动三个 userservice 服务的实例进行测试,启动成功之后,我们可以在 Nacos 控制台看到如下信息:
image-20240320190015616

当我们搭建完成 Nacos 集群之后,在使用的时候几乎完全和我们之前所学的单节点 Nacos 的使用方法一致。


活动 8.1:

Nacos 配置中心使用


练习问题


  1. 当使用@Value 注解进行属性注入时,在需要的类上加上什么注解可以实现动态刷新配置。 ()
    a. @Service
    b. @Component
    c. @RestController
    d. @RefreshScope

    查看答案

    正确答案: d


  1. Nacos 配置中心的 artifactId 是下列的哪一个? ()
    a. spring-cloud-starter-alibaba-nacos-config
    b. spring-cloud-starter-alibaba-nacos-discovery
    c. spring-cloud-starter-Netflix-nacos-config
    d. spring-cloud-starter-alibaba-config-config

    查看答案

    正确答案: a


  1. 下列配置文件,哪个是 userservice 服务的多环境共享配置? ()
    a.userservice-dev.yaml
    b.application-dev.yaml
    c.userservice.yaml
    d.userservice-common.yaml

    查看答案

    正确答案: c


  1. 搭建 Nacos 集群,最少需要几个 Nacos 实例节点?
    a. 1
    b. 2
    c. 3
    d. 4

    查看答案

    正确答案: c


小结

在本章中,您学习了:

  • Nacos 配置中心快速入门

  • 配置热更新

    • 使用@Value 注解进行注入
    • 使用@Autowire 注解进行注入
  • 多环境配置共享

    • 配置优先级
  • 搭建 Nacos 集群

    • 初始化数据库
    • 配置 Nacos 节点
    • 启动 Nacos 集群
    • Nginx 反向代理

Views: 23

SpringSecurity-OAuth2+JWT+SpringCloudGateway实现统一鉴权管理

SpringSecurityJava学习笔记

版权 本文为时间海绵原创文章,转载无需和我联系,但请注明来自博客 https://blog.hzchendou.com

一、SpringSecurity 入门

介绍

SpringSecurity 是Spring 全家桶中的安全框架,为了解决“用户身份认证”、“资源访问鉴权”这两个核心问题,SpringSecurity提供了一整套安全框架,基于安全框架,用户可以自定义身份认证、资源鉴权功能,例如:手机验证码登录、基于RDBC鉴权等,本文章主要介绍如何创建基于SpringSecurity项目。

项目创建

项目源码已上传到Gitee: 地址。

项目依赖

基于 SpringBoot 创建SpringSecurity 可以实现开箱即用功能,引入依赖项:

- SpringBoot依赖

 <parent>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-parent</artifactId>
    <version>2.7.0</version>
    <relativePath/> <!-- lookup parent from repository -->
  </parent>

- Spring MVC 依赖(搭建基于 http 协议的web项目)

<dependency>
   <groupId>org.springframework.boot</groupId>
   <artifactId>spring-boot-starter-web</artifactId>
</dependency>

- Spring Security 依赖

<dependency>
   <groupId>org.springframework.boot</groupId>
   <artifactId>spring-boot-starter-security</artifactId>
</dependency>

详细 pom 文件可以参见源码:https://gitee.com/hzchendou/spring-security-demo/blob/lesson1/pom.xml

项目模块

创建简单mvc API,代码如下:

/**
 * hello 访问控制器
 * @Date: 2022-05-23 11:27
 * @since: 1.0
 */
@RequestMapping("/anonymity")
@RestController
public class AnonymityController {

    @RequestMapping("/hello")
    public ResultVO test() {
        return ResultVO.success("hello world");
    }
}

项目启动

自此完成项目配置,基于SpringBoot 自动装配功能可以帮助我们完成大部分配置,引入依赖后会帮助创建一个基础运行框架,配置了一些默认配置项,运行项目后看到如下日志:

  .   ____          _            __ _ _
 /\\ / ___'_ __ _ _(_)_ __  __ _ \ \ \ \
( ( )\___ | '_ | '_| | '_ \/ _` | \ \ \ \
 \\/  ___)| |_)| | | | | || (_| |  ) ) ) )
  '  |____| .__|_| |_|_| |_\__, | / / / /
 =========|_|==============|___/=/_/_/_/
 :: Spring Boot ::                (v2.7.0)

2022-05-23 12:23:13.584  INFO 8538 --- [           main] c.h.b.demo.springsecurity.Application    : Starting Application using Java 1.8.0_211 on hzchendoudeMac-mini.local with PID 8538 (/Users/chendou/repo/hzchendou/learning/springsecurity/target/classes started by chendou in /Users/chendou/repo/hzchendou/learning/springsecurity)
2022-05-23 12:23:13.586  INFO 8538 --- [           main] c.h.b.demo.springsecurity.Application    : No active profile set, falling back to 1 default profile: "default"
2022-05-23 12:23:14.338  INFO 8538 --- [           main] o.s.b.w.embedded.tomcat.TomcatWebServer  : Tomcat initialized with port(s): 8080 (http)
2022-05-23 12:23:14.344  INFO 8538 --- [           main] o.apache.catalina.core.StandardService   : Starting service [Tomcat]
2022-05-23 12:23:14.344  INFO 8538 --- [           main] org.apache.catalina.core.StandardEngine  : Starting Servlet engine: [Apache Tomcat/9.0.63]
2022-05-23 12:23:14.426  INFO 8538 --- [           main] o.a.c.c.C.[Tomcat].[localhost].[/]       : Initializing Spring embedded WebApplicationContext
2022-05-23 12:23:14.426  INFO 8538 --- [           main] w.s.c.ServletWebServerApplicationContext : Root WebApplicationContext: initialization completed in 806 ms
2022-05-23 12:23:14.666  WARN 8538 --- [           main] .s.s.UserDetailsServiceAutoConfiguration : 

Using generated security password: ab60d0d9-a34b-4aee-ad31-e8881672c6a0

This generated password is for development use only. Your security configuration must be updated before running your application in production.

2022-05-23 12:23:14.742  INFO 8538 --- [           main] o.s.s.web.DefaultSecurityFilterChain     : Will secure any request with [org.springframework.security.web.session.DisableEncodeUrlFilter@20eaeaf8, org.springframework.security.web.context.request.async.WebAsyncManagerIntegrationFilter@748ac6f3, org.springframework.security.web.context.SecurityContextPersistenceFilter@7affc159, org.springframework.security.web.header.HeaderWriterFilter@72eb6200, org.springframework.security.web.csrf.CsrfFilter@52bf7bf6, org.springframework.security.web.authentication.logout.LogoutFilter@66de00f2, org.springframework.security.web.authentication.UsernamePasswordAuthenticationFilter@163042ea, org.springframework.security.web.authentication.ui.DefaultLoginPageGeneratingFilter@479b5066, org.springframework.security.web.authentication.ui.DefaultLogoutPageGeneratingFilter@68f6e55d, org.springframework.security.web.authentication.www.BasicAuthenticationFilter@1d8b0500, org.springframework.security.web.savedrequest.RequestCacheAwareFilter@1682c08c, org.springframework.security.web.servletapi.SecurityContextHolderAwareRequestFilter@3fd05b3e, org.springframework.security.web.authentication.AnonymousAuthenticationFilter@6fff46bf, org.springframework.security.web.session.SessionManagementFilter@76ececd, org.springframework.security.web.access.ExceptionTranslationFilter@67e25252, org.springframework.security.web.access.intercept.FilterSecurityInterceptor@52b46d52]
2022-05-23 12:23:14.783  INFO 8538 --- [           main] o.s.b.w.embedded.tomcat.TomcatWebServer  : Tomcat started on port(s): 8080 (http) with context path ''
2022-05-23 12:23:14.791  INFO 8538 --- [           main] c.h.b.demo.springsecurity.Application    : Started Application in 1.471 seconds (JVM running for 1.869)

会生成一串用户密码,这是SpringSecurity 帮助学习的默认配置,后续将会讲解,

启动完成在浏览器输入访问地址:http://localhost:8080/anonymity/hello

网页会自动跳转到 http://localhost:8080/login

输入用户名:user

输入密码:在日志中的一串字符串, 这里是 *ab60d0d9-a34b-4aee-ad31-e8881672c6a0*(由程序自动生成,每次生成内容不一样)

登录成功后跳转到指定地址,得到内容如下:

{"code":200,"data":"hello world","message":null}

至此完成SpringSecurity项目搭建,SpringSecurity 提供了默认配置,默认组织匿名访问接口。

总结

  1. SpringSecurity 项目搭建很方便,结合 SpringBoot 进行使用可以快速完成基础框架搭建,同时提供默认配置,不需要任何配置即可完成项目资源保护
  2. SpringSecurity 提供了 用户身份鉴定(用户登录), 以及用户访问权限控制(判断是否拥有权限访问项目接口)

上述内容帮助完成搭建基础项目,当然这样的程序无法满足实际项目需求,我们需要自定义认证(登录方式)以及 鉴权(权限控制)流程,下一篇我们将在此基础上自定义登录方式

特别声明:项目采用最新SpringSecurity版本:5..7.1,版本升级带来了一点新变化,可能与老版本由一点不同,但是核心理念是一致的

二、SpringSecurity 自定义手机验证登录方式

简介

在上一篇文章中,我们介绍了如何搭建一套基于SpringSecuity的项目框架,并且进行了演示,本文将继续扩展项目功能,实现自定义用户登录功能。

项目源码仓库:Gitee

代码分支:lesson2

原理介绍

SpringSecurity 提供了web服务项目相关的安全配置,通常我们使用 Spring MVC进行开发(基于Servlet 容器技术实现,现在 Spring 提供了 WebFlux 技术可以提高系统吞吐量,两者都是基于 HTTP协议开发的web服务,MVC提供的是阻塞I/O,WebFlux 提供非阻塞 I/O),Servlet 容器中提供了两种核心组件:

  1. Filter
  2. Servlet

Filter 简介

Filter 组件可以实现过滤器功能,Http 请求达到时,Filter 优先接收到请求信息,并且可以依据业务逻辑对请求提前进行处理,例如,CORS(浏览器的同源请求策略, 详细信息参见:阮一峰网络日志), 对于非同源的请求,项目方可以按照要求选择拒绝或者接受请求,接口定义如下:

public interface Filter {
    /// 初始化
    public default void init(FilterConfig filterConfig) throws ServletException {}
    /// 过滤
    public void doFilter(ServletRequest request, ServletResponse response,
            FilterChain chain) throws IOException, ServletException;
    /// 销毁
    public default void destroy() {}
}

在 public void doFilter(ServletRequest request, ServletResponse response, FilterChain chain) throws IOException, ServletException; 方法中可以获取到http 请求信息,并且可以阻断 Http 请求,防止调用实际的业务逻辑代码,例如 实现基于IP黑名单过滤器,发现请求IP 在系统的 IP黑名单中,可以直接返回错误信息阻止请求继续执行。

Servlet 简介

Servlet 组件用于接收Http请求信息,并依据请求信息进行处理,项目的业务逻辑在 Servlet 中进行处理,接口定义如下:

public interface Servlet {
    /// 初始化方法
    public void init(ServletConfig config) throws ServletException;
    /// 获取配置
    public ServletConfig getServletConfig();
    //// 业务处理
    public void service(ServletRequest req, ServletResponse res)
            throws ServletException, IOException;
    /// 获取基础信息
    public String getServletInfo();
    /// 销毁
    public void destroy();
}

其中最主要的是 public void service(ServletRequest req, ServletResponse res) throws ServletException, IOException; 业务逻辑代码在此处进行调用处理(Spring MVC 中的重要组件 DispatcherServlet 是Servlet 子类,通过 service 方法接收并处理 Http 请求)

SpringSecurity原理

通过上述Servlet 技术简单讲解,我们知道Filter主要用于实现过滤功能,这些功能与业务逻辑关系不大,可以在请求进入业务逻辑之前进行拦截处理,保障系统稳定运行,SpringSecurity正是通过一系列“Filter组件”来实现安全过滤功能(在执行业务逻辑之前对Http请求进行身份校验和权限控制),SpringSecurity 中的两个主要功能分别是:

  • 身份校验:对当前发起请求的用户(可能是真实用户,也可能是网络爬虫或者是恶意攻击者)进行身份识别,主要解决你是谁的问题
  • 权限控制:对当前访问资源进行权限控制(管理后台功能只对管理员开发,普通用户无法访问),主要解决你是否有权限访问资源

通过上述两个功能点可以实现系统的访问控制安全,对于不符合要求的请求,直接返回错误信息,阻止不安全的资源访问。本文重点讲解身份校验,权限控制将在后续进行分析。

用户名密码登录分析

在之前文章中使用“user”用户进行了登录,这是SpringSecurity提供的默认用户密码登录实现:"org.springframework.security.web.authentication.UsernamePasswordAuthenticationFilter",这是一个Filter 子类,可以实现Filter过滤功能,核心代码如下:

public class UsernamePasswordAuthenticationFilter extends AbstractAuthenticationProcessingFilter {
    /// 默认匹配 POST /login 请求
    private static final AntPathRequestMatcher DEFAULT_ANT_PATH_REQUEST_MATCHER = new AntPathRequestMatcher("/login",
            "POST");

    public UsernamePasswordAuthenticationFilter() {
        super(DEFAULT_ANT_PATH_REQUEST_MATCHER);
    }

    public UsernamePasswordAuthenticationFilter(AuthenticationManager authenticationManager) {
        super(DEFAULT_ANT_PATH_REQUEST_MATCHER, authenticationManager);
    }
    //// 在 doFilter 方法中调用该方法实现过滤
    @Override
    public Authentication attemptAuthentication(HttpServletRequest request, HttpServletResponse response)
            throws AuthenticationException {
        /// 判断请求方法是否支持
        if (this.postOnly && !request.getMethod().equals("POST")) {
            throw new AuthenticationServiceException("Authentication method not supported: " + request.getMethod());
        }
        String username = obtainUsername(request);
        username = (username != null) ? username.trim() : "";
        String password = obtainPassword(request);
        password = (password != null) ? password : "";
        /// 包装成 用户密码Token
        UsernamePasswordAuthenticationToken authRequest = UsernamePasswordAuthenticationToken.unauthenticated(username,
                password);
        // Allow subclasses to set the "details" property
        /// 设置请求信息,这是一些额外的信息,例如用户IP地址等信息,与核心校验逻辑关系不大
        setDetails(request, authRequest);
        //// 调用 AuthenticationManager 进行身份验证,成功返回 Authentication 对象,失败抛出异常
        return this.getAuthenticationManager().authenticate(authRequest);
    }
}

我们来分析一下 attemptAuthentication 方法,执行的逻辑如下:

  1. 判断Http请求是否为用户密码登录请求(主要看请求路径是否为 /login 并且为 POST方法)
  2. 获取请求中的用户名密码信息(登录参数信息)
  3. 委托给AuthenticationManager组件进行身份验证
  4. 返回成功或者是错误信息

可以理解为Filter中并没有承担核心的身份信息校验责任,主要完成校验请求是否为用户名密码请求,如果是提取出相关参数,委托给AuthenticationManager组件校验身份,如果成功返回Authentication对象。这里有几个关键的类:

  • UsernamePasswordAuthenticationToken:保存用户名密码信息(是Authentication的子类)
  • Authentication:代表待验证信息或者是已验证完成后的身份信息(可以是未验证的信息也可以是已验证的身份信息,通过方法boolean isAuthenticated() 返回值判断是为已验证信息)
  • AuthenticationManager:验证管理器负责对待验证信息内容进行验证,验证成功返回身份信息,失败返回错误信息

UsernamePasswordAuthenticationToken和Authentication都是数据模型类,不存在处理逻辑,AuthenticationManager是主要的验证逻辑处理类,在SpringSecurity 中提供了ProviderManager实现类,核心代码如下:

public class ProviderManager implements AuthenticationManager, MessageSourceAware, InitializingBean{
   /// 身份校验处理器
   private List<AuthenticationProvider> providers = Collections.emptyList();

   @Override
    public Authentication authenticate(Authentication authentication) throws AuthenticationException {
        Class<? extends Authentication> toTest = authentication.getClass();
        AuthenticationException lastException = null;
        AuthenticationException parentException = null;
        Authentication result = null;
        Authentication parentResult = null;
        /// 循环使用Provider 来验证身份信息,只要有一个验证通过就算成功
        for (AuthenticationProvider provider : getProviders()) {
            /// 判断Provider是否支持验证Authentication子类类型,例如前面的UsernamepasswordAuthenticationToken
            if (!provider.supports(toTest)) {
                continue;
            }
            try {
                /// 使用具体的验证器进行验证,验证通过返回具体验证信息
                result = provider.authenticate(authentication);
                if (result != null) {
                    break;
                }
            }
            catch (AccountStatusException | InternalAuthenticationServiceException ex) {
                throw ex;
            }
            catch (AuthenticationException ex) {
                lastException = ex;
            }
        }
        /// 如果验证器无法验证,并且存在父级验证器那么使用父级验证器进行验证
        if (result == null && this.parent != null) {
            // Allow the parent to try.
            try {
                parentResult = this.parent.authenticate(authentication);
                result = parentResult;
            }
            catch (ProviderNotFoundException ex) {
            }
            catch (AuthenticationException ex) {
                parentException = ex;
                lastException = ex;
            }
        }
        /// 判断是否存在已验证结果,存在返回验证信息,不存在抛出一样信息
        if (result != null) {
            if (this.eraseCredentialsAfterAuthentication && (result instanceof CredentialsContainer)) {

                ((CredentialsContainer) result).eraseCredentials();
            }

            if (parentResult == null) {
                this.eventPublisher.publishAuthenticationSuccess(result);
            }

            return result;
        }
        throw lastException;
    }
}

在上述代码中最主要的是private List providers = Collections.emptyList();属性信息,ProviderManager委托该属性循环处理Authentication子类对象,直到验证通过或者是全部不通过, AuthenticationProvider接口定义如下:

public interface AuthenticationProvider {

    ///对待验证信息进行验证
    Authentication authenticate(Authentication authentication) throws AuthenticationException;

    //// 判断当前验证器是否支持对该类型验证信息进行校验处理
    boolean supports(Class<?> authentication);

}

SpringSecurity中提供了对UsernamepasswordAuthenticationToken参数验证的AuthenticationProvider子类DaoAuthenticationProvider,相关接口实现如下:

  • 判断是否支持方法
/// 判断待验证参数authentication是否为UsernamePasswordAuthenticationToken类型或者是其子类
public boolean supports(Class<?> authentication) {
    return (UsernamePasswordAuthenticationToken.class.isAssignableFrom(authentication));
}
  • 身份验证逻辑
@Override
    public Authentication authenticate(Authentication authentication) throws AuthenticationException {

        String username = determineUsername(authentication);
        boolean cacheWasUsed = true;
        UserDetails user = this.userCache.getUserFromCache(username);
        if (user == null) {
            cacheWasUsed = false;
            try {
                /// 依据用户名以及参数信息查找用户信息
                user = retrieveUser(username, (UsernamePasswordAuthenticationToken) authentication);
            }
            catch (UsernameNotFoundException ex) {
                /// 这里为了隐藏用户不存在错误,会对该错误进行包装,抛出新错误
                if (!this.hideUserNotFoundExceptions) {
                    throw ex;
                }
                throw new BadCredentialsException(this.messages
                        .getMessage("AbstractUserDetailsAuthenticationProvider.badCredentials", "Bad credentials"));
            }
        }
        try {
            //// 信息校验前执行
            this.preAuthenticationChecks.check(user);
            //// 校验用户密码是否正确
            additionalAuthenticationChecks(user, (UsernamePasswordAuthenticationToken) authentication);
        }
        catch (AuthenticationException ex) {
            if (!cacheWasUsed) {
                throw ex;
            }
            //// 如果使用的是缓存,那么进行绕过缓存再次验证防止缓存信息过期
            cacheWasUsed = false;
            user = retrieveUser(username, (UsernamePasswordAuthenticationToken) authentication);
            this.preAuthenticationChecks.check(user);
            additionalAuthenticationChecks(user, (UsernamePasswordAuthenticationToken) authentication);
        }
        /// 校验结束处理
        this.postAuthenticationChecks.check(user);
        if (!cacheWasUsed) {
            this.userCache.putUserInCache(user);
        }
        Object principalToReturn = user;
        if (this.forcePrincipalAsString) {
            principalToReturn = user.getUsername();
        }
        /// 验证成功返回成功信息
        return createSuccessAuthentication(principalToReturn, authentication, user);
    }

整体流程图如下所示:

img

自定义验证流程

通过上述分析我们可以知道,自定义一个身份验证逻辑需要实现以下三个组件:

  1. 自定义验证参数类型:Authentication
  2. 自定义拦截过滤器:Filter
  3. 自定义特定验证参数类型验证器:AuthenticationProvider

下面我们将实现常用的手机验证码验证登录功能。

自定义验证参数

通过分析UsernamePasswordAuthenticationToken组件,我们知道该Token主要包装验证参数信息,方便后续使用,实现逻辑如下:

public class PhoneCodeAuthenticationToken extends AbstractAuthenticationToken {

    private final Object principal;

    private Object credentials;
    ///未验证参数构造器
    private PhoneCodeAuthenticationToken(String phone, String code) {
        super(null);
        this.principal = phone;
        this.credentials = code;
        /// 设置是否验证:false-未验证,true-已验证
        super.setAuthenticated(false);
    }
    ///已验证参数构造器
    /// authorities代表取得打权限信息
    private PhoneCodeAuthenticationToken(Object principal,
            Collection<? extends GrantedAuthority> authorities) {
        super(authorities);
        this.principal = principal;
         /// 设置是否验证:false-未验证,true-已验证
        super.setAuthenticated(true);
    }

    /// 未验证Token
    public PhoneCodeAuthenticationToken unAuthToken(String phone, String code) {
        return new PhoneCodeAuthenticationToken(phone, code);
    }

    ////已验证Token
    public PhoneCodeAuthenticationToken authToken(Object principal,
            Collection<? extends GrantedAuthority> authorities) {
        return new PhoneCodeAuthenticationToken(principal, authorities);
    }

    @Override
    public Object getCredentials() {
        return credentials;
    }

    @Override
    public Object getPrincipal() {
        return principal;
    }
}

自定义拦截器

通过分析UsernamePasswordAuthenticationFilter组件,我们知道拦截器主要完成三个功能:

  1. 拦截特定请求
  2. 解析参数
  3. 委托给验证器进行验证处理

手机验证码登录拦截POST/phone/login请求,解析参数,包装成PhoneCodeAuthenticationToken对象,最后委托给AuthenticationManager组件验证,具体代码参见Gitee仓库:地址

运行验证

程序启动后发起POST请求,参数信息:

  • phone:15000000000
  • code:888888

请求成功后返回用户信息(SpringSecurity默认配置会将登录成功请求跳转到 / 路径):

{
    "code": 200,
    "data": {
        "username": "15000000000",
        "phone": "15000000000",
        "roles": [
            "ROLE_USER"
        ]
    },
    "message": null
}

至此完成手机验证码登录功能

我们使用POST请求,SpringSecurity默认提供csrf保护,会拦截 POST请求,因此需要禁用

总结

  • SpringSecurity 使用Servlet容器组件Filter功能进行请求拦截,实现身份校验以及权限控制
  • SpringSecurity 使用AuthenticationManager来实现身份校验功能(实际上你可以在Filter中直接完成身份验证功能,但是这种硬编码方式会增加程序耦合性,后期维护/扩展不方便)
  • SpringSecurity 中的AuthenticationManager委托多个AuthenticationProvider对请求参数进行校验
  • 自定义手机验证码验证流程需要实现三个类:
    • 继承Filter的PhoneCodeAuthenticationFilter, 对手机验证码登录请求进行拦截,并解析处请求参数信息,最后委托给AuthenticationManager进行身份校验
    • 继承Authentication的PhoneCodeAuthenticationToken
    • 登录时存放请求参数信息:手机号和验证码
    • 登录成功后存放用户信息:用户名、手机号、权限等
    • 继承AuthenticationProvider的PhoneCodeAuthenticationProvider,对请求参数进行验证,验证通过返回用户信息

有过SpringSecurity开发经验的同学会发现仓库中的代码使用HttpSecurity进行配置的方式与之前的方式不同,这是SpringSecurity官方在新版中推荐使用的方式,老版本的配置方式将会被遗弃,目前两种方式都可以使用

参考文档

三、SpringSecurity 动态权限访问控制

简介

在先前文章中我们搭建了SpringSecurity项目,并且讲解了自定义登录方式需要做哪些工作,如果你感兴趣可以前往博客阅读文章以及代码,在本文将继续讲解如何实现动态权限控制。

代码仓库:Gitee

代码分支:lesson3

目标

Web项目通常都有前台和后台服务,前台服务面向目标客户,后台服务为项目方提供管理和数据分析服务,因此不同的用户需要赋予不同的角色,例如前台用户角色为USER,后台用户为ADMIN,USER允许访问"/user/hello"接口,ADMIN允许访问"/admin/hello"接口,但是USER不能访问。这是项目必须有的基本功能,同时访问规则也会不断变化,例如: 有一个用户昵称功能,初期只允许会员用户(可以理解为拥有角色VIP的用户)使用,后期产品决定全员都可以使用,这种需求也很常见,如果采用硬编码的方式那么会导致频繁修改代码,测试、发布,增加额外工作量,如果可以动态配置接口访问权限,那么就能减少很多工作量,SpringSecurity框架提供了扩展点,基于这些扩展点可以很方便的实现动态权限控制访问功能,我们再来回顾一下需求:

  • 基于角色进行接口权限控制
  • 访问接口需要的角色可以动态配置

原理分析

通过上一篇文章我们知道SpringSecurity基于Filter实现身份验证和权限控制功能,SpringSecurity提供了实现类FilterSecurityInterceptor对访问路径进行权限控制,核心代码逻辑如下:

public void invoke(FilterInvocation filterInvocation) throws IOException, ServletException {
        ///此处省略无关逻辑
        /// 在这里执行权限控制逻辑
        InterceptorStatusToken token = super.beforeInvocation(filterInvocation);
        try {
            filterInvocation.getChain().doFilter(filterInvocation.getRequest(), filterInvocation.getResponse());
        }
        finally {
            super.finallyInvocation(token);
        }
        super.afterInvocation(token, null);
}

在访问实际业务逻辑之前调用父级方法beforeInvocation进行权限判断,如果权限不符合要求,直接抛出异常阻止访问实际业务逻辑,核心代码如下:

protected InterceptorStatusToken beforeInvocation(Object object) {
        //此处省略无关代码
        //// 这里获取与访问路径相关的权限信息,例如:/user/hello 对应 ROLE_USER 角色,当然一个路径可能对应多个权限
        Collection<ConfigAttribute> attributes = this.obtainSecurityMetadataSource().getAttributes(object);
        if (CollectionUtils.isEmpty(attributes)) {
            /// 这里注意如果对应的路径在系统中没有配置权限或者是获取方法没有处理这种请求会导致放行,特别注意
            return null; // no further work post-invocation
        }
        /// 未登录用户直接返回未验证错误
        if (SecurityContextHolder.getContext().getAuthentication() == null) {
            credentialsNotFound(this.messages.getMessage("AbstractSecurityInterceptor.authenticationNotFound",
                    "An Authentication object was not found in the SecurityContext"), object, attributes);
        }
        //// 获取验证信息
        Authentication authenticated = authenticateIfRequired();
        // Attempt authorization
        /// 判断用户是否拥有访问权限
        attemptAuthorization(object, attributes, authenticated);
        /// 这里实现了类似Linux su 命令,将当前用户暂时赋予另外一个用户运行权限,可以先忽略不看
        // Attempt to run as a different user
        Authentication runAs = this.runAsManager.buildRunAs(authenticated, object, attributes);
        if (runAs != null) {
            SecurityContext origCtx = SecurityContextHolder.getContext();
            SecurityContext newCtx = SecurityContextHolder.createEmptyContext();
            newCtx.setAuthentication(runAs);
            SecurityContextHolder.setContext(newCtx);

            if (this.logger.isDebugEnabled()) {
                this.logger.debug(LogMessage.format("Switched to RunAs authentication %s", runAs));
            }
            // need to revert to token.Authenticated post-invocation
            return new InterceptorStatusToken(origCtx, true, attributes, object);
        }
        this.logger.trace("Did not switch RunAs authentication since RunAsManager returned null");
        // no further work post-invocation
        return new InterceptorStatusToken(SecurityContextHolder.getContext(), false, attributes, object);

    }

这里有两个重点内容:

  1. 通过this.obtainSecurityMetadataSource().getAttributes(object);方法来获取访问所需要的权限信息
  2. 通过attemptAuthorization(object, attributes, authenticated)方法对访问所需的权限以及用户身份信息进行决策,判断是否允许访问

路径权限分析

上述的this.obtainSecurityMetadataSource()方法返回SecurityMetadataSource类型对象,该接口核心代码如下:

public interface SecurityMetadataSource extends AopInfrastructureBean {
    /// 依据object获取权限信息,我们可以把ConfigAttribute理解为String类型,在我们系统中可以理解保存着角色信息,例如ROLE_USER
    Collection<ConfigAttribute> getAttributes(Object object) throws IllegalArgumentException;
    ///获取系统中配置的所有权限信息,用于后续验证器判断是否支持该类型决策
    Collection<ConfigAttribute> getAllConfigAttributes();
    /// object 类型,用于判断SecurityMetadataSource支持解析的object类型
    boolean supports(Class<?> clazz);
}

可以看出SecurityMetadataSource的主要作用是给出当前访问需要哪些权限,方便后续判断,可以理解为一个数据源,用来获取访问权限列表

访问权限控制分析

这里我们需要重点查看方法attemptAuthorization(object, attributes, authenticated);包含对用户访问控制权限进行判断,核心代码如下:

private void attemptAuthorization(Object object, Collection<ConfigAttribute> attributes,
            Authentication authenticated) {
        try {
            /// 委托accessDecisionManager进行决策判断
            this.accessDecisionManager.decide(authenticated, object, attributes);
        }
        catch (AccessDeniedException ex) {
            /// 异常请求直接向上抛出异常信息
            throw ex;
        }
}

这个方法很简单,就是委托accessDecisionManager来进行访问决策,我们来看一下这个接口的核心代码:

public interface AccessDecisionManager {
    //// 对访问进行决策,判断是否有权限
    void decide(Authentication authentication, Object object, Collection<ConfigAttribute> configAttributes)
            throws AccessDeniedException, InsufficientAuthenticationException;
    /// 查看访问控制器是否支持该类型决策
    boolean supports(ConfigAttribute attribute);
    /// 查看访问控制器是否支持特定类型,这个类型就是 上面方法中object对应的类型
    boolean supports(Class<?> clazz);
}

接口也很简单,有点类似AuthenticationProvider接口,调用decide方法,如果允许访问那么不进行任何处理,如果不允许访问就抛出异常信息。

代码实现梳理分析

上述核心逻辑很简单,但是实现逻辑有点绕,不要紧我们画个流程图再来梳理一遍(觉得绕主要是不相关代码对理解造成了困扰,还有可能就是被这种俄罗斯套娃形式绕晕了)

img

通过上述分析我们可以发现,权限控制需要两个核心功能:

  1. 访问路径所需要的权限(实现接口SecurityMetadataSource)
  2. 依据用户权限、路径所需权限进行决策判断是否允许访问(实现接口AccessDecisionManager)

完成上述功能后,将这些功能组装成FilterSecurityInterceptor类型对象,然后放置到SpringSecurity过滤链中实现过滤功能

代码实现

直接动手实现动态权限控制

实现路径权限获取

这里为了更加贴近实际项目,将提供一个RoleService作为数据源,实现代码如下:

@Service
public class RoleService {
    public List<ConfigAttribute> roles = new ArrayList<>();
    private Map<String, List<ConfigAttribute>> urlRoleMaps = new HashMap<>();
    @PostConstruct
    public void init() {
        /// 初始化数据
        roles.addAll(SecurityConfig.createList("ROLE_USER", "ROLE_ADMIN", "ROLE_VIP"));
        urlRoleMaps.put("/", SecurityConfig.createList("ROLE_USER", "ROLE_ADMIN"));
        urlRoleMaps.put("/user/hello", SecurityConfig.createList("ROLE_USER", "ROLE_ADMIN"));
        urlRoleMaps.put("/user/nickname", SecurityConfig.createList("ROLE_VIP"));
        urlRoleMaps.put("/admin/hello", SecurityConfig.createList("ROLE_ADMIN"));
    }
    /// 获取所有角色信息
    public Collection<ConfigAttribute> getAllRoles() {
        return Collections.unmodifiableList(roles);
    }
    ///依据请求路径查询所需权限
    public Collection<ConfigAttribute> getRoleByPath(String path) {
        Collection<ConfigAttribute> roles = urlRoleMaps.get(path);
        if (roles == null) {
            return Collections.EMPTY_LIST;
        }
        return Collections.unmodifiableCollection(roles);
    }
}

代码很简单,就是初始化数据,提供路径与权限对应的数据服务,在实际项目中通常从数据库中获取这些信息。

下面编写RolePermissionMetadataSource接口的实现类,代码如下:

/// FilterInvocationSecurityMetadataSource 是SecurityMetadataSource的子接口,实际上就是 SecurityMetadataSource,没有扩展任何方法
public class RolePermissionMetadataSource implements FilterInvocationSecurityMetadataSource {
    @Autowired
    private RoleService roleService;
    @Override
    public Collection<ConfigAttribute> getAttributes(Object object)
            throws IllegalArgumentException {
        FilterInvocation invocation = (FilterInvocation) object;
        String url = invocation.getRequestUrl();
        /// 通过请求路径获取访问路径所需的权限列表
        Collection<ConfigAttribute> roles = roleService.getRoleByPath(url);
        if (roles != null && roles.size() > 0) {
            return roles;
        }
        //没有匹配上的资源,禁止访问,设置不存在的访问权限
        // 通过之前的分析知道,如果这里返回空,将会直接放行,运行登录用户访问,这是有风险的
        return SecurityConfig.createList(RoleEnums.ROLE_REFUSE.name());
    }
    @Override
    public Collection<ConfigAttribute> getAllConfigAttributes() {
        return roleService.getAllRoles();
    }
    @Override
    public boolean supports(Class<?> clazz) {
        return FilterInvocation.class.isAssignableFrom(clazz);
    }
}

实现路径访问控制决策类

继承接口AccessDecisionManager,核心代码如下:

public class PathAccessDecisionManager implements AccessDecisionManager {
    ///拒绝访问权限名称
    private static final String BASE_REFUSE_NAME = RoleEnums.ROLE_REFUSE.name();
    @Override
    public void decide(Authentication authentication, Object object,
            Collection<ConfigAttribute> configAttributes)
            throws AccessDeniedException, InsufficientAuthenticationException {
        Iterator<ConfigAttribute> iterator = configAttributes.iterator();
        //进行权限匹配,如果用户拥有资源权限那么进行放行操作
        while (iterator.hasNext()) {
            ConfigAttribute ca = iterator.next();
            // 当前请求需要的权限
            String needRole = ca.getAttribute();
            if (RoleEnums.ROLE_ANONYMOUS.name().equalsIgnoreCase(needRole)) {
                return;
            }
            if (BASE_REFUSE_NAME.equalsIgnoreCase(needRole)) {
                if (authentication instanceof AnonymousAuthenticationToken) {
                    //匿名用户
                    throw new AccessDeniedException("资源信息不存在");
                } else {
                    //登录用户
                    throw new AccessDeniedException("权限不足!");
                }
            }
            // 当前用户所具有的权限
            Collection<? extends GrantedAuthority> authorities = authentication.getAuthorities();
            for (GrantedAuthority authority : authorities) {
                if (authority.getAuthority().equalsIgnoreCase(needRole)) {
                    return;
                }
            }
        }
        //如果当前请求没有验证,返回未验证异常
        if (authentication instanceof AnonymousAuthenticationToken) {
            throw new AccessDeniedException("用户未登录");
        }
        throw new AccessDeniedException("权限不足!");
    }
    @Override
    public boolean supports(ConfigAttribute attribute) {
        return true;
    }
    @Override
    public boolean supports(Class<?> clazz) {
        return FilterInvocation.class.isAssignableFrom(clazz);
    }
}

组装Filter

我们将上述实现类与FilterSecurityInterceptor进行组装,实现权限动态过滤:

///动态权限控制 Filter, 默认会拦截所有请求进行权限判断
private FilterSecurityInterceptor filterSecurityInterceptor() {
    FilterSecurityInterceptor interceptor = new FilterSecurityInterceptor();
    /// 由于包含Spring Bean,因此需要注入实现,而不是直接new
    interceptor.setSecurityMetadataSource(rolePermissionMetadataSource);
    interceptor.setAccessDecisionManager(new PathAccessDecisionManager());
    return interceptor;
}
///加入到SpringSecurity过滤链中
httpSecurity.addFilterBefore(roleAuthFilter, FilterSecurityInterceptor.class);

运行验证

我们代码里创建了三个用户:

  1. 15000000000, 拥有:ROLE_USER
  2. 15666666666, 拥有:ROLE_USER、ROLE_VIP
  3. 15888888888, 拥有:ROLE_USER、ROLE_ADMIN

程序运行完成后使用15000000000进行手机验证码登录:

- POST http://localhost:8080/phone/login?phone=15000000000&code=888888

- 返回:

{
    "code": 200,
    "data": {
        "username": "15000000000",
        "phone": "15000000000",
        "roles": [
            "ROLE_USER"
        ]
    },
    "message": null
}

可以看到拥有ROLE_USER权限,那么我们访问 http://localhost:8080/user/hello, 返回:

{
    "code": 200,
    "data": "Hello User",
    "message": null
}

访问http://localhost:8080/admin/hello, 返回:

{
    "code": 400,
    "message": "请求受限"
}

我们看到以上结果符合预期,同理可以使用15888888888用户进行同样的访问操作,在这里我们就不做过多介绍,大家有兴趣可以下载代码自行运行测试,文章开头有代码地址。

如果要修改路径对应的权限,那么只要修改RoleService中的数据即可实现权限动态配置。

总结

通过上述文章分析,我们已经完成权限动态配置,当然运行中展现的JSON数据是配置了对应处理器处理的结果,细节处理请前往代码仓库下载源码自行查看。

为了完成动态权限我们需要完成三个步骤,实现两个接口,步骤如下:

  1. 实现路径权限数据访问接口(实现SecurityMetadataSource)
  2. 实现访问控制决策接口(实现AccessDecisionManager)
  3. 组装Filter并加入到过滤链中(FilterSecurityInterceptor)

熟练掌握上述步骤,实现动态权限控制将不再是难题。

在SpringSecurity 中不仅提供了FilterSecurityInterceptor实现类来对访问进行权限控制,在SpringSecyrity 5.4版本中还提供了AuthorizationFilter实现类来实现相同功能,具体实现方式自行前往仓库进行查看

参考文档

四、SpringSecurity OAuth2统一授权服务

代码

代码仓库:地址

代码分支:lesson4

简介

在先前文章中我们实战演练了在SpringBoot单体应用中使用SpringSecurity开发自定义登录流程以及动态权限控制,有兴趣的同学可以前往博客阅读SpringSecurity相关文章(所有代码都已上传到Gitee仓库, 每篇文章都有一个专属分支)。随着业务的扩张,单体应用无法满足业务需求,微服务是当前大型商业服务的主流架构,在Java领域中,SpringCloud 全家桶是微服务主流开发框架,Spring Cloud Alibaba 是在Spring Cloud 的基础上进行扩展,目的是为了更好的搭配使用Alibaba生态中的微服务组件(Nacos、Sentinel、Seata等),具体内容可查看文档。

微服务框架如下所示:

img

上图中的服务集群代表具体业务服务,微服务下的权限控制是为了实现服务集群的安全访问,每个服务集群包含1~N个服务,如果每个服务都定义一套安全策略,那么后期维护将会是一个大工程,因此需要统一安全策略,实现全局安全访问。

OAuth2

OAuth是一个关于授权(authorization)的开放网络标准,2.0版本在全世界得到广泛应用,网上有很多讲解OAuth2的文档,如果你对OAuth没有概念,可以查看往期文章:理解OAuth2.0、OAuth2.0的一个简单解释、OAuth2.0的四种调用方式。

假设我们是一家公司,公司内部有以下业务部门:

  • 微信:提供聊天、账户系统,同时对外提供账户授权登录服务
  • 电子支付:提供在线支付服务
  • 王者荣耀:提供王者荣耀游戏服务

img

为了实现各个业务系统之间相互调用,需要一套授权系统,对内外系统提供统一的授权访问服务,OAuth协议能够满足上述要求,只需要一套授权服务就能同时满足内外系统的调用要求。

SpringSecurity OAuth

OAuth 涉及四个角色:

  • 用户:实际拥有资源所有权的使用者,例如:张三
  • 客户端:提供应用功能的程序客户端,可以是APP形式、web形式,例如:时间海绵博客()
  • 授权服务:实现OAuth2授权协议,对外提供授权服务,例如:微信开发平台
  • 资源服务:对外提供资源访问服务,例如:我们使用微信的微信扫码登录时为网站提供用户信息的微信服务

下面以微信登录为例:

img

传统的模式中,我们默认web客户端是可信任的,所以没有用户授权的过程,可以访问任何数据。在OAuth协议中,客户端默认是不可信任的,需要进行授权处理。无论是内部客户端还是外部客户端都需要得到授权服务器的认证,但是内部和外部客户端可以使用不同的授权认证方式(比如最严格的授权码方式和最简单的客户端方式)。

授权服务器

使用 spring-security-oauth2 搭建授权服务器,对外提供以下功能:

  • /oauth/authorize 获取授权码
  • /oauth/token 提供客户端(client_credentials)、简化(implicit)、密码(password)、刷新Token(refresh_token),授权码(authorization_code)方式获取access_token
  • /oauth/check_token 依据access_token获取授权信息

同时需要管理客户端信息,客户端信息包含以下属性:

  • clientId 客户端Id
  • clientSecret 客户端密码
  • resourceId 资源Id,用于指定可以访问哪些资源服务
  • authorizedGrantTypes 授权模式,客户端(client_credentials)、简化(implicit)、密码(password)、授权码(authorization_code)
  • scopes 授权范围,可以指定客户端访问权限,这个在协议中没有明确指明作用,各个授权服务可以基于业务自行处理
  • authorities 权限,客户端授权模式下需要配置,其它模式下可以不配置
  • redirectUri 授权结果跳转地址

在先前的文章中提到了UserDetailsService用于查询用户信息,在OAuth中需要提供ClientDetailsService来查询客户端信息,代码如下:

public class AuthClientDetailService implements ClientDetailsService {
    private ClientDetailsService clientDetailsService;
    public AuthClientDetailService() {
        InMemoryClientDetailsServiceBuilder builder = new InMemoryClientDetailsServiceBuilder();
        builder.withClient("blog")// client_id
                .secret(new BCryptPasswordEncoder().encode("blog"))
                .resourceIds("blog", "resource")
                .authorizedGrantTypes("authorization_code", "password", "client_credentials", "implicit", "refresh_token")// 该client允许的授权类型 authorization_code,password,refresh_token,implicit,client_credentials
                .scopes("all", "user")// 允许的授权范围
                .autoApprove(false) //加上验证回调地址
                .authorities("blog")
                .accessTokenValiditySeconds(60 * 60 * 2) // 令牌默认有效期2小时
                .refreshTokenValiditySeconds(60 * 60 * 24 * 3) // 刷新令牌默认有效期3天
                .redirectUris("https://blog.hzchendou.com");
        try {
            this.clientDetailsService = builder.build();
        } catch (Exception e) {
            System.exit(-1);
        }
    }
    @Override
    public ClientDetails loadClientByClientId(String clientId) throws ClientRegistrationException {
        return this.clientDetailsService.loadClientByClientId(clientId);
    }
}

同时需要引入AuthorizationServerConfigurer对授权服务进行配置,代码如下:

@Configuration
@EnableAuthorizationServer
public class AuthorizationServerConfig extends AuthorizationServerConfigurerAdapter {
    @Autowired
    private AuthClientDetailService authClientDetailService;
    @Autowired
    private AuthenticationManager authenticationManager;
    /**
     * 配置客户端信息
     *
     * @param clients
     * @throws Exception
     */
    @Override
    public void configure(ClientDetailsServiceConfigurer clients) throws Exception {
        clients.withClientDetails(authClientDetailService);
    }
    /**
     * 配置OAuth token相关配置
     *
     * @param endpoints
     * @throws Exception
     */
    @Override
    public void configure(AuthorizationServerEndpointsConfigurer endpoints) throws Exception {
        endpoints
                .authenticationManager(authenticationManager)/// 用于密码模式验证时需要提供客户端身份验证
                .reuseRefreshTokens(false);
    }
    @Override
    public void configure(AuthorizationServerSecurityConfigurer security) throws Exception {
        security.authenticationEntryPoint(new AuthAuthenticationEntryPoint());
        security.tokenKeyAccess("permitAll()")//oauth/token_key是公开
                .checkTokenAccess("permitAll()");//oauth/check_token公开
        /// 配置使用 客户端id 和密码的方式进行登录(使用明文传输,不推荐) 
        AuthClientCredentialsTokenEndpointFilter endpointFilter = new AuthClientCredentialsTokenEndpointFilter(security);
        endpointFilter.afterPropertiesSet();
        endpointFilter.setAuthenticationEntryPoint(new AuthAuthenticationEntryPoint());
        // 客户端认证之前的过滤器
        security.addTokenEndpointAuthenticationFilter(endpointFilter);
    }
}

使用@EnableAuthorizationServer注解配置会自动配置AuthorizationEndpoint、TokenEndpoint、CheckTokenEndpoint接口,提供OAuth授权服务。

至此完成授权服务搭建

资源服务器搭建

资源服务器也是使用spring-security-oauth2进行搭建,通过上面的介绍,我们知道资源服务器需要识别access_token来获取用户授权的信息内容,配置信息如下:

@EnableResourceServer
@Configuration
public class SpringSecurityResourceServerConfig extends ResourceServerConfigurerAdapter {
    public static final String RESOURCE_ID = "resource";
    @Autowired
    TokenStore tokenStore;
    @Override
    public void configure(HttpSecurity httpSecurity) throws Exception {
        httpSecurity.formLogin().disable();
        httpSecurity.exceptionHandling()
                .accessDeniedHandler(new PathAccessDeniedHandler())
                .authenticationEntryPoint(new AuthAuthenticationEntryPoint());
        httpSecurity.authorizeRequests()
                .antMatchers("/admin/**").hasAuthority("admin")
                .antMatchers("/user/**").hasAuthority("user")
                .and().csrf().disable()
                .sessionManagement().sessionCreationPolicy(SessionCreationPolicy.STATELESS);
    }

    @Override
    public void configure(ResourceServerSecurityConfigurer resources) {
        //当前资源服务 id,用于校验授权信息是否能够访问
        resources.resourceId(RESOURCE_ID)
                .tokenStore(tokenStore)
                /// 设置Token服务,用于识别access_token授权信息
                .tokenServices(tokenService())//验证令牌的服务
                .stateless(true);
    }

    //资源服务令牌解析服务
    @Bean
    public ResourceServerTokenServices tokenService() {
        //使用远程服务请求授权服务器校验token,必须指定校验token 的url、client_id,client_secret
        RemoteTokenServices service = new RemoteTokenServices();
        service.setCheckTokenEndpointUrl("http://localhost:8081/oauth/check_token");
        service.setClientId("blog");
        service.setClientSecret("blog");
        return service;
    }
}

通过配置@EnableResourceServer注解,将会在Filter过滤链中添加OAuth2AuthenticationProcessingFilter过滤器拦截请求,依据access_token解析授权信息, 查看请求头是否有Authorization属性,该属性代表access_token信息。

OAuth服务验证

启动授权服务以及资源服务,访问 POST http://localhost:8081/oauth/token,请求参数如下所示(使用密码模式访问):

client_id:blog
client_secret:blog
grant_type:password
username:admin
password:admin

返回结果:

{
    "access_token": "db9898be-2aef-4f86-9486-b3735db6403e",
    "token_type": "bearer",
    "refresh_token": "44d9aeff-8839-4a9b-b6ea-1ea6310bab5e",
    "expires_in": 6971,
    "scope": "all user"
}

启动资源服务,访问 GET http://localhost:8082/admin/hello,请求头信息如下:

Authorization:Bearer db9898be-2aef-4f86-9486-b3735db6403e

返回信息如下:

{
    "code": 200,
    "data": "Hello Admin"
}

说明访问成功

总结

SpringSecurity中的功能都是通过组装Filter链来完成特定功能实现,

  • SpringSecurity OAuth授权服务需要提供对外服务,因此还提供了AuthorizationEndpoint、TokenEndpoint、CheckTokenEndpoint等接口模块,当然你也可以用自己的实现来替换这些服务。
  • SpringSecurity OAuth 资源服务提供了OAuth2AuthenticationProcessingFilter过滤器来解析access_token授权信息,获得权限信息,后续的权限验证流程与先前的SpringSecurity单体应用是一致的

参考文档

五、SpringSecurity OAuth2扩展自定义授权模式

代码

代码仓库:地址

代码分支:lesson5

简介

在上一篇文章中,我们使用SpringSecurity OAuth2搭建了一套授权服务,对业务系统进行统一授权管理。OAuth提供了四种授权方式:

  • 授权码模式(authorization_code)
  • 简化模式(implicit)
  • 客户端(client_credentials)
  • 密码(password)

在实际业务中上述四种模式不能满足所有要求,例如业务系统接入了短信验证码登录方式,需要进行扩展满足业务需求

手机验证码登录

原理分析

SpringSecurity OAuth在使用@EnableAuthorizationServer注解会自动装配TokenEndpoint对象,这个对象会提供一个POST /oauth/token接口,我们以密码授权模式分析调用流程,如下所示:

img

用户在时间海绵博客发起用户名密码登录请求,时间海绵博客服务器端接收到请求后,调用OAuth协议中的密码授权模式发送请求到OAuth授权服务器,请求信息如下:

client_id:blog
client_secret:blog
grant_type:password
username:admin
password:admin

在ClientCredentialsTokenEndpointFilter过滤器中对客户端信息(client_id和client_secret)进行校验,授权成功后才能访问TokenEndpoint接口(这里要注意,对于OAuth授权服务器来说,过滤链主要完成对客户端信息的校验,用户信息在TokenEndpoint中进行校验,这是因为不同的授权模式关注的用户信息类型不同,需要具体问题具体分析)。

在TokenEndpoint中依据验证的客户端信息以及请求的授权模式进行对比,校验客户端是否有权限进行特定类型授权请求,校验通过后委托给TokenGranter组件进行具体授权模式校验,如上图所示,ResourceOwnerPasswordTokenGranter负责对密码授权模式请求校验,SpringSecurity OAuth还提供了以下实现来校验授权请求:

  • AuthorizationCodeTokenGranter负责校验授权码模式(authorization_code)请求
  • ClientCredentialsTokenGranter负责校验客户端模式(client_credentials)请求
  • ImplicitTokenGranter负责校验简化模式(implicit)请求
  • ResourceOwnerPasswordTokenGranter负责校验密码模式(password)请求

TokenGranter校验成功将返回AccessToken信息,后续客户端可以使用AccessToken信息获取到授权用户信息完成对应操作。

通过上述可以知道,如果要扩展实现短信验证码模式,需要自定义实现TokenGranter组件来校验手机验证码授权请求,TokenGranter定义如下所示:

public interface TokenGranter {
    //// 对授权类型以及请求参数进行处理,如果成功则返回AccessToken信息
    OAuth2AccessToken grant(String grantType, TokenRequest tokenRequest);
}

手机验证码模式代码实现

实现一个继承ToeknGranter接口的类,通过分析已有的TokenGranter子类,我们可以很容易实现,定义一个SmsCodeTokenGranter类,代码实现如下所示:

public class SmsCodeTokenGranter extends AbstractTokenGranter {

    /// 授权模式类型,需要与请求字段grant_type的值相等才会进入处理
    private static final String GRANT_TYPE = "sms_code";
    /// 验证手机与验证码信息是否匹配,这里只是简单的进行匹配处理,判断是否是否为15000000000,验证码是否为:888888
    private final PhoneSmsCodeService phoneSmsCodeService;

    public SmsCodeTokenGranter(AuthorizationServerTokenServices tokenServices,
            ClientDetailsService clientDetailsService,
            OAuth2RequestFactory requestFactory, PhoneSmsCodeService phoneSmsCodeService
    ) {
        super(tokenServices, clientDetailsService, requestFactory, GRANT_TYPE);
        this.phoneSmsCodeService = phoneSmsCodeService;
    }

    @Override
    protected OAuth2Authentication getOAuth2Authentication(ClientDetails client, TokenRequest tokenRequest) {

        Map<String, String> parameters = new LinkedHashMap(tokenRequest.getRequestParameters());

        String mobile = parameters.get("mobile"); // 手机号
        String code = parameters.get("code"); // 短信验证码
        ///对参数进行基础校验
        if (StringUtils.isEmpty(mobile) || StringUtils.isEmpty(code)) {
            throw new InvalidGrantException("授权请求参数异常");
        }
        /// 校验手机验证码是否符合要求
        if (!phoneSmsCodeService.checkSmsCode(mobile, code)) {
            throw new InvalidGrantException("授权请求参数异常");
        }

        /// 这里为了简单直接硬编码写入用户信息,通常需要在数据库中取出用户相关信息
        List<GrantedAuthority> roles = new ArrayList<>();
        roles.add( new SimpleGrantedAuthority("user"));
        User user = new User(mobile, "", roles);
        /// 授权模式
        UsernamePasswordAuthenticationToken userAuth = new UsernamePasswordAuthenticationToken(user, null, roles);
        userAuth.setAuthenticated(true);
        OAuth2Request storedOAuth2Request = this.getRequestFactory()
                .createOAuth2Request(client, tokenRequest);
        return new OAuth2Authentication(storedOAuth2Request, userAuth);
    }
}
继承AbstractTokenGranter抽象类(该类继承了TokenGranter接口),我们只需要关注核心逻辑"验证手机验证码信息是否正确",校验正确将返回一个OAuth2Authentication对象,这个对象包含了用户信息以及OAuth2请求信息。
完成这一步后,我们需要将SmsCodeTokenGranter装配到TokenEndpoint组件中对sms_code授权类型进行校验处理,配置逻辑如下:
    public void configure(AuthorizationServerEndpointsConfigurer endpoints) throws Exception {
        endpoints
                .authenticationManager(authenticationManager)/// 用于密码模式验证时需要提供客户端身份验证信息
                .reuseRefreshTokens(false);
        //// 取出系统中的四种模式
        List<TokenGranter> granters = new ArrayList<>(Arrays.asList(endpoints.getTokenGranter()));
        /// 添加手机验证码的授权模式
        granters.add(new SmsCodeTokenGranter(endpoints.getTokenServices(), endpoints.getClientDetailsService(), endpoints.getOAuth2RequestFactory(), phoneSmsCodeService));
        /// 这是一个组装模式,实现了TokenGranter接口,循环调用List中的TokenGranter组件进行校验处理,直到返回验证成功信息或者是异常信息
        CompositeTokenGranter compositeTokenGranter = new CompositeTokenGranter(granters);
        endpoints.tokenGranter(compositeTokenGranter);
    }

上述在AuthorizationServerConfigurerAdapter配置类中进行,具体参见代码。

运行校验

需要注意,这是一个新的授权模式,因此需要先授权客户端拥有手机短信验证模式请求权限,配置客户端的authorizedGrantTypes属性包含sms_code权限(这一步参见代码)。

启动服务,发送手机验证码授权请求POST /oauth/token,请求参数:

client_id:blog
client_secret:blog
grant_type:sms_code
mobile:15000000000
code:888888

返回参数信息:

{
    "access_token": "cf4243c1-085e-4f82-b733-03fb38c90a7c",
    "token_type": "bearer",
    "refresh_token": "01474dd1-c5c5-436b-98b1-8f146cde2391",
    "expires_in": 7199,
    "scope": "all user"
}

手机短信验证码授权模式验证成功,具体代码参见代码仓库。

总结

  • 扩展自定义授权模式,需要继承TokenGranter接口,实现具体校验逻辑
  • 将实现的自定义授权TokenGraner类装配到TokenEndpoint组件中
  • 特别注意需要将新的授权模式grantType信息授权给指定的客户端,不然客户端无法发送自定义授权模式请求,例如本案例中的sms_code请求

参考文档

六、SpringSecurity OAuth2 + SpringCloud Gateway实现统一鉴权管理

代码

代码仓库:地址

代码分支:lesson6

简介

在先前文章中,我们使用SpringSecurity OAuth2搭建了一套基于OAuth2协议的授权系统,并扩展了手机验证码授权模式。在微服务架构下,网关承担着流量入口的角色,所有的请求都要先经过网关,然后由网关负责转发到具体的服务,因此可以在网关实现统一鉴权,网关对请求中的权限进行鉴定,然后将权限信息转发到具体的资源服务,在资源服务中只需要简单校验请求中的权限信息即可(查看信息是否有效),整体流程如下所示:

img

统一鉴权

SpringCloud Gateway网关

我们在上一篇的基础上引入网关服务,在这里使用SpringCloud Gateway组件进行搭建,引入依赖:

<dependency>
  <groupId>org.springframework.cloud</groupId>
  <artifactId>spring-cloud-starter-gateway</artifactId>
</dependency>

网关在OAuth2授权协议中承担着资源服务的角色,对请求进行身份鉴定和访问权限控制,身份鉴定需要访问OAuth2授权服务,因此需要引入OAuth2资源服务以及客户端依赖:

<dependency>
  <groupId>org.springframework.security</groupId>
  <artifactId>spring-security-oauth2-resource-server</artifactId>
</dependency>

<dependency>
  <groupId>org.springframework.boot</groupId>
  <artifactId>spring-boot-starter-oauth2-client</artifactId>
</dependency>

通过之前的文章,我们可以知道SpringSecurity 通过组装一系列的Filter来完成身份验证和权限访问控制功能,但是SpringCloud Gateway使用了新技术框架Reactive Stack(响应式编程),在Spring中提供了Spring WebFlux模块支持响应式编程,传统的Spring MVC都是基于阻塞I/O编程,而Spring WebFlux是基于非阻塞I/O,我们不再这里讨论这两个的区别,只需要知道WebFlux特别适合I/O密集型性应用,网关就是典型的I/O密集应用(网络I/O处理频繁)。SpringSecurity对WebFlux提供了支持,在WebFlux中WebFilter组件承担着与Filter相似的功能。

我们在先前的应用中通过HttpSecurity组件来组装SpringSecurity功能,在这里要使用新的组件ServerHttpSecurity来组装SpringSecurity功能,配置如下所示:

///启用WebFlux下的SpringSecurity配置
@EnableWebFluxSecurity
public class ResourceServerConfig {
    //// 访问权限验证
    @Autowired
    AuthManagerHandler authManagerHandler;
    //// 无权限访问处理器
    @Autowired
    AccessDeniedHandler accessDeniedHandler;
    /// 登录信息失效处理器
    @Autowired
    LoginLoseHandler loginLoseHandler;
    ////访问白名单,对白名单路径可以实现匿名访问
    @Autowired
    private WhiteUrlProperties whiteUrlProperties;

    @Bean
    public SecurityWebFilterChain springSecurityFilterChain(ServerHttpSecurity http) {
        http.oauth2ResourceServer()
                /// 这里配置对令牌的校验,从OAuth2授权服务中获取令牌对应的授权信息
                .opaqueToken()
                ///令牌校验地址,用于校验令牌是否有效,已经令牌对应的授权信息
                .introspectionUri("http://localhost:8081/oauth/check_token")
                //// 客户端信息
                .introspectionClientCredentials("blog", "blog")
                .and()
                .accessDeniedHandler(accessDeniedHandler)
                .authenticationEntryPoint(loginLoseHandler)
                .and().authorizeExchange()
                .pathMatchers(HttpMethod.OPTIONS).permitAll() //o
                .pathMatchers("/**").access(authManagerHandler)
                .anyExchange().authenticated()
                .and()
                .addFilterBefore(securityGlobalFilter(whiteUrlProperties), SecurityWebFiltersOrder.FIRST)
                .cors().disable().csrf().disable();
        return http.build();
    }
    /// 该过滤器实现将获取到的授权信息转发到下游服务中,方便后续校验
    public WebFilter securityGlobalFilter(WhiteUrlProperties properties) {
        return new SecurityGlobalFilter(properties);
    }

}

网关路由配置以及其它细节信息可以前往代码仓库进行查看,在此不做过多解释。

资源服务器

资源服务器也需要做一些调整,不需要对请求进行严格的访问控制,只需要校验网关传递的授权信息即可,然后将授权信息放入到SecurityContext中方便后续处理,同时需要注意在资源服务中还是使用Spring MVC框架进行处理(Spring WebFlux可以提高系统吞吐量,但是也会增加编程难度,例如原先的线程变量将不适用,具体需要考量整体编程人员掌握的技术栈来做决定)。

这里的资源服务器不再依赖OAuth授权服务,因此可以移除@EnableResourceServer配置(不直接参与权限控制,只需要校验上游传递的授权信息是否有效即可),同时增加对上游SpringCloud Gateway传递的授权信息进行解析处理,增加自定义SecurityAuthTokenFilter组件:

public class SecurityAuthTokenFilter extends OncePerRequestFilter {
    private static final String AUTH_TOKEN_NAME = "token";
    @Override
    protected void doFilterInternal(HttpServletRequest request, HttpServletResponse response,
            FilterChain filterChain) throws ServletException, IOException {
        String token = request.getHeader(AUTH_TOKEN_NAME);
        if (StringUtils.isEmpty(token)) {
            /// 继续处理
            filterChain.doFilter(request, response);
            return;
        }
        ///....省略处理细节,具体前往代码仓库进行查看
        //// 创建自定义的Authentication对象,必须申明为已授权,也就是isAuthenticated()方法返回为true
        BlogAuthentication authentication = new BlogAuthentication(userId, clientId, authorities);
        //....省略处理细节,具体前往代码仓库进行查看
        /// 将授权信息放入到SecurityContext中,方便后续使用
        SecurityContextHolder.getContext().setAuthentication(authentication);
        filterChain.doFilter(request, response);
    }
}

运行验证

分别运行Gateway网关服务、OAuth授权服务、Resource资源服务

  • Gateway 8080端口
  • OAuth 8081端口
  • Resource 8082端口

授权登录

使用密码模式进行授权登录,发送请求POST http://localhost:8080/blog-oauth/oauth/token,请求参数:

client_id:blog
client_secret:blog
grant_type:password
username:admin
password:admin

返回结果:

{
    "access_token": "4aace702-cc9d-4a92-b507-9b65f192a65f",
    "token_type": "bearer",
    "refresh_token": "a50a6cff-97b0-4d0f-b3d2-e0fdcee6f142",
    "expires_in": 5591,
    "scope": "all user"
}

资源访问

使用得到的access_token访问资源服务器中的/admin/hello接口,发送请求GET http://localhost:8080/blog-resource/admin/hello,请求头中携带参数:

Authorization:Bearer 4aace702-cc9d-4a92-b507-9b65f192a65f

返回结果:

{
    "code": 200,
    "data": "Hello Admin"
}
访问其他权限的接口,发送请求GET http://localhost:8080/blog-resource/user/hello,请求头中携带参数:
Authorization:Bearer 4aace702-cc9d-4a92-b507-9b65f192a65f

返回结果:

{
    "code": 400,
    "message": "无权限访问"
}

至此得到期望的访问结果,实现了统一权限控制

总结

  • SpringCloud Gateway使用WebFlux技术进行开发
  • SpringSecurity提供了@EnableWebFluxSecurity来支持WebFlux
  • SpringSecurity使用ReactiveSecurityContextHolder.getContext()来实现SecurityContextHolder功能

参考文档

七、SpringSecurity OAuth2 + JWT + SpringCloud Gateway实现统一鉴权管理

代码

代码仓库:地址

代码分支: lesson7

简介

在上一篇文章中,我们使用SpringSecurity OAuth2 + SpringCloud Gateway搭建了一套符合微服务架构的授权系统,在Gateway网关实现统一身份鉴定、访问权限控制,同时将授权信息下发到下游业务服务中,下游业务服务只需要关注核心业务逻辑。上述架构依赖于auth授权服务器,每一次业务请求都需要使用access_token请求auth授权服务器来获取用户授权信息,如果access_token自带授权信息,那么网关只需要鉴别access_token有效信息,这将会降低系统对auth授权服务器的依赖,JWT(JSON Web Token)将是很好的选择。

JWT

我们这里不详细介绍JWT,有兴趣的同学可以查看阮一峰老师的文章:JWT入门教程。JWT定义了一种数据结构,它由三部分组成:

  • Header,头部,定义了签名算法,令牌类型
  • Payload,负载,是一个JSON对象,包含实际应用中使用的数据,例如用户名,用户角色,注意这部分内容是不加密的,因此不能包含保密信息
  • Signature,签名,用于验证JWT是否有效,防止信息内容篡改

JWT内容是不加密的,可以使用在线工具解码信息,查看内容。例如有一个JWT格式Token:

eyJhbGciOiJSUzI1NiIsInR5cCI6IkpXVCJ9.eyJhdWQiOlsicmVzb3VyY2UiLCJibG9nIl0sImV4X3VzZXJuYW1lIjoiYWRtaW4iLCJ1c2VyX25hbWUiOiJhZG1pbiIsInNjb3BlIjpbImFsbCIsInVzZXIiXSwiZXhwIjoxNjU0NTk0MTc2LCJhdXRob3JpdGllcyI6WyJhZG1pbiJdLCJqdGkiOiI2NTdjZmU0Yi05ZDBlLTRhNTUtYjJjOS1iZWE3MTA2YWJkYWIiLCJjbGllbnRfaWQiOiJibG9nIn0.CGQTlvCCwGWIuJBy_qNeX2YBEYYTy6W1FPXOll75P1jdEyvi_TDiTLE4AO2Fa9vtgdWKrtywgGi4kFWZw8mcRFmhVfl9ehdoPcN5Hmdnz-ybJuLWh0i1k0xqg6MsZryTR1wAweEggZkHsIdCZfOw-yPZFTKuhAgVL4d-12Uthb4

在线解析后得到信息如下:

img

项目改造优化

优化分析

在先前文章中,我们将授权信息保存在auth授权服务器中,客户端需要通过请求auth授权服务器来获取授权信息,如果使用JWT,并且在JWT中保存相关授权信息,那么可以直接解析JWT就可以获取授权信息(需要验证JWT是有有效)。

在先前的项目中,我们使用SpringSecurity OAuth2默认配置来创建access_token,实际上SpringSecurity OAuth2提供了对JWT格式access_token支持,我们需要更改access_token的生成方式,因此需要修改auth授权服务中的Token生成方式,同时需要对Gateway网关服务中的Token解析方式进行修改

生成JWT格式Token

使用SpringSecurity OAuth2时,如果没有配置TokenService对象,将会默认使用DefaultTokenServices组件来管理access_token, 使用UUID.randomUUID().toString()方式生成access_token和refresh_token,因此token中不包含任何信息。我们需要配置新的TokenService对象来生成JWT格式Token。

SpringSecurity OAuth2提供了以下组件来生成JWT格式Token:

  • JwtTokenStore实现了TokenStore接口,用来管理access_token和refresh_token
  • JwtAccessTokenConverter实现了TokenEnhancer, AccessTokenConverter接口,可以生成JWT格式token

JWT签名算法我们选用安全性更高的非对称加密算法:RSA(在代码auth/src/test/java/com/hzchendou/blog/demo/RSAKeyTest中提供生成RSA Key方法),配置TokenService:

@Bean
public TokenStore tokenStore() {
   return new JwtTokenStore(accessTokenConverter());
}
@Bean
public JwtAccessTokenConverter accessTokenConverter() {
  JwtAccessTokenConverter converter = new JwtAccessTokenConverter();
  converter.setKeyPair(keyPair()); //非对称秘钥,具体参见代码
  return converter;
}

@Bean
public AuthorizationServerTokenServices tokenService() {
   DefaultTokenServices service = new DefaultTokenServices();
   service.setClientDetailsService(authClientDetailService);
   service.setSupportRefreshToken(true);
   service.setTokenStore(tokenStore);
   //令牌增强
   TokenEnhancerChain tokenEnhancerChain = new TokenEnhancerChain();
   List<TokenEnhancer> tokenEnhancers = new ArrayList<>();
   tokenEnhancers.add(tokenEnhancer);
   tokenEnhancers.add(accessTokenConverter);
   tokenEnhancerChain.setTokenEnhancers(tokenEnhancers);

   service.setTokenEnhancer(tokenEnhancerChain);
   service.setAccessTokenValiditySeconds(60 * 60 * 2); // 令牌默认有效期2小时
   service.setRefreshTokenValiditySeconds(60 * 60 * 24 * 3); // 刷新令牌默认有效期3天
   return service;
}

将TokenService配置到OAuth服务中:

endpoints.tokenServices(tokenService());//令牌管理服务

解析JWT格式Token

在网关中需要配置JWT格式解析器,使用JwtAuthenticationConverter来解析JWT中的字段:

@Bean
public Converter<Jwt, ? extends Mono<? extends AbstractAuthenticationToken>> jwtAuthenticationConverter() {
  JwtGrantedAuthoritiesConverter jwtGrantedAuthoritiesConverter = new JwtGrantedAuthoritiesConverter();
  jwtGrantedAuthoritiesConverter.setAuthorityPrefix("");
  jwtGrantedAuthoritiesConverter.setAuthoritiesClaimName("authorities");
  JwtAuthenticationConverter jwtAuthenticationConverter = new JwtAuthenticationConverter();
  jwtAuthenticationConverter.setJwtGrantedAuthoritiesConverter(jwtGrantedAuthoritiesConverter);
  return new ReactiveJwtAuthenticationConverterAdapter(jwtAuthenticationConverter);
}

因为使用RSA签名算法,因此在Gateway中需要配置RSA公钥来验证Token有效性, 有多种方式可以配置JWT解析验证器来验证JWT的有效性:

方法一、SpringSecurity OAuth2提供的方式:

配置public key信息来验证JWT有效性,在配置文件中配置(配置获取公约的接口地址):

spring:
  security:
    oauth2:
      resourceserver:
        jwt:
          jwk-set-uri: http://localhost:8081/key/public-key /// 这里需要在oauth授权服务器中配置接口

@RestController
public class KeyController {
    @Autowired
    KeyPair keyPair;
    //获取公钥
    @GetMapping("/key/public-key")
    public Map<String, Object> getPublicKey() {
        RSAPublicKey publicKey = (RSAPublicKey) keyPair.getPublic();
        RSAKey key = new RSAKey.Builder(publicKey).build();
        return new JWKSet(key).toJSONObject();
    }
}

方法二、直接配置PublicKey方式(直接将公钥写入到配置文件中进行读取):

我们这里直接将public key 配置到gateway网关中:

jwt:
  rsa:
    publickey: MIGfMA0GCSqGSIb3DQEBAQUAA4GNADCBiQKBgQCEfoWyfxqYz6j6tczCoELJfCwxpC+iHox7YEvz6slxNworp+CQAC86qt4Rx14lijoufiBMol0/mAABlG1lv3K1LOgQGcwueZDY5nk0uabOWv787moVbQHRTQoAwMIeSDPQ3SgSoEFyHM6Jj/We7XUpAyQEXKk9AabAvywEk2u9ewIDAQAB

然后手动创建JwtDecoder:

@Slf4j
@Configuration
public class TokenConfig {
    @Value("${jwt.rsa.publickey}")
    private String publicKey;
    public RSAPublicKey rsaPublicKey() {
        try {
            return (RSAPublicKey)RSAUtils.decodePublicKey(publicKey);
        } catch (Exception ex) {
            log.error("生成 KeyPair 失败", ex);
            System.exit(-1);
            return null;
        }
    }
    @Bean
    public NimbusReactiveJwtDecoder nimbusReactiveJwtDecoder() {
        return NimbusReactiveJwtDecoder.withPublicKey(rsaPublicKey())
                .signatureAlgorithm(SignatureAlgorithm.from("RS256")).build();
    }
}

还有一步需要配置,JWT Token解析后的类型是JwtAuthenticationToken,因此需要修改SecurityGlobalFilter中ReactiveSecurityContextHolder.getContext()方法返回的authentication类型(具体参见代码)

运行校验

分别运行auth、resource、gateway服务:

  • gateway - 8080
  • auth - 8081
  • resource - 8082

发起OAuth2密码授权请求: POST http://localhost:8080/blog-oauth/oauth/token,请求参数:

client_id:blog
client_secret:blog
grant_type:password
username:admin
password:admin

请求结果:

{
    "access_token": "eyJhbGciOiJSUzI1NiIsInR5cCI6IkpXVCJ9.eyJhdWQiOlsicmVzb3VyY2UiLCJibG9nIl0sImV4X3VzZXJuYW1lIjoiYWRtaW4iLCJ1c2VyX25hbWUiOiJhZG1pbiIsInNjb3BlIjpbImFsbCIsInVzZXIiXSwiZXhwIjoxNjU0NTk4MzY3LCJhdXRob3JpdGllcyI6WyJhZG1pbiJdLCJqdGkiOiIyMmVlOTU4My00Y2U5LTRmNzEtOGI4MS02YjViNzdmYWNlYWEiLCJjbGllbnRfaWQiOiJibG9nIn0.atWwzwpCK1ycjf3-EkPUYs4DMqO7rGPIMwMjHKS3FrTKRjMW5DHkQjtilG2EB8qGNBlwQJo0xAnQ_RNMzOjVojGxyb-TUPCubqODnmnYhuee0ho2TurDT5YzfO-Ypkv2SDqEm6Kw38m-oV_93NofGtKNJD1or2kwdoZe6kn4qgw",
    "token_type": "bearer",
    "refresh_token": "eyJhbGciOiJSUzI1NiIsInR5cCI6IkpXVCJ9.eyJhdWQiOlsicmVzb3VyY2UiLCJibG9nIl0sImV4X3VzZXJuYW1lIjoiYWRtaW4iLCJ1c2VyX25hbWUiOiJhZG1pbiIsInNjb3BlIjpbImFsbCIsInVzZXIiXSwiYXRpIjoiMjJlZTk1ODMtNGNlOS00ZjcxLThiODEtNmI1Yjc3ZmFjZWFhIiwiZXhwIjoxNjU0ODUwMzY3LCJhdXRob3JpdGllcyI6WyJhZG1pbiJdLCJqdGkiOiJjMzZhNDQ2Mi00ZmE3LTQ2OGUtODNiMS1iYzk4MDFjMzBjMWEiLCJjbGllbnRfaWQiOiJibG9nIn0.IyBeBQMjU-KYGIvvlQTrTkEtrPmTjLZIl1oFvyK0vytOlOFaE4Q5tMOLf1lt1UaBpmi2Tz4ElQSc6EMYX_OKmbyEHSidYxseUr8gE5MVM1raqOPCnR0Dyn7okQ0NvArOB9JuxLTXSa3NoSM3OxRQm2sUS55e6FKpifZ2q7xgGnY",
    "expires_in": 7199,
    "scope": "all user",
    "ex_username": "admin",
    "jti": "22ee9583-4ce9-4f71-8b81-6b5b77faceaa"
}

发起资源服务器请求:POST http://localhost:8080/blog-resource/admin/hello,请求头携带token参数:

Authorization:Bearer eyJhbGciOiJSUzI1NiIsInR5cCI6IkpXVCJ9.eyJhdWQiOlsicmVzb3VyY2UiLCJibG9nIl0sImV4X3VzZXJuYW1lIjoiYWRtaW4iLCJ1c2VyX25hbWUiOiJhZG1pbiIsInNjb3BlIjpbImFsbCIsInVzZXIiXSwiZXhwIjoxNjU0NTk2OTQyLCJhdXRob3JpdGllcyI6WyJhZG1pbiJdLCJqdGkiOiJkOWRhZjhmYS1jOTY4LTQ2YzQtODMyMi1kOTQxNGU3YWZhY2UiLCJjbGllbnRfaWQiOiJibG9nIn0.Jkz_Tlk1W7opXspihyDkp1VFcIXu3ebfVPkshjFcKpktPmqkUIA4D2aWF5A13fq5QGUIDQKf89rVeGHaFfer657J7kqaax2qNT6yuNgmQAu4C8VQkG01VLDsOa-m9xaZnqR_--Af-Z7FbwpZNOT2pBuyP4M3efnMmGhRQQjB4hQ

返回结果:

{
    "code": 200,
    "data": "Hello Admin"
}

结果符合预期,到此完成JWT + SpringSecurity OAuth2 + SpringCloud Gateway 统一权限访问控制功能

总结

  • JWT自带用户信息,只需要验证Token有效性
  • OAuth授权服务添加JWT TokenService返回JWT格式token
  • Gateway网关服务添加JWTAuthenticationConverter解析JWT信息,同时添加NimbusReactiveJwtDecoder对JWT有效性进行验证

参考文档

Views: 403