通信服务
邮件服务
Spring Mail 集成
Spring Boot 通过 spring-boot-starter-mail 提供对 JavaMailSender 的自动配置,底层封装了 JavaMail (Jakarta Mail)。
引入依赖
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-mail</artifactId>
</dependency>配置邮件客户端
以阿里云邮件推送服务为例:
spring:
mail:
host: smtp.aliyun.com
port: 465
username: noreply@example.com
password: your-smtp-password
protocol: smtps
default-encoding: UTF-8
properties:
mail:
smtp:
auth: true
socketFactory:
class: javax.net.ssl.SSLSocketFactory
ssl:
enable: true
starttls:
enable: false
timeout: 10000
connectiontimeout: 10000
writetimeout: 10000常用 SMTP 服务器地址:
| 服务商 | SMTP 地址 | SSL 端口 | 备注 |
|---|---|---|---|
| 阿里云邮件推送 | smtp.aliyun.com | 465 | 需开启 SMTP 密码 |
| 腾讯企业邮 | smtp.exmail.qq.com | 465 | 需客户端专用密码 |
| QQ 邮箱 | smtp.qq.com | 465 | 需授权码 |
| 163 邮箱 | smtp.163.com | 465 | 需开启 SMTP 服务 |
| Gmail | smtp.gmail.com | 465 | 需 App Passwords |
| AWS SES | email-smtp.region.amazonaws.com | 465 | 需 SMTP 凭证 |
注入 JavaMailSender
@Service
public class MailService {
@Autowired
private JavaMailSender mailSender;
@Value("${spring.mail.username}")
private String from;
}邮件类型对比
简单文本邮件
public void sendSimpleMail(String to, String subject, String content) {
SimpleMailMessage message = new SimpleMailMessage();
message.setFrom(from);
message.setTo(to);
message.setSubject(subject);
message.setText(content);
mailSender.send(message);
}MIME 富文本邮件(HTML)
public void sendHtmlMail(String to, String subject, String htmlContent) throws MessagingException {
MimeMessage message = mailSender.createMimeMessage();
MimeMessageHelper helper = new MimeMessageHelper(message, true, "UTF-8");
helper.setFrom(from);
helper.setTo(to);
helper.setSubject(subject);
helper.setText(htmlContent, true); // true 表示启用 HTML
mailSender.send(message);
}HTML 模板邮件(Thymeleaf)
结合 Thymeleaf 模板引擎渲染邮件正文,支持变量注入、条件判断、循环渲染。
依赖:
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-thymeleaf</artifactId>
</dependency>模板文件 templates/mail/welcome.html:
<!DOCTYPE html>
<html xmlns:th="http://www.thymeleaf.org">
<head>
<meta charset="UTF-8" />
<title>欢迎邮件</title>
<style>
.container { max-width: 600px; margin: 0 auto; font-family: Arial, sans-serif; }
.header { background: #1890ff; color: #fff; padding: 20px; text-align: center; }
.content { padding: 20px; line-height: 1.8; }
.footer { background: #f5f5f5; padding: 10px; text-align: center; font-size: 12px; color: #999; }
</style>
</head>
<body>
<div class="container">
<div class="header">
<h1 th:text="${title}">欢迎加入</h1>
</div>
<div class="content">
<p th:text="'尊敬的 ' + ${username} + ':'">尊敬的 用户:</p>
<p th:text="${message}">欢迎信息</p>
<a th:href="${actionUrl}" style="display:inline-block;padding:10px 30px;background:#1890ff;color:#fff;text-decoration:none;border-radius:4px;">
<span th:text="${actionText}">立即验证</span>
</a>
</div>
<div class="footer">
<p>如果您未注册,请忽略此邮件</p>
</div>
</div>
</body>
</html>Java 代码渲染模板:
@Service
public class ThymeleafMailService {
@Autowired
private JavaMailSender mailSender;
@Autowired
private SpringTemplateEngine templateEngine;
@Value("${spring.mail.username}")
private String from;
public void sendTemplateMail(String to, String username, String verifyUrl) throws MessagingException {
Context context = new Context();
context.setVariable("title", "邮箱验证");
context.setVariable("username", username);
context.setVariable("message", "请点击下方按钮完成邮箱验证,链接有效期为 24 小时。");
context.setVariable("actionUrl", verifyUrl);
context.setVariable("actionText", "验证邮箱");
String html = templateEngine.process("mail/welcome", context);
MimeMessage message = mailSender.createMimeMessage();
MimeMessageHelper helper = new MimeMessageHelper(message, true, "UTF-8");
helper.setFrom(from);
helper.setTo(to);
helper.setSubject("邮箱验证");
helper.setText(html, true);
mailSender.send(message);
}
}带附件邮件
public void sendAttachmentMail(String to, String subject, String text, String attachmentPath) throws MessagingException {
MimeMessage message = mailSender.createMimeMessage();
MimeMessageHelper helper = new MimeMessageHelper(message, true, "UTF-8");
helper.setFrom(from);
helper.setTo(to);
helper.setSubject(subject);
helper.setText(text, true);
FileSystemResource file = new FileSystemResource(new File(attachmentPath));
helper.addAttachment(file.getFilename(), file);
mailSender.send(message);
}支持同时添加多个附件,也可以从 InputStream 添加:
helper.addAttachment("report.pdf", () -> new FileInputStream(reportPath));批量邮件 + 异步发送
批量邮件需要两个关键处理:多收件人和异步执行,避免阻塞主线程。
多收件人发送
public void sendBatchMail(String[] tos, String subject, String htmlContent) throws MessagingException {
MimeMessage message = mailSender.createMimeMessage();
MimeMessageHelper helper = new MimeMessageHelper(message, true, "UTF-8");
helper.setFrom(from);
helper.setTo(tos); // 收件人数组
helper.setSubject(subject);
helper.setText(htmlContent, true);
mailSender.send(message);
}注意:使用 setTo(String[]) 时所有收件人能看到彼此地址。如需隐藏,应使用密送:
helper.setTo("placeholder@example.com");
helper.setBcc(recipientList.toArray(new String[0]));异步发送
配置线程池:
@Configuration
public class MailAsyncConfig {
@Bean("mailTaskExecutor")
public Executor mailTaskExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(2);
executor.setMaxPoolSize(5);
executor.setQueueCapacity(100);
executor.setThreadNamePrefix("mail-sender-");
executor.setRejectedExecutionHandler(new CallerRunsPolicy());
executor.initialize();
return executor;
}
}异步发送服务:
@Service
public class AsyncMailService {
@Autowired
private JavaMailSender mailSender;
@Autowired
@Qualifier("mailTaskExecutor")
private Executor mailTaskExecutor;
@Value("${spring.mail.username}")
private String from;
public void sendAsync(String to, String subject, String htmlContent) {
mailTaskExecutor.execute(() -> {
try {
MimeMessage message = mailSender.createMimeMessage();
MimeMessageHelper helper = new MimeMessageHelper(message, true, "UTF-8");
helper.setFrom(from);
helper.setTo(to);
helper.setSubject(subject);
helper.setText(htmlContent, true);
mailSender.send(message);
log.info("邮件发送成功: {}", to);
} catch (Exception e) {
log.error("邮件发送失败: {}", to, e);
}
});
}
public void sendBatchAsync(List<String> toList, String subject, String htmlContent) {
toList.forEach(to -> sendAsync(to, subject, htmlContent));
}
}使用 @Async 注解(推荐方式):
@Slf4j
@Service
public class AsyncMailService {
@Autowired
private JavaMailSender mailSender;
@Value("${spring.mail.username}")
private String from;
@Async("mailTaskExecutor")
public void sendMailAsync(String to, String subject, String htmlContent) {
try {
MimeMessage message = mailSender.createMimeMessage();
MimeMessageHelper helper = new MimeMessageHelper(message, true, "UTF-8");
helper.setFrom(from);
helper.setTo(to);
helper.setSubject(subject);
helper.setText(htmlContent, true);
mailSender.send(message);
log.info("邮件发送成功: {}", to);
} catch (Exception e) {
log.error("邮件发送失败: {}", to, e);
}
}
}邮件队列 + 重试机制
高可靠场景下,邮件应写入消息队列,由独立消费者异步发送,并支持失败重试。
基于 Redis 的邮件队列
@Component
public class MailQueueProducer {
@Autowired
private StringRedisTemplate redisTemplate;
private static final String MAIL_QUEUE_KEY = "mail:queue";
private static final String MAIL_RETRY_KEY = "mail:retry";
public void enqueue(MailTask task) {
String json = JSON.toJSONString(task);
redisTemplate.opsForList().rightPush(MAIL_QUEUE_KEY, json);
}
}@Slf4j
@Component
public class MailQueueConsumer {
@Autowired
private StringRedisTemplate redisTemplate;
@Autowired
private JavaMailSender mailSender;
@Value("${spring.mail.username}")
private String from;
private static final String MAIL_QUEUE_KEY = "mail:queue";
private static final String MAIL_RETRY_KEY = "mail:retry";
private static final int MAX_RETRY = 3;
@Scheduled(fixedDelay = 1000)
public void consume() {
String json = redisTemplate.opsForList().leftPop(MAIL_QUEUE_KEY);
if (json == null) return;
MailTask task = JSON.parseObject(json, MailTask.class);
try {
MimeMessage message = mailSender.createMimeMessage();
MimeMessageHelper helper = new MimeMessageHelper(message, true, "UTF-8");
helper.setFrom(from);
helper.setTo(task.getTo());
helper.setSubject(task.getSubject());
helper.setText(task.getContent(), true);
mailSender.send(message);
log.info("邮件发送成功: {}", task.getTo());
} catch (Exception e) {
task.incrementRetry();
if (task.getRetryCount() < MAX_RETRY) {
// 延迟重试,使用 ZSet 按时间排序
redisTemplate.opsForZSet()
.add(MAIL_RETRY_KEY, json, System.currentTimeMillis() + 30000);
log.warn("邮件发送失败,加入重试队列: {}", task.getTo());
} else {
log.error("邮件发送失败,已达最大重试次数: {}", task.getTo());
// 写入死信队列,人工介入
redisTemplate.opsForList().rightPush("mail:dead", json);
}
}
}
@Scheduled(fixedDelay = 5000)
public void retry() {
Set<String> tasks = redisTemplate.opsForZSet()
.rangeByScore(MAIL_RETRY_KEY, 0, System.currentTimeMillis());
if (tasks == null || tasks.isEmpty()) return;
tasks.forEach(json -> {
redisTemplate.opsForZSet().remove(MAIL_RETRY_KEY, json);
redisTemplate.opsForList().rightPush(MAIL_QUEUE_KEY, json);
});
}
}@Data
public class MailTask {
private String to;
private String subject;
private String content;
private int retryCount;
public void incrementRetry() {
this.retryCount++;
}
}基于 RabbitMQ 的邮件队列
@Configuration
public class MailRabbitConfig {
public static final String EXCHANGE = "mail.exchange";
public static final String QUEUE = "mail.queue";
public static final String ROUTING_KEY = "mail.send";
public static final String DLX_EXCHANGE = "mail.dlx";
public static final String DLX_QUEUE = "mail.dlx.queue";
@Bean
public Queue mailQueue() {
return QueueBuilder.durable(QUEUE)
.withArgument("x-dead-letter-exchange", DLX_EXCHANGE)
.withArgument("x-dead-letter-routing-key", "mail.dead")
.withArgument("x-message-ttl", 60000)
.build();
}
@Bean
public DirectExchange mailExchange() {
return new DirectExchange(EXCHANGE);
}
@Bean
public Binding mailBinding() {
return BindingBuilder.bind(mailQueue()).to(mailExchange()).with(ROUTING_KEY);
}
// 死信队列
@Bean
public Queue dlxQueue() {
return QueueBuilder.durable(DLX_QUEUE).build();
}
@Bean
public DirectExchange dlxExchange() {
return new DirectExchange(DLX_EXCHANGE);
}
@Bean
public Binding dlxBinding() {
return BindingBuilder.bind(dlxQueue()).to(dlxExchange()).with("mail.dead");
}
}@RabbitListener(queues = "mail.queue")
public void handleMailTask(MailTask task, Channel channel, Message message) throws IOException {
try {
MimeMessage mimeMsg = mailSender.createMimeMessage();
MimeMessageHelper helper = new MimeMessageHelper(mimeMsg, true, "UTF-8");
helper.setFrom(from);
helper.setTo(task.getTo());
helper.setSubject(task.getSubject());
helper.setText(task.getContent(), true);
mailSender.send(mimeMsg);
channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
log.info("邮件发送成功: {}", task.getTo());
} catch (Exception e) {
log.error("邮件发送失败: {}", task.getTo(), e);
// 重试指定次数后进入死信队列
channel.basicReject(message.getMessageProperties().getDeliveryTag(), false);
}
}邮件发送常见问题
被识别为垃圾邮件
| 原因 | 解决方案 |
|---|---|
| IP 信誉度低 | 使用专业邮件推送服务(阿里云邮件推送、SendGrid、AWS SES) |
| 内容含敏感词 | 避免全大写标题、过多感叹号、营销类词汇 |
| 缺乏 SPF 记录 | 在 DNS 中添加 SPF TXT 记录 |
| 缺乏 DKIM 签名 | 配置 DKIM 密钥,对邮件头部签名 |
| 缺乏 DMARC 策略 | 配置 DMARC 记录,指定处理策略 |
| 发送频率过高 | 控制发送速率,设置合理的 QPS |
| 收件人地址质量差 | 清理无效地址,支持退订链接 |
SPF 配置
SPF (Sender Policy Framework) 在 DNS 中声明哪些 IP 允许发送该域名的邮件。
TXT 记录: v=spf1 include:spf.aliyun.com ~allinclude:spf.aliyun.com:允许阿里云邮件服务器代发~all:软失败(标记但不受理)-all:硬失败(建议拒收)?all:中立
DKIM 配置
DKIM (DomainKeys Identified Mail) 使用数字签名验证邮件发件人。
TXT 记录: k=rsa; p=MIGfMA0GCSqGSIb3DQEBAQUAA4...配置步骤:
- 在邮件推送服务商获取 DKIM 私钥(服务商管理)和公钥(用户配置 DNS)
- 在 DNS 中添加
{selector}._domainkey.example.com的 TXT 记录 - 等待 DNS 生效后,服务商会自动对发出的邮件进行 DKIM 签名
DMARC 配置
DMARC (Domain-based Message Authentication, Reporting and Conformance) 基于 SPF 和 DKIM 的结果制定处理策略。
TXT 记录 (_dmarc.example.com):
v=DMARC1; p=quarantine; rua=mailto:dmarc-reports@example.com; pct=100; sp=quarantine参数说明:
p:策略,none(仅监控)、quarantine(标记为垃圾邮件)、reject(直接拒收)rua:聚合报告接收邮箱ruf:失败报告接收邮箱pct:采样百分比sp:子域名策略
建议上线流程:先 p=none 监控一段时间,确认无误后再提升为 p=quarantine 或 p=reject。
其他常见问题
- 535 Authentication Failed:用户名密码错误,或未开启 SMTP 服务
- 550 Mailbox not found:收件人地址不存在
- Timeout 异常:增大
timeout/connectiontimeout配置,或检查网络连通性 - Connection refused:端口错误或 SMTP 服务器 IP 被限制
- 附件过大:建议不超过 10MB,大文件可先上传对象存储再提供下载链接
# 推荐的超时配置
spring:
mail:
properties:
mail:
smtp:
timeout: 30000
connectiontimeout: 15000
writetimeout: 30000对象存储
云服务商对比与选型
| 特性 | 阿里云 OSS | 腾讯云 COS | AWS S3 |
|---|---|---|---|
| 存储类型 | 标准 / 低频 / 归档 / 冷归档 | 标准 / 低频 / 归档 / 深度归档 | Standard / IA / Glacier / Deep Archive |
| SDK 成熟度 | ★★★★★ | ★★★★ | ★★★★★ |
| 国内访问速度 | 极佳 | 极佳 | 一般 |
| 全球节点 | 丰富 | 较丰富 | 最丰富 |
| 计费模式 | 按量 / 资源包 | 按量 / 资源包 | 按量 / 预留容量 |
| CDN 集成 | 阿里云 CDN | 腾讯云 CDN | CloudFront |
| S3 兼容 API | 支持 | 支持 | 原生 |
| 外网下行流量费 | 较高 | 中等 | 较高 |
| 合规认证 | 等保三级等 | 等保三级等 | SOC2/ISO27001 等 |
选型建议:
- 国内业务:首选阿里云 OSS,SDK 最完善、节点最丰富;如已使用腾讯云生态(CDB/CLB 等),可优先 COS
- 海外业务:AWS S3 是行业标准,SDK 和工具链极其成熟
- 多云容灾:各厂商均支持 S3 兼容 API,可封装统一接口层实现无缝切换
- 成本敏感:低频/归档存储成本显著低于标准存储,适合备份、日志场景
Spring Boot 集成 SDK 配置
以阿里云 OSS 为例(各厂商 SDK 配置模式基本一致):
引入依赖
<!-- 阿里云 OSS -->
<dependency>
<groupId>com.aliyun.oss</groupId>
<artifactId>aliyun-sdk-oss</artifactId>
<version>3.17.4</version>
</dependency>
<!-- AWS S3 -->
<dependency>
<groupId>software.amazon.awssdk</groupId>
<artifactId>s3</artifactId>
<version>2.25.0</version>
</dependency>
<!-- 腾讯云 COS -->
<dependency>
<groupId>com.qcloud</groupId>
<artifactId>cos_api</artifactId>
<version>5.6.155</version>
</dependency>配置属性
oss:
aliyun:
endpoint: oss-cn-hangzhou.aliyuncs.com
access-key-id: your-access-key-id
access-key-secret: your-access-key-secret
bucket-name: your-bucket
domain: https://your-bucket.oss-cn-hangzhou.aliyuncs.com
aws:
region: ap-northeast-1
access-key-id: your-access-key-id
secret-access-key: your-secret-access-key
bucket: your-bucket
domain: https://your-bucket.s3.ap-northeast-1.amazonaws.com
tencent:
region: ap-guangzhou
secret-id: your-secret-id
secret-key: your-secret-key
bucket: your-bucket-1250000000
domain: https://your-bucket.cos.ap-guangzhou.myqcloud.com配置类
阿里云 OSS 客户端配置:
@Configuration
@ConfigurationProperties(prefix = "oss.aliyun")
@Data
public class AliyunOssConfig {
private String endpoint;
private String accessKeyId;
private String accessKeySecret;
private String bucketName;
private String domain;
@Bean
public OSS ossClient() {
return new OSSClientBuilder()
.build(endpoint, accessKeyId, accessKeySecret);
}
}AWS S3 客户端配置:
@Configuration
@ConfigurationProperties(prefix = "oss.aws")
@Data
public class AwsS3Config {
private String region;
private String accessKeyId;
private String secretAccessKey;
private String bucket;
private String domain;
@Bean
public S3Client s3Client() {
return S3Client.builder()
.region(Region.of(region))
.credentialsProvider(StaticCredentialsProvider.create(
AwsBasicCredentials.create(accessKeyId, secretAccessKey)))
.build();
}
}文件上传
普通上传
@Service
public class OssService {
@Autowired
private OSS ossClient;
@Value("${oss.aliyun.bucket-name}")
private String bucketName;
/**
* 上传文件流
*/
public String upload(String objectKey, InputStream inputStream) {
ossClient.putObject(bucketName, objectKey, inputStream);
return getObjectUrl(objectKey);
}
/**
* 上传本地文件
*/
public String uploadFile(String objectKey, File file) {
PutObjectRequest request = new PutObjectRequest(bucketName, objectKey, file);
ossClient.putObject(request);
return getObjectUrl(objectKey);
}
/**
* 上传字节数组
*/
public String uploadBytes(String objectKey, byte[] bytes) {
ossClient.putObject(bucketName, objectKey, new ByteArrayInputStream(bytes));
return getObjectUrl(objectKey);
}
private String getObjectUrl(String objectKey) {
return domain + "/" + objectKey;
}
}断点续传
适用于上传中断后可从中断点继续,避免重新上传已完成的部分。
public String resumeUpload(String objectKey, File file) {
UploadFileRequest request = new UploadFileRequest(bucketName, objectKey);
request.setUploadFile(file.getAbsolutePath());
// 分片大小(字节),推荐 1MB
request.setPartSize(1 * 1024 * 1024);
// 并发线程数
request.setTaskNum(5);
// 开启断点续传
request.setEnableCheckpoint(true);
// 断点记录文件存放目录
request.setCheckpointFile("upload-checkpoint/" + objectKey + ".cp");
UploadFileResult result = ossClient.uploadFile(request);
return getObjectUrl(objectKey);
}分片上传(原生 API)
手动控制分片上传流程,适合大文件或自定义并发策略。
public String multipartUpload(String objectKey, File file) {
// 1. 初始化分片上传
InitiateMultipartUploadRequest initRequest = new InitiateMultipartUploadRequest(bucketName, objectKey);
InitiateMultipartUploadResult initResult = ossClient.initiateMultipartUpload(initRequest);
String uploadId = initResult.getUploadId();
// 2. 分片上传
List<PartETag> partETags = new ArrayList<>();
long partSize = 5 * 1024 * 1024L; // 每个分片 5MB
long fileLength = file.length();
int partCount = (int) Math.ceil((double) fileLength / partSize);
try (FileInputStream fis = new FileInputStream(file)) {
for (int i = 0; i < partCount; i++) {
long skipBytes = (long) i * partSize;
fis.skip(skipBytes - (i > 0 ? partSize : 0));
long size = Math.min(partSize, fileLength - skipBytes);
byte[] data = new byte[(int) size];
fis.read(data);
UploadPartRequest uploadRequest = new UploadPartRequest();
uploadRequest.setBucketName(bucketName);
uploadRequest.setKey(objectKey);
uploadRequest.setUploadId(uploadId);
uploadRequest.setPartNumber(i + 1);
uploadRequest.setInputStream(new ByteArrayInputStream(data));
uploadRequest.setPartSize(size);
UploadPartResult result = ossClient.uploadPart(uploadRequest);
partETags.add(result.getPartETag());
}
} catch (IOException e) {
// 上传失败,取消分片上传
ossClient.abortMultipartUpload(new AbortMultipartUploadRequest(bucketName, objectKey, uploadId));
throw new RuntimeException("分片上传失败", e);
}
// 3. 完成分片上传
CompleteMultipartUploadRequest completeRequest = new CompleteMultipartUploadRequest(
bucketName, objectKey, uploadId, partETags);
ossClient.completeMultipartUpload(completeRequest);
return getObjectUrl(objectKey);
}文件下载
流式下载
public void download(String objectKey, OutputStream outputStream) {
OSSObject ossObject = ossClient.getObject(bucketName, objectKey);
try (InputStream in = ossObject.getObjectContent()) {
IOUtils.copy(in, outputStream);
} catch (IOException e) {
throw new RuntimeException("文件下载失败", e);
}
}
public byte[] downloadAsBytes(String objectKey) {
OSSObject ossObject = ossClient.getObject(bucketName, objectKey);
try (InputStream in = ossObject.getObjectContent()) {
return IOUtils.toByteArray(in);
} catch (IOException e) {
throw new RuntimeException("文件下载失败", e);
}
}签名 URL 临时授权
生成有时效性的下载链接,适用于临时授权给第三方下载私有文件。
/**
* 生成签名 URL
* @param objectKey 对象键
* @param expiration 过期时间(秒)
*/
public String generatePresignedUrl(String objectKey, int expirationSeconds) {
Date expiration = new Date(System.currentTimeMillis() + expirationSeconds * 1000L);
GeneratePresignedUrlRequest request = new GeneratePresignedUrlRequest(bucketName, objectKey);
request.setExpiration(expiration);
request.setMethod(HttpMethod.GET);
URL url = ossClient.generatePresignedUrl(request);
return url.toString();
}
/**
* 生成带自定义文件名的下载签名 URL
*/
public String generatePresignedDownloadUrl(String objectKey, String fileName, int expirationSeconds) {
Date expiration = new Date(System.currentTimeMillis() + expirationSeconds * 1000L);
GeneratePresignedUrlRequest request = new GeneratePresignedUrlRequest(bucketName, objectKey);
request.setExpiration(expiration);
request.setMethod(HttpMethod.GET);
// 设置响应头,强制下载并指定文件名
ResponseHeaderOverrides headers = new ResponseHeaderOverrides();
headers.setContentDisposition("attachment; filename=\"" + fileName + "\"");
request.setResponseHeaders(headers);
URL url = ossClient.generatePresignedUrl(request);
return url.toString();
}AWS S3 签名 URL 示例:
public String generateS3PresignedUrl(String key, int expirationSeconds) {
GetObjectRequest getObjectRequest = GetObjectRequest.builder()
.bucket(bucket)
.key(key)
.build();
LocalDateTime expiration = LocalDateTime.now().plusSeconds(expirationSeconds);
PresignedGetObjectRequest presignedRequest = s3Presigner.presignGetObject(r -> r
.signatureDuration(Duration.ofSeconds(expirationSeconds))
.getObjectRequest(getObjectRequest));
return presignedRequest.url().toString();
}防盗链
Referer 白名单
# OSS 控制台设置
防盗链:
白名单:
- http://www.example.com
- https://www.example.com
- http://*.example.com
是否允许空 Referer: false通过 SDK 配置 Bucket 防盗链:
public void setReferer(List<String> referers, boolean allowEmpty) {
SetBucketRefererRequest request = new SetBucketRefererRequest(bucketName);
request.setRefererList(new ArrayList<>(referers));
request.setAllowEmptyReferer(allowEmpty);
ossClient.setBucketReferer(request);
}CNAME + 自定义域名
将 OSS Bucket 绑定自定义域名,通过 CDN 加速并隐藏真实 OSS 地址。
OSS 控制台:
Bucket 配置:
域名管理:
- 绑定自定义域名: files.example.com
- CNAME 记录: files.example.com → your-bucket.oss-cn-hangzhou.aliyuncs.com
- 开启 CDN 加速: 阿里云 CDN / 腾讯云 CDN / CloudFront私有读写 + 签名鉴权
将 Bucket 设为私有读写,所有访问必须通过签名 URL 或 Authorization Header。
Bucket ACL:
权限: 私有(private)
读写权限: 只有 Bucket Owner 可读写
访问方式:
- 下载: Pre-signed URL(临时签名)
- 上传: Pre-signed URL 或 Authorization Header
- 前端直传: PostObject + 签名策略(form 表单)前端直传(服务端签名,客户端直传 OSS)示例:
/**
* 生成前端直传签名
*/
public Map<String, String> generateUploadSignature(String dir) {
// 过期时间
Date expiration = new Date(System.currentTimeMillis() + 3600 * 1000L);
PolicyConditions conditions = new PolicyConditions();
conditions.addConditionItem("bucket", bucketName);
conditions.addConditionItem(PolicyConditions.COND_CONTENT_LENGTH_RANGE, 0, 1048576000);
conditions.addConditionItem(PolicyConditions.COND_STARTS_WITH, "key", dir);
String postPolicy = ossClient.generatePostPolicy(expiration, conditions);
byte[] binaryData = postPolicy.getBytes(StandardCharsets.UTF_8);
String encodedPolicy = BinaryUtil.toBase64String(binaryData);
String postSignature = ossClient.calculatePostSignature(postPolicy);
Map<String, String> resp = new HashMap<>();
resp.put("accessId", accessKeyId);
resp.put("policy", encodedPolicy);
resp.put("signature", postSignature);
resp.put("dir", dir);
resp.put("host", domain);
resp.put("expire", String.valueOf(expiration.getTime() / 1000));
return resp;
}CDN 加速配置
CDN 加速流程:
1. 在 CDN 控制台添加加速域名(如 cdn.example.com)
2. 源站类型选择 OSS Bucket 外网域名
3. CDN 控制台生成 CNAME(如 cdn.example.com.w.kunluncan.com)
4. 在 DNS 中将 cdn.example.com CNAME 到 CDN 分配的域名
5. OSS Bucket 开启 CDN 回源鉴权(CDN 回源时自动添加 Authorization Header)Spring Boot 中切换 CDN 域名:
oss:
aliyun:
# 内网访问(同区域 ECS 免流量费)
internal-endpoint: oss-cn-hangzhou-internal.aliyuncs.com
# 外网直接访问
external-endpoint: oss-cn-hangzhou.aliyuncs.com
# CDN 加速域名(推荐线上使用)
cdn-domain: https://cdn.example.com@Service
public class CdnOssService {
@Autowired
private OSS ossClient;
@Value("${oss.aliyun.cdn-domain}")
private String cdnDomain;
@Value("${oss.aliyun.bucket-name}")
private String bucketName;
/**
* 返回 CDN 加速 URL(有 CDN 时优先使用)
*/
public String getAccessUrl(String objectKey, boolean useCdn) {
if (useCdn) {
return cdnDomain + "/" + objectKey;
}
return domain + "/" + objectKey;
}
}大文件上传策略
| 文件大小 | 推荐策略 | 说明 |
|---|---|---|
| < 10MB | 普通上传 | 一次请求,简单直接 |
| 10MB ~ 100MB | 分片上传 | 并发上传分片,容错性好 |
| 100MB ~ 1GB | 断点续传(SDK) | SDK 内置 checkpoint,支持续传 |
| > 1GB | 分片上传 + 本地分片 | 先在客户端切分,逐片上传,最后合并 |
客户端分片上传策略:
public class LargeFileUploader {
private static final long CHUNK_SIZE = 5 * 1024 * 1024; // 5MB
private final OssService ossService;
private final Executor executor;
private final AtomicInteger successCount = new AtomicInteger(0);
private final int totalParts;
public LargeFileUploader(OssService ossService, long fileSize) {
this.ossService = ossService;
this.executor = Executors.newFixedThreadPool(5);
this.totalParts = (int) Math.ceil((double) fileSize / CHUNK_SIZE);
}
public CompletableFuture<Void> uploadInChunks(String objectKey, File file) {
// 初始化分片上传
String uploadId = ossService.initMultipartUpload(objectKey);
List<CompletableFuture<PartETag>> futures = new ArrayList<>();
for (int i = 0; i < totalParts; i++) {
int partNumber = i + 1;
long startPos = (long) i * CHUNK_SIZE;
long partSize = Math.min(CHUNK_SIZE, file.length() - startPos);
CompletableFuture<PartETag> future = CompletableFuture.supplyAsync(() -> {
try (RandomAccessFile raf = new RandomAccessFile(file, "r")) {
raf.seek(startPos);
byte[] data = new byte[(int) partSize];
raf.readFully(data);
return ossService.uploadPart(objectKey, uploadId, partNumber, data);
} catch (IOException e) {
throw new RuntimeException("分片上传失败", e);
}
}, executor);
futures.add(future);
}
return CompletableFuture.allOf(futures.toArray(new CompletableFuture[0]))
.thenApply(v -> {
List<PartETag> etags = futures.stream()
.map(CompletableFuture::join)
.collect(Collectors.toList());
ossService.completeMultipartUpload(objectKey, uploadId, etags);
return null;
})
.exceptionally(e -> {
ossService.abortMultipartUpload(objectKey, uploadId);
throw new RuntimeException("大文件分片上传失败", e);
});
}
}即时通讯
WebSocket 基础
WebSocket 是 HTML5 定义的全双工通信协议,在单个 TCP 连接上提供双向实时数据传输。与 HTTP 不同,WebSocket 支持服务器主动向客户端推送消息。
连接建立过程:
Client Server
│ │
├── HTTP Upgrade Request ──────────────→ │ GET /ws
│ Upgrade: websocket │ Sec-WebSocket-Key: xxxx
│ Sec-WebSocket-Version: 13 │
│ │
│ ←── HTTP 101 Switching Protocols ──────┤ Sec-WebSocket-Accept: yyyy
│ │
├─────── 全双工通信开始 ──────────────────┤
│ ←→ Data Frame (text/binary) │STOMP 协议 + Spring WebSocket
STOMP (Simple Text Oriented Messaging Protocol) 是基于 WebSocket 的子协议,定义了消息的格式、目的地和路由规则。
引入依赖
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-websocket</artifactId>
</dependency>配置 WebSocket 端点
@Configuration
@EnableWebSocketMessageBroker
public class WebSocketConfig implements WebSocketMessageBrokerConfigurer {
/**
* 注册 STOMP 端点,客户端通过此路径建立连接
*/
@Override
public void registerStompEndpoints(StompEndpointRegistry registry) {
registry.addEndpoint("/ws") // WebSocket 连接端点
.setAllowedOriginPatterns("*") // 跨域(生产环境应限定域名)
.withSockJS(); // 降级方案(浏览器不支持 WebSocket 时用轮询)
}
/**
* 配置消息代理
*/
@Override
public void configureMessageBroker(MessageBrokerRegistry registry) {
// 应用前缀:客户端发送消息到服务端时使用的前缀
registry.setApplicationDestinationPrefixes("/app");
// 简单代理:服务端直接广播消息到订阅了 /topic 或 /queue 的客户端
// /topic → 广播模式(所有订阅者收到)
// /queue → 点对点模式(仅指定用户收到)
registry.enableSimpleBroker("/topic", "/queue");
// 点对点用户消息前缀
registry.setUserDestinationPrefix("/user");
}
}@MessageMapping 双向通信
服务端接收客户端消息:
@Controller
public class ChatController {
@Autowired
private SimpMessagingTemplate messagingTemplate;
/**
* 客户端发送消息到 /app/chat.sendMessage
* 服务端处理后广播到 /topic/public
*/
@MessageMapping("/chat.sendMessage")
@SendTo("/topic/public")
public ChatMessage sendMessage(ChatMessage message) {
message.setTimestamp(Instant.now().toString());
return message;
}
/**
* 点对点发送(私聊)
* 客户端发送到 /app/chat.private
*/
@MessageMapping("/chat.private")
public void sendPrivateMessage(PrivateMessage message, Principal principal) {
// 只发送给指定用户
messagingTemplate.convertAndSendToUser(
message.getTargetUser(), // 收件人用户名
"/queue/private", // 目的队列
message // 消息体
);
}
/**
* 用户加入通知
*/
@MessageMapping("/chat.join")
@SendTo("/topic/public")
public ChatMessage join(UserJoin join, SimpMessageHeaderAccessor headerAccessor) {
headerAccessor.getSessionAttributes()
.put("username", join.getUsername());
ChatMessage message = new ChatMessage();
message.setType(MessageType.JOIN);
message.setSender(join.getUsername());
message.setContent(join.getUsername() + " 加入了聊天室");
return message;
}
}消息体定义:
@Data
public class ChatMessage {
private MessageType type;
private String content;
private String sender;
private String timestamp;
}
public enum MessageType {
CHAT,
JOIN,
LEAVE
}
@Data
public class PrivateMessage {
private String content;
private String sender;
private String targetUser;
}客户端 JavaScript 示例:
const socket = new SockJS('/ws');
const stompClient = Stomp.over(socket);
stompClient.connect({}, function (frame) {
console.log('Connected: ' + frame);
// 订阅公共频道
stompClient.subscribe('/topic/public', function (message) {
const msg = JSON.parse(message.body);
displayMessage(msg);
});
// 订阅私信
stompClient.subscribe('/user/queue/private', function (message) {
const msg = JSON.parse(message.body);
displayPrivateMessage(msg);
});
});
// 发送公共消息
function sendMessage(content) {
stompClient.send('/app/chat.sendMessage', {}, JSON.stringify({
content: content,
sender: currentUser,
type: 'CHAT'
}));
}
// 发送私信
function sendPrivateMessage(targetUser, content) {
stompClient.send('/app/chat.private', {}, JSON.stringify({
content: content,
sender: currentUser,
targetUser: targetUser
}));
}WebSocket 鉴权
握手拦截器 + Token 验证
在 WebSocket 握手阶段验证用户身份:
public class AuthChannelInterceptor implements ChannelInterceptor {
@Override
public Message<?> preSend(Message<?> message, MessageChannel channel) {
StompHeaderAccessor accessor = StompHeaderAccessor.wrap(message);
switch (accessor.getCommand()) {
case CONNECT:
// 从请求头获取 Token
String token = accessor.getFirstNativeHeader("Authorization");
if (token == null || token.isEmpty()) {
throw new AuthenticationCredentialsNotFoundException("缺少认证令牌");
}
try {
// 验证 Token,解析用户信息
UserPrincipal user = validateToken(token.replace("Bearer ", ""));
accessor.setUser(user);
accessor.getSessionAttributes()
.put("userId", user.getUserId());
} catch (Exception e) {
throw new AuthenticationCredentialsNotFoundException("Token 验证失败");
}
break;
case SUBSCRIBE:
// 验证订阅权限(只能订阅自己的私信队列)
String destination = accessor.getDestination();
UserPrincipal user = (UserPrincipal) accessor.getUser();
if (destination.startsWith("/user/") && user != null) {
// 防止用户订阅其他人的队列
}
break;
}
return message;
}
private UserPrincipal validateToken(String token) {
// JWT 验证逻辑
Claims claims = Jwts.parser()
.verifyWith(secretKey)
.build()
.parseSignedClaims(token)
.getPayload();
return new UserPrincipal(claims.getSubject(), claims.get("userId", Long.class));
}
}@Data
@AllArgsConstructor
public class UserPrincipal implements Principal {
private String name; // 用户名
private Long userId; // 用户 ID
}注册拦截器:
@Configuration
@EnableWebSocketMessageBroker
public class SecureWebSocketConfig implements WebSocketMessageBrokerConfigurer {
@Override
public void configureClientInboundChannel(ChannelRegistration registration) {
registration.interceptors(new AuthChannelInterceptor());
}
@Override
public void registerStompEndpoints(StompEndpointRegistry registry) {
registry.addEndpoint("/ws")
.setAllowedOriginPatterns("*")
.withSockJS();
}
@Override
public void configureMessageBroker(MessageBrokerRegistry registry) {
registry.enableSimpleBroker("/topic", "/queue");
registry.setApplicationDestinationPrefixes("/app");
registry.setUserDestinationPrefix("/user");
}
}客户端携带 Token 连接:
const token = 'Bearer ' + getAuthToken();
const socket = new SockJS('/ws');
const stompClient = Stomp.over(socket);
stompClient.connect({
'Authorization': token
}, function (frame) {
console.log('Authenticated');
});WebSocket 集群
单机版 SimpleBroker 无法跨进程广播消息,集群环境下需要使用外部消息代理。
Redis 消息代理
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-redis</artifactId>
</dependency>@Configuration
@EnableWebSocketMessageBroker
public class RedisWebSocketConfig implements WebSocketMessageBrokerConfigurer {
@Override
public void registerStompEndpoints(StompEndpointRegistry registry) {
registry.addEndpoint("/ws")
.setAllowedOriginPatterns("*")
.withSockJS();
}
@Override
public void configureMessageBroker(MessageBrokerRegistry registry) {
// 使用 Redis 作为消息代理
registry.enableStompBrokerRelay("/topic", "/queue")
.setRelayHost("redis")
.setRelayPort(61613) // Redis STOMP 端口
.setSystemLogin("admin")
.setSystemPasscode("password")
.setClientLogin("guest")
.setClientPasscode("guest");
registry.setApplicationDestinationPrefixes("/app");
registry.setUserDestinationPrefix("/user");
}
}启用 Redis 的 STOMP 支持需要在 Redis 配置中加载 STOMP 模块或使用 Redis Stack。
RabbitMQ STOMP
RabbitMQ 提供原生 STOMP 插件,更适合生产级消息路由。
RabbitMQ 端启用 STOMP:
# 启用 STOMP 插件
rabbitmq-plugins enable rabbitmq_stomp
# 启用 Web STOMP(可选,支持浏览器直连 RabbitMQ)
rabbitmq-plugins enable rabbitmq_web_stompSpring Boot 配置:
@Configuration
@EnableWebSocketMessageBroker
public class RabbitWebSocketConfig implements WebSocketMessageBrokerConfigurer {
@Override
public void registerStompEndpoints(StompEndpointRegistry registry) {
registry.addEndpoint("/ws")
.setAllowedOriginPatterns("*")
.withSockJS();
}
@Override
public void configureMessageBroker(MessageBrokerRegistry registry) {
// RabbitMQ STOMP 代理
registry.enableStompBrokerRelay("/topic", "/queue")
.setRelayHost("rabbitmq.example.com")
.setRelayPort(61613) // STOMP 插件端口
.setSystemLogin("guest")
.setSystemPasscode("guest")
.setVirtualHost("/");
registry.setApplicationDestinationPrefixes("/app");
registry.setUserDestinationPrefix("/user");
}
}集群架构:
Client A ─→ WebSocket ─→ App Instance 1 ─→ RabbitMQ STOMP ─→ Topic Exchange
│
Client B ─→ WebSocket ─→ App Instance 2 ──────────────────────────────┘当任意实例向 /topic/xxx 发送消息时,RabbitMQ 将消息广播到所有订阅了该 Topic 的实例。
聊天空闲检测
心跳机制
STOMP 协议内置心跳检测,由客户端和服务端协商周期。
const socket = new SockJS('/ws');
const stompClient = Stomp.over(socket);
// 配置心跳:每 10 秒发送心跳,每 10 秒等待服务器心跳
stompClient.heartbeat.outgoing = 10000; // 10s
stompClient.heartbeat.incoming = 10000; // 10s
stompClient.connect({}, function (frame) {
console.log('Connected with heartbeats');
});服务端心跳配置:
@Override
public void configureWebSocketTransport(WebSocketTransportRegistration registration) {
registration
.setSendTimeLimit(15 * 1000) // 发送超时 15s
.setSendBufferSizeLimit(512 * 1024) // 发送缓存 512KB
.setMessageSizeLimit(128 * 1024); // 消息大小限制 128KB
}重连机制
客户端自动重连:
let stompClient = null;
let reconnectTimer = null;
function connect() {
const socket = new SockJS('/ws');
stompClient = Stomp.over(socket);
// 关闭自动重连日志
stompClient.debug = null;
stompClient.connect({
'Authorization': 'Bearer ' + getAuthToken()
}, function (frame) {
console.log('连接成功');
// 重新订阅所有频道
subscribeChannels();
// 清除重连定时器
if (reconnectTimer) {
clearTimeout(reconnectTimer);
reconnectTimer = null;
}
}, function (error) {
console.error('连接断开,准备重连:', error);
scheduleReconnect();
});
}
function scheduleReconnect() {
if (reconnectTimer) return;
reconnectTimer = setTimeout(function () {
console.log('尝试重连...');
connect();
}, 5000); // 5 秒后重连
}
function disconnect() {
if (stompClient !== null) {
stompClient.disconnect();
}
if (reconnectTimer) {
clearTimeout(reconnectTimer);
reconnectTimer = null;
}
}
// 指数退避重连
function scheduleReconnectWithBackoff(attempt) {
const delay = Math.min(1000 * Math.pow(2, attempt), 30000);
setTimeout(function () {
console.log('第 ' + (attempt + 1) + ' 次重连...');
connect();
}, delay);
}
connect();服务端空闲检测
@Component
public class WebSocketSessionManager {
private final Map<String, SessionInfo> sessions = new ConcurrentHashMap<>();
/**
* 注册会话
*/
public void registerSession(String sessionId, String userId) {
sessions.put(sessionId, new SessionInfo(sessionId, userId, Instant.now()));
}
/**
* 更新心跳时间
*/
public void updateHeartbeat(String sessionId) {
SessionInfo info = sessions.get(sessionId);
if (info != null) {
info.setLastHeartbeat(Instant.now());
}
}
/**
* 检测空闲连接(超过 60 秒无心跳则关闭)
*/
@Scheduled(fixedRate = 30000)
public void checkIdleSessions() {
Instant threshold = Instant.now().minus(60, ChronoUnit.SECONDS);
sessions.entrySet().removeIf(entry -> {
SessionInfo info = entry.getValue();
if (info.getLastHeartbeat().isBefore(threshold)) {
log.warn("关闭空闲连接: sessionId={}, userId={}",
info.getSessionId(), info.getUserId());
// 发送关闭通知并清理
return true;
}
return false;
});
}
@Data
@AllArgsConstructor
private static class SessionInfo {
private String sessionId;
private String userId;
private Instant lastHeartbeat;
}
}消息推送方案对比
WebSocket vs SSE vs 短轮询 vs 长轮询 vs MQTT
方案概述
| 方案 | 协议 | 通信模式 | 浏览器支持 | 适用场景 |
|---|---|---|---|---|
| WebSocket | ws/wss | 全双工 | HTML5 原生 | 即时聊天、协同编辑、实时游戏 |
| Server-Sent Events (SSE) | http/https | 服务器单向推送 | EventSource API | 通知推送、股票行情、日志流 |
| 短轮询 | http/https | 客户端定时请求 | 全兼容 | 简单状态刷新、兼容老旧浏览器 |
| 长轮询 | http/https | 客户端保持连接 | 全兼容 | 兼容性要求高的实时场景 |
| MQTT | mqtt/mqtts | 发布/订阅 | 需 MQTT.js | IoT 物联网、弱网络环境 |
详细对比
| 特性 | WebSocket | SSE | 短轮询 | 长轮询 | MQTT |
|---|---|---|---|---|---|
| 方向 | 双向 | 服务端到客户端 | 客户端请求/服务端响应 | 双向模拟 | 双向 |
| 实时性 | 极高 | 高 | 取决于轮询间隔 | 较高 | 高 |
| 延迟 | 10-50ms | 50-200ms | 1s-30s | 200ms-2s | 50-200ms |
| 连接开销 | 高(需握手升级) | 低(单 HTTP 连接) | 极低(短连接) | 中 | 低 |
| 服务端资源 | 高(需维持长连接) | 中 | 低(无长连接) | 高(需挂起请求) | 低 |
| 断线重连 | 需手动实现 | 自动重连 | 天然支持 | 需手动实现 | 内置 QoS 机制 |
| 二进制支持 | 原生支持 | 需 base64 编码 | 支持 | 支持 | 原生支持 |
| 跨域 | 需服务端配置 | 需服务端配置 | 天然支持 | 天然支持 | 需代理转发 |
| 浏览器兼容 | IE10+ | 除 IE/Edge 外 | 全兼容 | 全兼容 | 需额外库 |
| 协议复杂度 | 中 | 低 | 极低 | 低 | 中 |
WebSocket 适用场景
- 双工通信:即时聊天、客服系统、协作白板
- 低延迟要求:在线游戏、金融行情、实时竞拍
- 频繁消息交换:协同编辑、共享光标
优点: 实时性极强、双向通信、消息协议灵活(STOMP / 自定义)
缺点: 连接维护成本高(需心跳保活)、水平扩展需外置消息代理
SSE 适用场景
- 单向推送:通知中心、消息提醒、系统公告
- 流式数据:日志实时输出、AI 对话流式响应、新闻推送
- 数据更新频繁:股票/加密货币行情、监控面板
@RestController
public class SseController {
private final Map<String, SseEmitter> emitters = new ConcurrentHashMap<>();
@GetMapping("/sse/subscribe/{userId}")
public SseEmitter subscribe(@PathVariable String userId) {
SseEmitter emitter = new SseEmitter(3600_000L); // 1 小时超时
emitters.put(userId, emitter);
emitter.onCompletion(() -> emitters.remove(userId));
emitter.onTimeout(() -> emitters.remove(userId));
return emitter;
}
/**
* 推送通知给指定用户
*/
public void pushNotification(String userId, Notification notification) {
SseEmitter emitter = emitters.get(userId);
if (emitter != null) {
try {
emitter.send(SseEmitter.event()
.name("notification")
.data(notification));
} catch (IOException e) {
emitters.remove(userId);
}
}
}
/**
* 广播给所有订阅者
*/
public void broadcast(String eventName, Object data) {
List<String> disconnected = new ArrayList<>();
emitters.forEach((userId, emitter) -> {
try {
emitter.send(SseEmitter.event()
.name(eventName)
.data(data));
} catch (IOException e) {
disconnected.add(userId);
}
});
disconnected.forEach(emitters::remove);
}
}客户端:
const eventSource = new EventSource('/sse/subscribe/' + userId);
eventSource.addEventListener('notification', function (event) {
const notification = JSON.parse(event.data);
showNotification(notification);
});
eventSource.onerror = function () {
console.log('SSE 连接断开,重新连接...');
// 浏览器会自动重连
};短轮询适用场景
- 老旧系统兼容:需要支持 IE8/9 等低级浏览器
- 低频率刷新:系统状态监控(每分钟刷新一次)
- 非实时场景:代办提醒、邮件未读数
function pollStatus() {
setInterval(async () => {
const response = await fetch('/api/status');
const status = await response.json();
updateUI(status);
}, 5000); // 每 5 秒轮询一次
}长轮询适用场景
- 兼容性优先:无法使用 WebSocket/SSE 但追求实时性
- 过渡方案:在升级到 WebSocket 前的中间方案
@RestController
public class LongPollingController {
@GetMapping("/poll/messages")
public DeferredResult<List<Message>> pollMessages(
@RequestParam String userId,
@RequestParam(defaultValue = "0") long lastMessageId) {
DeferredResult<List<Message>> deferredResult = new DeferredResult<>(30_000L);
// 如果有新消息立即返回,否则等待
List<Message> newMessages = messageService.getNewMessages(userId, lastMessageId);
if (!newMessages.isEmpty()) {
deferredResult.setResult(newMessages);
} else {
// 注册异步监听,有新消息时触发
messageService.registerListener(userId, lastMessageId, deferredResult);
deferredResult.onTimeout(() ->
deferredResult.setResult(Collections.emptyList()));
}
return deferredResult;
}
}MQTT 适用场景
- IoT 物联网:传感器数据采集、智能设备控制
- 弱网络环境:移动端、卫星通信、网络不稳定场景
- 低带宽需求:消息头极小(最小仅 2 字节)
MQTT 架构:
┌─────────┐ ┌─────────────┐ ┌─────────┐
│ Publisher │ ←→ │ MQTT Broker │ ←→ │ Subscriber │
│ (传感器) │ │ (如 EMQX) │ │ (后端服务) │
└─────────┘ └─────────────┘ └─────────┘
↑↓
┌────────────┐
│ WebSocket │
│ 桥接 │
└────────────┘技术选型决策流程
是否需要双向实时通信?
├─ 是 → WebSocket(推荐)/ MQTT(IoT 场景)
│ ↓
│ 是否高频小消息 + 弱网络?→ MQTT
│ 否则 → WebSocket + STOMP
│
└─ 否 → 是否服务器主动推送?
├─ 是 → SSE(现代浏览器)/ 长轮询(兼容)
│ ↓
│ 需要兼容老旧浏览器?→ 长轮询
│ 否则 → SSE
│
└─ 否 → 短轮询(简单定时刷新)混合架构示例
实际项目中常采用多方案混合架构:
┌──────────────┐ ┌─────────────────────┐
│ Web App │ ←→ │ WebSocket (STOMP) │ ← 即时消息、协同
│ │ ←→ │ SSE │ ← 通知推送、行情
│ │ ←→ │ REST (短轮询降级) │ ← 降级方案
└──────────────┘ └─────────────────────┘
┌──────────────┐ ┌─────────────────────┐
│ Mobile App │ ←→ │ WebSocket (STOMP) │ ← 即时消息
│ │ ←→ │ MQTT │ ← 离线消息、推送
└──────────────┘ └─────────────────────┘
┌──────────────┐ ┌─────────────────────┐
│ IoT Devices │ ←→ │ MQTT (EMQX Broker) │ ← 设备数据采集
└──────────────┘ └─────────────────────┘