SaaS 开放平台对接规范 / 第三方接口统一管理
1. SaaS 开放平台概述
1.1 开放平台架构
SaaS 开放平台作为生态系统的核心枢纽,为第三方开发者提供标准化的接入能力和管理界面。整体架构由以下核心组件构成:
| 组件 | 职责 | 技术选型参考 |
|---|---|---|
| 开发者门户 | 开发者注册、应用管理、文档查阅、调试工具 | Vue/React + VitePress |
| API 网关 | 统一入口、路由转发、限流熔断、协议转换 | Spring Cloud Gateway / Kong |
| 认证中心 | AK/SK 签发、JWT 颁发、OAuth2 授权 | Spring Authorization Server / Keycloak |
| 沙箱环境 | 隔离的测试环境、模拟数据、Mock 接口 | 独立部署 + Mock Server |
| 计量计费 | 调用次数统计、套餐管理、用量预警 | Prometheus + 自定义计费服务 |
| 开发者文档 | OpenAPI 规范、SDK 下载、示例代码 | VitePress / Swagger UI |
架构分层示意图:
text
+---------------------------------------------------------------------+
| 第三方开发者/合作伙伴 |
+---------------------------------------------------------------------+
|
HTTPS / HTTP
|
+---------------------------------------------------------------------+
| API 网关 (统一入口层) |
| 路由转发 / 限流熔断 / IP 白名单 / 请求签名校验 / 日志采集 |
+---------------------------------------------------------------------+
| | |
v v v
+----------------+ +------------------+ +------------------+
| 认证中心 | | 核心业务服务 | | 计量计费服务 |
| AK/SK 管理 | | 数据开放 API | | 调用量统计 |
| JWT 签发 | | 业务开放 API | | 套餐管理 |
| OAuth2 授权 | | 能力开放 API | | 账单生成 |
+----------------+ +------------------+ +------------------+
| | |
v v v
+---------------------------------------------------------------------+
| 后端基础设施层 |
| 数据库 / Redis / MQ / 配置中心 / 注册中心 / 日志中心 |
+---------------------------------------------------------------------+1.2 开放能力分类
| 能力类型 | 说明 | 示例 |
|---|---|---|
| 数据开放 | 开放平台积累的业务数据,供开发者查询和分析 | 订单数据、用户画像、商品信息 |
| 业务开放 | 将核心业务流程封装为 API,支持开发者调用 | 下单、支付、退款、物流查询 |
| 能力开放 | 开放平台的技术能力和算法能力 | 智能推荐、图像识别、风控评分 |
| 硬件开放 | 开放 IoT 设备、打印设备等硬件能力 | 云打印、智能硬件控制 |
1.3 开发者生态
1.3.1 应用市场
应用市场是开发者发布应用、商家选购应用的平台,核心流程如下:
text
开发者提交应用
|
v
应用审核 (自动化 + 人工)
|
v
审核通过 -> 上架应用市场
|
v
商家选购 -> 授权应用 -> 开通服务
|
v
应用运行 (调用开放平台 API)
|
v
违规/到期 -> 下架/停用1.3.2 应用审核
应用审核包含以下检查项:
| 审核项 | 检查内容 | 审核方式 |
|---|---|---|
| 应用信息 | 应用名称、图标、描述完整性 | 自动 |
| 权限申请 | 申请的 API 权限是否合理 | 自动 + 人工 |
| 安全审查 | 回调地址白名单、IP 白名单 | 自动 |
| 功能测试 | 应用功能是否符合描述 | 人工 |
| 合规审查 | 是否涉及违规内容、数据合规 | 人工 |
1.3.3 应用上架/下架
java
/**
* 应用状态机
*/
public enum AppStatus {
DRAFT(0, "草稿"),
PENDING_REVIEW(1, "待审核"),
REVIEWING(2, "审核中"),
REJECTED(3, "审核驳回"),
APPROVED(4, "已通过"),
PUBLISHED(5, "已上架"),
OFF_SHELF(6, "已下架"),
DISABLED(7, "已停用");
private final int code;
private final String description;
}
/**
* 应用审核服务
*/
@Service
public class AppReviewService {
@Autowired
private AppMapper appMapper;
@Transactional
public void submitForReview(Long appId) {
App app = appMapper.selectById(appId);
if (app.getStatus() != AppStatus.DRAFT && app.getStatus() != AppStatus.REJECTED) {
throw new BusinessException("当前状态不允许提交审核");
}
app.setStatus(AppStatus.PENDING_REVIEW);
app.setSubmitTime(LocalDateTime.now());
appMapper.updateById(app);
// 发送审核通知
notifyReviewer(app);
}
@Transactional
public void approve(Long appId, String reviewer) {
App app = appMapper.selectById(appId);
app.setStatus(AppStatus.APPROVED);
app.setReviewer(reviewer);
app.setReviewTime(LocalDateTime.now());
appMapper.updateById(app);
}
@Transactional
public void publish(Long appId) {
App app = appMapper.selectById(appId);
app.setStatus(AppStatus.PUBLISHED);
app.setPublishTime(LocalDateTime.now());
appMapper.updateById(app);
}
@Transactional
public void offShelf(Long appId, String reason) {
App app = appMapper.selectById(appId);
app.setStatus(AppStatus.OFF_SHELF);
app.setOffShelfReason(reason);
app.setOffShelfTime(LocalDateTime.now());
appMapper.updateById(app);
// 通知开发者
notifyDeveloper(app.getDeveloperId(), "您的应用已下架,原因: " + reason);
}
}2. API 统一管理
2.1 API 注册与版本管理
2.1.1 API 注册
每个开放 API 需要在 API 管理平台注册元数据,包括接口路径、请求方式、参数定义、权限等级等。
java
/**
* API 定义实体
*/
@Data
@TableName("open_api_definition")
public class ApiDefinition {
@TableId(type = IdType.AUTO)
private Long id;
private String apiName; // API 名称
private String apiPath; // 请求路径,如 /v1/orders
private String httpMethod; // GET / POST / PUT / DELETE
private String version; // 版本号,如 v1、v2
private String category; // 分类,如 order、user、payment
private String description; // 接口描述
private Integer authType; // 鉴权方式:1-AK/SK, 2-JWT, 3-OAuth2
private Integer rateLimitQps; // 默认 QPS 限制
private Integer timeoutMs; // 超时时间(毫秒)
private Boolean requireSsl; // 是否强制 HTTPS
private String status; // ENABLED / DEPRECATED / DISABLED
private String requestSchema; // 请求参数 JSON Schema
private String responseSchema; // 响应参数 JSON Schema
private String errorCodes; // 错误码定义 JSON
private LocalDateTime createdAt;
private LocalDateTime updatedAt;
}2.1.2 版本策略
| 版本类型 | 格式 | 兼容性 | 生命周期 |
|---|---|---|---|
| 主版本 | v1、v2 | 不兼容,允许 Breaking Change | 长期维护 |
| 次版本 | v1.0、v1.1 | 向后兼容,新增字段或接口 | 6 个月 |
| 修订版本 | v1.0.0、v1.0.1 | 完全兼容,仅修复 Bug | 3 个月 |
版本管理策略:
text
v1 发布 -> v1.1 发布(兼容 v1)-> v2 发布(不兼容 v1)
|
v1 标记废弃 (deprecated)
|
v1 进入维护期(仅修复安全问题)
|
v1 下架 ( sunset ),返回 410 Gonejava
/**
* API 版本路由
* 通过请求路径中的版本号分发到不同的处理器
*/
@RestController
public class OrderApiController {
/**
* v1 版本 - 旧版订单查询
*/
@GetMapping("/v1/orders/{id}")
public OrderV1Response getOrderV1(@PathVariable Long id) {
Order order = orderService.getById(id);
return OrderV1Response.from(order); // 旧版响应格式
}
/**
* v2 版本 - 新版订单查询,优化了响应结构
*/
@GetMapping("/v2/orders/{id}")
public OrderV2Response getOrderV2(@PathVariable Long id) {
Order order = orderService.getById(id);
return OrderV2Response.from(order); // 新版响应格式,增加字段
}
}
/**
* 版本废弃通知
*/
@Scheduled(cron = "0 0 0 * * ?") // 每日执行
public void checkDeprecatedApis() {
List<ApiDefinition> deprecatedApis = apiDefinitionMapper.selectList(
new LambdaQueryWrapper<ApiDefinition>()
.eq(ApiDefinition::getStatus, "DEPRECATED")
.lt(ApiDefinition::getUpdatedAt, LocalDateTime.now().minusMonths(6)));
for (ApiDefinition api : deprecatedApis) {
// 通知已订阅的应用开发者
List<AppSubscription> subscribers = subscriptionMapper.selectByApiId(api.getId());
for (AppSubscription sub : subscribers) {
notificationService.notifyDeveloper(
sub.getDeveloperId(),
"您使用的 API [" + api.getApiName() + "] 将于 30 天后下架,请尽快迁移");
}
// 标记为即将下架
api.setStatus("SUNSETTING");
apiDefinitionMapper.updateById(api);
}
}2.2 API 文档规范
2.2.1 OpenAPI 3.0 标准
所有开放 API 必须遵循 OpenAPI 3.0 规范定义,通过注解自动生成文档。
java
@RestController
@RequestMapping("/v1/orders")
@Tag(name = "订单管理", description = "订单创建、查询、退款等接口")
public class OrderOpenApiController {
@PostMapping
@Operation(summary = "创建订单", description = "第三方开发者通过此接口创建交易订单")
@ApiResponses(value = {
@ApiResponse(responseCode = "200", description = "创建成功",
content = @Content(schema = @Schema(implementation = CreateOrderResponse.class))),
@ApiResponse(responseCode = "400", description = "参数错误",
content = @Content(schema = @Schema(implementation = ErrorResponse.class))),
@ApiResponse(responseCode = "401", description = "鉴权失败"),
@ApiResponse(responseCode = "429", description = "请求频率超限")
})
public Response<CreateOrderResponse> createOrder(
@Valid @RequestBody CreateOrderRequest request,
@RequestHeader("X-Access-Key") String accessKey,
@RequestHeader("X-Signature") String signature,
@RequestHeader("X-Timestamp") Long timestamp) {
// 1. 签名校验
signatureService.verify(accessKey, timestamp, signature, request);
// 2. 权限校验
permissionService.checkApiPermission(accessKey, "order:create");
// 3. 业务处理
Order order = orderService.create(request.toOrder());
return Response.success(CreateOrderResponse.from(order));
}
}对应的 OpenAPI 规范片段(自动生成):
yaml
openapi: 3.0.0
info:
title: SaaS 开放平台 API
version: v1
description: SaaS 开放平台提供订单、支付、物流等标准化接口
servers:
- url: https://api.example.com
description: 生产环境
- url: https://sandbox-api.example.com
description: 沙箱环境
paths:
/v1/orders:
post:
summary: 创建订单
tags:
- 订单管理
parameters:
- name: X-Access-Key
in: header
required: true
schema:
type: string
- name: X-Signature
in: header
required: true
schema:
type: string
- name: X-Timestamp
in: header
required: true
schema:
type: integer
requestBody:
content:
application/json:
schema:
$ref: '#/components/schemas/CreateOrderRequest'
responses:
'200':
description: 创建成功
content:
application/json:
schema:
$ref: '#/components/schemas/CreateOrderResponse'2.2.2 在线调试
通过 Swagger UI 或 Knife4j 提供在线调试功能,开发者无需编写代码即可测试 API。
yaml
# application.yml - Knife4j 配置
knife4j:
enable: true
setting:
enableSwaggerModels: true
enableDocumentManage: true
enableReloadCacheParameter: true
enableVersion: true
enableRequestCache: true
enableFooter: false2.2.3 SDK 生成
使用 OpenAPI Generator 自动生成多语言 SDK:
xml
<!-- pom.xml - Maven 插件配置 -->
<plugin>
<groupId>org.openapitools</groupId>
<artifactId>openapi-generator-maven-plugin</artifactId>
<version>7.2.0</version>
<executions>
<execution>
<goals>
<goal>generate</goal>
</goals>
<configuration>
<inputSpec>${project.basedir}/src/main/resources/openapi.yaml</inputSpec>
<generatorName>java</generatorName>
<output>${project.build.directory}/generated-sdk</output>
<packageName>com.example.openapi.sdk</packageName>
<library>okhttp-gson</library>
</configuration>
</execution>
</executions>
</plugin>bash
# 使用 CLI 生成 SDK
# Java
openapi-generator-cli generate -i openapi.yaml -g java -o sdk/java
# Python
openapi-generator-cli generate -i openapi.yaml -g python -o sdk/python
# JavaScript
openapi-generator-cli generate -i openapi.yaml -g javascript -o sdk/javascript
# Go
openapi-generator-cli generate -i openapi.yaml -g go -o sdk/go2.3 API 路由与负载均衡
API 网关根据请求路径、Header 等条件将请求路由到后端服务。
yaml
# Spring Cloud Gateway 路由配置
spring:
cloud:
gateway:
routes:
# 开放平台 API 路由
- id: openapi-v1-orders
uri: lb://order-service
predicates:
- Path=/v1/orders/**
filters:
- StripPrefix=0
- name: RequestRateLimiter
args:
key-resolver: "#{@apiKeyResolver}"
redis-rate-limiter.replenishRate: 100
redis-rate-limiter.burstCapacity: 200
- id: openapi-v1-payments
uri: lb://payment-service
predicates:
- Path=/v1/payments/**
filters:
- StripPrefix=0
- id: openapi-v2-orders
uri: lb://order-service-v2
predicates:
- Path=/v2/orders/**
filters:
- StripPrefix=02.4 API 限流降级
2.4.1 限流维度
| 限流维度 | 说明 | 典型配置 |
|---|---|---|
| 应用级 | 每个应用的整体调用量限制 | 免费版 1000 QPS |
| 用户级 | 每个开发者的调用量限制 | 单个开发者 5000 QPS |
| API 级 | 单个 API 的调用量限制 | 下单接口 200 QPS |
| IP 级 | 来源 IP 的调用量限制 | 单 IP 1000 QPS |
2.4.2 Sentinel 限流实现
java
/**
* API 限流切面
*/
@Aspect
@Component
public class OpenApiRateLimitAspect {
private static final Logger log = LoggerFactory.getLogger(OpenApiRateLimitAspect.class);
@Around("@annotation(rateLimit)")
public Object doRateLimit(ProceedingJoinPoint pjp, RateLimit rateLimit) throws Throwable {
// 获取当前请求的上下文
RequestContext context = RequestContextHolder.get();
// 应用级限流
String appResource = "app:" + context.getAppId();
try (Entry entry = SphU.entry(appResource, EntryType.IN)) {
// API 级限流
String apiResource = "api:" + context.getApiPath();
try (Entry apiEntry = SphU.entry(apiResource, EntryType.IN)) {
return pjp.proceed();
} catch (BlockException e) {
log.warn("API 限流触发: appId={}, api={}", context.getAppId(), context.getApiPath());
throw new RateLimitException("API 调用频率超限,请稍后重试");
}
} catch (BlockException e) {
log.warn("应用级限流触发: appId={}", context.getAppId());
throw new RateLimitException("应用调用量已达上限");
}
}
}
/**
* 动态限流规则配置
* 根据套餐绑定不同的限流策略
*/
@Service
public class RateLimitRuleService {
@Autowired
private RedisTemplate<String, String> redisTemplate;
/**
* 为应用加载限流规则
*/
public void loadRateLimitRules(String appId, String planCode) {
// 根据套餐获取限流配置
PlanConfig planConfig = getPlanConfig(planCode);
// 应用级规则
FlowRule appRule = new FlowRule();
appRule.setResource("app:" + appId);
appRule.setGrade(RuleConstant.FLOW_GRADE_QPS);
appRule.setCount(planConfig.getAppQps());
appRule.setLimitApp("default");
FlowRuleManager.loadRules(List.of(appRule));
// API 级规则存储在 Redis 中,由网关加载
String ruleKey = "ratelimit:rules:" + appId;
redisTemplate.opsForValue().set(ruleKey, JSON.toJSONString(planConfig.getApiRules()));
}
private PlanConfig getPlanConfig(String planCode) {
switch (planCode) {
case "free":
return new PlanConfig(100, 50, 10);
case "professional":
return new PlanConfig(1000, 500, 100);
case "enterprise":
return new PlanConfig(10000, 5000, 1000);
default:
throw new BusinessException("未知套餐: " + planCode);
}
}
}
@Data
@AllArgsConstructor
public class PlanConfig {
private int appQps; // 应用级 QPS
private int userQps; // 用户级 QPS
private int defaultApiQps; // 默认单 API QPS
}2.5 API 鉴权
2.5.1 AK/SK HMAC-SHA256
AK/SK 鉴权是最常见的开放平台鉴权方式,Access Key 标识身份,Secret Key 用于签名。
签名流程:
text
1. 开发者使用 AK + SK 对请求参数进行 HMAC-SHA256 签名
2. 请求时携带 AK、Signature、Timestamp 等 Header
3. 服务端根据 AK 查找对应的 SK,重新计算签名并比对签名算法:
java
/**
* 服务端签名校验
*/
@Component
public class SignatureService {
@Autowired
private CredentialManager credentialManager;
private static final long SIGN_TIMEOUT_SECONDS = 300; // 5 分钟有效
/**
* 验证请求签名
*/
public void verify(String accessKey, Long timestamp, String signature, Object body) {
// 1. 校验时间戳防止重放攻击
long now = System.currentTimeMillis() / 1000;
if (Math.abs(now - timestamp) > SIGN_TIMEOUT_SECONDS) {
throw new AuthException("请求已过期");
}
// 2. 查询 AK 对应的 SK
AppCredential credential = credentialManager.getCredential(accessKey);
if (credential == null || credential.getStatus() != 1) {
throw new AuthException("无效的 AccessKey");
}
// 3. 重新计算签名
String expectedSign = calculateSignature(
credential.getAccessSecret(), timestamp, body);
// 4. 比对签名(使用 MessageDigest.isEqual 防止时序攻击)
if !MessageDigest.isEqual(
signature.getBytes(StandardCharsets.UTF_8),
expectedSign.getBytes(StandardCharsets.UTF_8)) {
throw new AuthException("签名验证失败");
}
}
/**
* 签名计算
* signature = HMAC-SHA256(secret, timestamp + "|" + body)
*/
public String calculateSignature(String secret, Long timestamp, Object body) {
String bodyStr = body == null ? "" : JSON.toJSONString(body);
String data = timestamp + "|" + bodyStr;
try {
Mac mac = Mac.getInstance("HmacSHA256");
SecretKeySpec keySpec = new SecretKeySpec(
secret.getBytes(StandardCharsets.UTF_8), "HmacSHA256");
mac.init(keySpec);
byte[] signBytes = mac.doFinal(data.getBytes(StandardCharsets.UTF_8));
return Base64.getEncoder().encodeToString(signBytes);
} catch (Exception e) {
throw new AuthException("签名计算失败", e);
}
}
/**
* 校验请求中的 Nonce 防止重放
*/
public boolean checkNonce(String nonce) {
String key = "nonce:" + nonce;
Boolean success = redisTemplate.opsForValue().setIfAbsent(key, "1",
Duration.ofSeconds(SIGN_TIMEOUT_SECONDS));
return Boolean.TRUE.equals(success);
}
}客户端签名示例:
java
/**
* 客户端签名工具
*/
public class OpenApiClient {
private final String accessKey;
private final String accessSecret;
private final String baseUrl;
private final OkHttpClient httpClient;
public OpenApiClient(String accessKey, String accessSecret, String baseUrl) {
this.accessKey = accessKey;
this.accessSecret = accessSecret;
this.baseUrl = baseUrl;
this.httpClient = new OkHttpClient.Builder()
.connectTimeout(10, TimeUnit.SECONDS)
.readTimeout(30, TimeUnit.SECONDS)
.build();
}
/**
* 发送签名请求
*/
public String call(String method, String path, Object body) {
long timestamp = System.currentTimeMillis() / 1000;
String bodyJson = body == null ? "" : JSON.toJSONString(body);
String signature = sign(timestamp, bodyJson);
Request.Builder builder = new Request.Builder()
.url(baseUrl + path)
.addHeader("X-Access-Key", accessKey)
.addHeader("X-Signature", signature)
.addHeader("X-Timestamp", String.valueOf(timestamp))
.addHeader("Content-Type", "application/json");
if (body != null) {
builder.method(method, RequestBody.create(bodyJson, MediaType.parse("application/json")));
} else {
builder.method(method, null);
}
try (Response response = httpClient.newCall(builder.build()).execute()) {
if (!response.isSuccessful()) {
String errorBody = response.body() != null ? response.body().string() : "";
throw new ApiException(response.code(), "API 调用失败: " + errorBody);
}
return response.body() != null ? response.body().string() : null;
} catch (IOException e) {
throw new ApiException("网络异常", e);
}
}
private String sign(long timestamp, String body) {
String data = timestamp + "|" + body;
try {
Mac mac = Mac.getInstance("HmacSHA256");
SecretKeySpec keySpec = new SecretKeySpec(
accessSecret.getBytes(StandardCharsets.UTF_8), "HmacSHA256");
mac.init(keySpec);
byte[] signBytes = mac.doFinal(data.getBytes(StandardCharsets.UTF_8));
return Base64.getEncoder().encodeToString(signBytes);
} catch (Exception e) {
throw new RuntimeException("签名失败", e);
}
}
}2.5.2 JWT 鉴权
适用于服务端到服务端的认证场景。
java
/**
* JWT Token 签发
*/
@Service
public class JwtTokenService {
@Value("${jwt.secret}")
private String secret;
@Value("${jwt.expiration-hours}")
private int expirationHours;
/**
* 为应用生成 JWT Token
*/
public String generateToken(AppCredential credential) {
long now = System.currentTimeMillis();
Date expiry = new Date(now + TimeUnit.HOURS.toMillis(expirationHours));
return Jwts.builder()
.setIssuer("SaaS-Platform")
.setSubject(credential.getAppId())
.claim("app_name", credential.getAppName())
.claim("developer_id", credential.getDeveloperId())
.claim("scopes", credential.getApiScopes())
.setIssuedAt(new Date(now))
.setExpiration(expiry)
.signWith(SignatureAlgorithm.HS256, secret.getBytes())
.compact();
}
/**
* 验证 JWT Token
*/
public Claims verifyToken(String token) {
try {
return Jwts.parser()
.setSigningKey(secret.getBytes())
.parseClaimsJws(token)
.getBody();
} catch (ExpiredJwtException e) {
throw new AuthException("Token 已过期");
} catch (JwtException e) {
throw new AuthException("无效的 Token");
}
}
/**
* 从 Token 中提取应用信息
*/
public TokenInfo parseToken(String token) {
Claims claims = verifyToken(token);
TokenInfo info = new TokenInfo();
info.setAppId(claims.getSubject());
info.setAppName((String) claims.get("app_name"));
info.setDeveloperId((String) claims.get("developer_id"));
info.setScopes((List<String>) claims.get("scopes"));
return info;
}
}2.5.3 OAuth2 授权
适用于需要获取用户授权的场景(如代用户操作)。
java
/**
* OAuth2 授权配置
*/
@Configuration
@EnableAuthorizationServer
public class OAuth2AuthorizationServerConfig extends AuthorizationServerConfigurerAdapter {
@Autowired
private AuthenticationManager authenticationManager;
@Autowired
private AppClientDetailsService clientDetailsService;
@Override
public void configure(ClientDetailsServiceConfigurer clients) throws Exception {
clients.withClientDetails(clientDetailsService);
}
@Override
public void configure(AuthorizationServerEndpointsConfigurer endpoints) {
endpoints
.authenticationManager(authenticationManager)
.tokenStore(tokenStore())
.accessTokenConverter(accessTokenConverter());
}
@Bean
public TokenStore tokenStore() {
return new RedisTokenStore(redisConnectionFactory);
}
@Bean
public JwtAccessTokenConverter accessTokenConverter() {
JwtAccessTokenConverter converter = new JwtAccessTokenConverter();
converter.setSigningKey(secret);
return converter;
}
}
/**
* OAuth2 应用注册
*/
@Service
public class OAuth2ClientService {
/**
* 注册 OAuth2 应用
*/
public void registerClient(OAuth2ClientRequest request) {
OAuth2Client client = new OAuth2Client();
client.setClientId(generateClientId());
client.setClientSecret(generateClientSecret());
client.setClientName(request.getClientName());
client.setRedirectUris(String.join(",", request.getRedirectUris()));
client.setGrantTypes("authorization_code,refresh_token");
client.setScopes(String.join(",", request.getScopes()));
client.setAccessTokenValidity(7200); // 2 小时
client.setRefreshTokenValidity(2592000); // 30 天
client.setStatus(1);
clientMapper.insert(client);
}
}3. SaaS 多租户架构
3.1 租户模型
3.1.1 隔离模式对比
| 隔离模式 | 数据隔离性 | 成本 | 运维复杂度 | 适用场景 |
|---|---|---|---|---|
| 独立数据库 | 最强 | 高 | 高 | 大型企业、金融合规要求高 |
| 共享数据库 + Schema 隔离 | 较强 | 中 | 中 | 中型客户,需要一定程度隔离 |
| 共享数据库 + 共享 Schema | 一般 | 低 | 低 | 小型客户,成本敏感 |
| 混合模式 | 灵活 | 按需 | 高 | 同时服务大中小客户 |
3.1.2 混合模式实现
java
/**
* 租户信息模型
*/
@Data
@TableName("saas_tenant")
public class Tenant {
@TableId
private String tenantId; // 租户 ID (唯一标识)
private String tenantName; // 租户名称
private String contactName; // 联系人
private String contactPhone; // 联系电话
private String domain; // 独立域名
private String planCode; // 套餐编码: free/professional/enterprise
private Integer isolationMode; // 隔离模式: 1-独立库, 2-Schema 隔离, 3-共享表
private String datasourceKey; // 数据源标识
private String dbName; // 数据库名 (独立库模式)
private String schemaName; // Schema 名 (Schema 隔离模式)
private Integer status; // 0-创建中, 1-运行中, 2-已冻结, 3-已删除
private LocalDateTime expireAt; // 过期时间
private LocalDateTime createdAt;
private LocalDateTime updatedAt;
}
/**
* 租户套餐信息
*/
@Data
@TableName("saas_tenant_plan")
public class TenantPlan {
@TableId
private Long id;
private String tenantId;
private String planCode; // 套餐编码
private String planName; // 套餐名称
private Integer maxUsers; // 最大用户数
private Integer maxStorageGb; // 最大存储 (GB)
private Integer maxApiQps; // 最大 API QPS
private Integer maxApps; // 最大应用数
private LocalDateTime startAt; // 生效时间
private LocalDateTime endAt; // 到期时间
private Integer autoRenew; // 自动续费
}3.2 租户路由
3.2.1 多数据源动态切换
java
/**
* 动态数据源路由器
* 根据 TenantContext 中保存的 tenantId,动态切换到对应的数据源
*/
@Component
public class DynamicDataSourceRouter extends AbstractRoutingDataSource {
@Autowired
private DataSourceConfigManager dataSourceConfigManager;
@Override
protected Object determineCurrentLookupKey() {
return TenantContext.getCurrentTenantId();
}
/**
* 根据租户 ID 获取数据源
*/
public DataSource getDataSource(String tenantId) {
// 先从缓存获取
DataSource ds = dataSourceCache.get(tenantId);
if (ds != null) {
return ds;
}
// 查询租户配置,创建数据源
Tenant tenant = tenantMapper.selectById(tenantId);
if (tenant == null) {
throw new TenantNotFoundException("租户不存在: " + tenantId);
}
ds = createDataSource(tenant);
dataSourceCache.put(tenantId, ds);
return ds;
}
private DataSource createDataSource(Tenant tenant) {
// 根据不同模式创建数据源
DataSourceConfig config = dataSourceConfigManager.getBaseConfig();
switch (tenant.getIsolationMode()) {
case 1: // 独立数据库
config.setUrl(buildJdbcUrl(tenant.getDbName()));
break;
case 2: // Schema 隔离
config.setUrl(buildJdbcUrlWithSchema(tenant.getSchemaName()));
break;
case 3: // 共享表
config.setUrl(buildJdbcUrl("shared_db"));
break;
}
return DataSourceBuilder.create()
.url(config.getUrl())
.username(config.getUsername())
.password(config.getPassword())
.driverClassName(config.getDriverClassName())
.build();
}
}3.2.2 TenantContext 传递
java
/**
* 租户上下文 - 使用 ThreadLocal 存储当前请求的租户信息
*/
public class TenantContext {
private static final ThreadLocal<String> CURRENT_TENANT = new ThreadLocal<>();
private static final ThreadLocal<Tenant> CURRENT_TENANT_INFO = new ThreadLocal<>();
public static void setTenantId(String tenantId) {
CURRENT_TENANT.set(tenantId);
}
public static String getCurrentTenantId() {
return CURRENT_TENANT.get();
}
public static void setTenantInfo(Tenant tenant) {
CURRENT_TENANT_INFO.set(tenant);
}
public static Tenant getCurrentTenantInfo() {
return CURRENT_TENANT_INFO.get();
}
public static void clear() {
CURRENT_TENANT.remove();
CURRENT_TENANT_INFO.remove();
}
}
/**
* 租户拦截器 - 从请求中解析租户信息,写入 TenantContext
*/
@Component
public class TenantInterceptor implements HandlerInterceptor {
@Autowired
private TenantService tenantService;
@Override
public boolean preHandle(HttpServletRequest request,
HttpServletResponse response, Object handler) {
// 从域名或 Header 中获取租户 ID
String tenantId = resolveTenantId(request);
if (tenantId == null) {
throw new TenantNotFoundException("无法识别租户");
}
// 查询租户信息
Tenant tenant = tenantService.getTenant(tenantId);
if (tenant == null || tenant.getStatus() != 1) {
throw new TenantNotFoundException("租户不可用");
}
// 设置租户上下文
TenantContext.setTenantId(tenantId);
TenantContext.setTenantInfo(tenant);
return true;
}
@Override
public void afterCompletion(HttpServletRequest request,
HttpServletResponse response, Object handler, Exception ex) {
TenantContext.clear();
}
/**
* 从请求中解析租户 ID
* 策略: Header > 子域名 > URL 路径
*/
private String resolveTenantId(HttpServletRequest request) {
// 1. 从 Header 获取
String tenantId = request.getHeader("X-Tenant-Id");
if (tenantId != null) {
return tenantId;
}
// 2. 从子域名获取 (tenant.example.com)
String host = request.getHeader("Host");
if (host != null && host.contains(".")) {
String subdomain = host.substring(0, host.indexOf('.'));
if (!"www".equals(subdomain) && !"api".equals(subdomain)) {
return subdomain;
}
}
// 3. 从 JWT Token 中获取
String authHeader = request.getHeader("Authorization");
if (authHeader != null && authHeader.startsWith("Bearer ")) {
String token = authHeader.substring(7);
try {
Claims claims = jwtTokenService.verifyToken(token);
return claims.get("tenant_id", String.class);
} catch (Exception e) {
// ignore
}
}
return null;
}
}3.3 租户初始化
租户初始化流程:
text
创建租户 -> 选择套餐 -> 支付 -> 初始化数据库 -> 配置域名 -> 完成java
/**
* 租户初始化服务
*/
@Service
public class TenantInitializationService {
@Autowired
private TenantMapper tenantMapper;
@Autowired
private DynamicDataSourceRouter dataSourceRouter;
@Autowired
private Flyway flyway;
@Transactional
public Tenant createTenant(CreateTenantRequest request) {
// 1. 创建租户记录
Tenant tenant = new Tenant();
tenant.setTenantId(generateTenantId());
tenant.setTenantName(request.getTenantName());
tenant.setContactName(request.getContactName());
tenant.setContactPhone(request.getContactPhone());
tenant.setPlanCode(request.getPlanCode());
tenant.setIsolationMode(determineIsolationMode(request.getPlanCode()));
tenant.setStatus(0); // 创建中
tenant.setExpireAt(LocalDateTime.now().plusMonths(1)); // 试用期
tenantMapper.insert(tenant);
// 2. 异步初始化租户资源
CompletableFuture.runAsync(() -> initializeTenantResources(tenant))
.exceptionally(ex -> {
log.error("租户初始化失败: tenantId={}", tenant.getTenantId(), ex);
tenant.setStatus(3); // 标记为失败
tenantMapper.updateById(tenant);
return null;
});
return tenant;
}
/**
* 初始化租户资源
*/
public void initializeTenantResources(Tenant tenant) {
try {
// 2.1 初始化数据库
initDatabase(tenant);
// 2.2 初始化 Redis 命名空间
initRedisNamespace(tenant);
// 2.3 初始化存储空间
initStorage(tenant);
// 2.4 配置域名 (可选)
if (tenant.getDomain() != null) {
configureDomain(tenant);
}
// 2.5 开通套餐
activatePlan(tenant);
// 2.6 标记租户为运行状态
tenant.setStatus(1);
tenantMapper.updateById(tenant);
log.info("租户初始化完成: tenantId={}", tenant.getTenantId());
} catch (Exception e) {
log.error("租户初始化异常: tenantId={}", tenant.getTenantId(), e);
throw e;
}
}
private void initDatabase(Tenant tenant) {
if (tenant.getIsolationMode() == 1) {
// 独立数据库模式: 创建新数据库
String dbName = "saas_tenant_" + tenant.getTenantId().replace("-", "_");
dataSourceConfigManager.createDatabase(dbName);
tenant.setDbName(dbName);
// 执行 Flyway 迁移
DataSource ds = dataSourceRouter.getDataSource(tenant.getTenantId());
Flyway.configure()
.dataSource(ds)
.locations("db/migration/tenant")
.baselineOnMigrate(true)
.load()
.migrate();
} else if (tenant.getIsolationMode() == 2) {
// Schema 隔离模式: 创建 Schema
String schemaName = "tenant_" + tenant.getTenantId().replace("-", "_");
dataSourceConfigManager.createSchema(schemaName);
tenant.setSchemaName(schemaName);
}
tenantMapper.updateById(tenant);
}
private void initRedisNamespace(Tenant tenant) {
// 为租户初始化 Redis key 前缀
String prefix = "tenant:" + tenant.getTenantId() + ":";
redisTemplate.opsForValue().set(prefix + "init", "1");
}
private void initStorage(Tenant tenant) {
// 初始化 OSS 存储目录
String bucketName = "saas-tenant-" + tenant.getTenantId();
ossClient.createBucket(bucketName);
}
private void activatePlan(Tenant tenant) {
TenantPlan plan = new TenantPlan();
plan.setTenantId(tenant.getTenantId());
plan.setPlanCode(tenant.getPlanCode());
plan.setPlanName(getPlanName(tenant.getPlanCode()));
plan.setStartAt(LocalDateTime.now());
plan.setEndAt(tenant.getExpireAt());
plan.setAutoRenew(0);
tenantPlanMapper.insert(plan);
}
private int determineIsolationMode(String planCode) {
switch (planCode) {
case "enterprise":
return 1; // 独立数据库
case "professional":
return 2; // Schema 隔离
case "free":
default:
return 3; // 共享表
}
}
}3.4 租户数据隔离实现
3.4.1 MyBatis Plus 拦截器
java
/**
* MyBatis Plus 租户拦截器
* 自动在 SQL 上追加 tenant_id 过滤条件
*/
@Component
public class TenantSqlInterceptor implements InnerInterceptor {
@Override
public void beforeQuery(Executor executor, MappedStatement ms,
Object parameter, RowBounds rowBounds,
ResultHandler resultHandler, BoundSql boundSql) {
String tenantId = TenantContext.getCurrentTenantId();
if (tenantId == null) {
return;
}
// 获取原始 SQL
String sql = boundSql.getSql();
// 判断是否需要追加租户过滤 (通过表名配置)
if (shouldAddTenantFilter(sql)) {
String newSql = addTenantFilter(sql, tenantId);
// 通过反射更新 BoundSql
ReflectionUtil.setFieldValue(boundSql, "sql", newSql);
}
}
@Override
public void beforeUpdate(Executor executor, MappedStatement ms,
Object parameter) {
String tenantId = TenantContext.getCurrentTenantId();
if (tenantId == null) {
return;
}
String sql = boundSql.getSql();
if (shouldAddTenantFilter(sql)) {
String newSql = addTenantFilter(sql, tenantId);
ReflectionUtil.setFieldValue(boundSql, "sql", newSql);
}
}
private String addTenantFilter(String sql, String tenantId) {
// 在 WHERE 条件后追加 AND tenant_id = 'xxx'
if (sql.contains("WHERE")) {
return sql.replace("WHERE", "WHERE tenant_id = '" + tenantId + "' AND ");
} else {
// 没有 WHERE 条件,追加
int lastFromIndex = sql.lastIndexOf("FROM");
int afterFrom = sql.indexOf(" ", lastFromIndex + 4);
String tablePart = sql.substring(lastFromIndex, afterFrom);
return sql.replace(tablePart,
tablePart + " WHERE tenant_id = '" + tenantId + "' ");
}
}
}3.4.2 行级数据隔离
java
/**
* 实体类中的租户字段
* 所有需要隔离的业务表都必须包含 tenant_id 字段
*/
@Data
@TableName("saas_order")
public class TenantOrder {
@TableId
private Long id;
private String tenantId; // 租户 ID (分片键)
private String orderNo;
private BigDecimal amount;
private Integer status;
private LocalDateTime createdAt;
private LocalDateTime updatedAt;
@TableField(exist = false)
private String tenantName; // 非持久化字段
}
/**
* Mapper 层自动填充租户 ID
*/
@Component
public class TenantMetaObjectHandler implements MetaObjectHandler {
@Override
public void insertFill(MetaObject metaObject) {
String tenantId = TenantContext.getCurrentTenantId();
if (tenantId != null) {
this.strictInsertFill(metaObject, "tenantId", String.class, tenantId);
}
this.strictInsertFill(metaObject, "createdAt", LocalDateTime.class, LocalDateTime.now());
this.strictInsertFill(metaObject, "updatedAt", LocalDateTime.class, LocalDateTime.now());
}
@Override
public void updateFill(MetaObject metaObject) {
this.strictUpdateFill(metaObject, "updatedAt", LocalDateTime.class, LocalDateTime.now());
}
}4. 第三方接口统一管理
4.1 第三方接口注册
4.1.1 接口元数据
java
/**
* 第三方接口定义
*/
@Data
@TableName("third_party_api")
public class ThirdPartyApi {
@TableId(type = IdType.AUTO)
private Long id;
private String provider; // 服务商: aliyun/tencent/amap/kdniao
private String apiName; // 接口名称: 短信发送/物流查询/地图编码
private String apiCode; // 接口编码: sms.send/logistics.query
private String urlTemplate; // URL 模板: https://{region}.aliyuncs.com
private String httpMethod; // 请求方法: GET/POST
private String contentType; // Content-Type: application/json
private Integer connectTimeout; // 连接超时 (ms)
private Integer readTimeout; // 读取超时 (ms)
private Integer maxRetries; // 最大重试次数
private Integer retryIntervalMs; // 重试间隔 (ms)
private String circuitBreakerRule; // 熔断规则 JSON
private Integer maxConcurrent; // 最大并发数
private String status; // ENABLED / DISABLED
private LocalDateTime createdAt;
private LocalDateTime updatedAt;
}4.1.2 URL 模板与动态参数
java
/**
* 第三方接口调用客户端
*/
@Component
public class ThirdPartyApiClient {
@Autowired
private ThirdPartyApiMapper apiMapper;
@Autowired
private CredentialManager credentialManager;
/**
* 执行第三方 API 调用
*/
public <T> T execute(String apiCode, ApiRequest request, Class<T> responseType) {
// 1. 获取接口定义
ThirdPartyApi api = getApiDefinition(apiCode);
// 2. 解析 URL 模板
String url = resolveUrlTemplate(api.getUrlTemplate(), request.getPathParams());
// 3. 构建请求
HttpMethod method = HttpMethod.valueOf(api.getHttpMethod());
HttpEntity<String> entity = buildHttpEntity(api, request);
// 4. 执行调用(含重试)
return executeWithRetry(api, () -> {
ResponseEntity<String> response = restTemplate.exchange(
url, method, entity, String.class);
return parseResponse(response, responseType);
});
}
private String resolveUrlTemplate(String template, Map<String, String> pathParams) {
if (pathParams == null || pathParams.isEmpty()) {
return template;
}
UriComponentsBuilder builder = UriComponentsBuilder.fromUriString(template);
for (Map.Entry<String, String> entry : pathParams.entrySet()) {
builder.replaceQueryParam(entry.getKey(), entry.getValue());
}
return builder.build().toUriString();
}
private HttpEntity<String> buildHttpEntity(ThirdPartyApi api, ApiRequest request) {
HttpHeaders headers = new HttpHeaders();
headers.setContentType(MediaType.parseMediaType(api.getContentType()));
// 注入认证信息
Credential credential = credentialManager.getCredential(api.getProvider());
headers.set("Authorization", "Bearer " + credential.getAccessToken());
// 注入自定义 Header
if (request.getHeaders() != null) {
request.getHeaders().forEach(headers::set);
}
String body = request.getBody() != null ? JSON.toJSONString(request.getBody()) : null;
return new HttpEntity<>(body, headers);
}
private <T> T executeWithRetry(ThirdPartyApi api, Supplier<T> supplier) {
int retries = 0;
long interval = api.getRetryIntervalMs();
while (true) {
try {
return supplier.get();
} catch (Exception e) {
retries++;
if (retries > api.getMaxRetries()) {
throw new ThirdPartyException("第三方接口调用失败: " + api.getApiCode(), e);
}
if (!isRetryable(e)) {
throw e;
}
try {
Thread.sleep(interval);
} catch (InterruptedException ie) {
Thread.currentThread().interrupt();
throw new ThirdPartyException("重试中断", ie);
}
interval *= 2; // 指数退避
}
}
}
}4.2 调用治理
4.2.1 Sentinel 熔断降级
java
/**
* 第三方接口 Sentinel 熔断配置
*/
@Configuration
public class ThirdPartySentinelConfig {
@PostConstruct
public void initDegradeRules() {
List<DegradeRule> rules = new ArrayList<>();
// 短信服务熔断规则: 异常比例 > 50% 时熔断 30 秒
DegradeRule smsRule = new DegradeRule("thirdparty:sms")
.setGrade(RuleConstant.DEGRADE_GRADE_EXCEPTION_RATIO)
.setCount(0.5)
.setTimeWindow(30)
.setMinRequestAmount(10);
rules.add(smsRule);
// 物流查询熔断规则: 慢调用比例 > 30% (响应时间 > 3s) 时熔断 15 秒
DegradeRule logisticsRule = new DegradeRule("thirdparty:logistics")
.setGrade(RuleConstant.DEGRADE_GRADE_RT)
.setCount(3000)
.setTimeWindow(15)
.setMinRequestAmount(5);
rules.add(logisticsRule);
// 地图服务熔断规则
DegradeRule mapRule = new DegradeRule("thirdparty:map")
.setGrade(RuleConstant.DEGRADE_GRADE_EXCEPTION_COUNT)
.setCount(5) // 5 分钟内异常数 > 5
.setTimeWindow(20)
.setMinRequestAmount(3);
rules.add(mapRule);
DegradeRuleManager.loadRules(rules);
}
}
/**
* 带熔断的第三方调用
*/
@Service
public class ThirdPartyService {
/**
* 短信发送 - 带熔断保护
*/
public SmsResponse sendSms(SmsRequest request) {
try (Entry entry = SphU.entry("thirdparty:sms")) {
return aliyunSmsClient.send(request);
} catch (BlockException e) {
// 熔断降级: 切到备用短信通道
log.warn("阿里云短信熔断,切换到腾讯云");
return tencentSmsClient.send(request);
}
}
/**
* 物流查询 - 带熔断保护
*/
public TrackResponse queryLogistics(String expCode, String expNo) {
try (Entry entry = SphU.entry("thirdparty:logistics")) {
return kdniaoClient.query(expCode, expNo);
} catch (BlockException e) {
// 熔断降级: 切到快递 100
log.warn("快递鸟熔断,切换到快递 100");
return kuaidi100Client.query(expCode, expNo);
}
}
}4.2.2 隔离舱与线程池隔离
java
/**
* 线程池隔离配置
* 不同第三方服务使用独立的线程池,防止某个服务的慢调用耗尽所有线程
*/
@Configuration
public class ThreadPoolIsolationConfig {
/**
* 短信服务线程池
*/
@Bean("smsThreadPool")
public ThreadPoolTaskExecutor smsThreadPool() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(5);
executor.setMaxPoolSize(10);
executor.setQueueCapacity(50);
executor.setThreadNamePrefix("sms-pool-");
executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
executor.initialize();
return executor;
}
/**
* 物流服务线程池
*/
@Bean("logisticsThreadPool")
public ThreadPoolTaskExecutor logisticsThreadPool() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(3);
executor.setMaxPoolSize(5);
executor.setQueueCapacity(30);
executor.setThreadNamePrefix("logistics-pool-");
executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
executor.initialize();
return executor;
}
/**
* 地图服务线程池
*/
@Bean("mapThreadPool")
public ThreadPoolTaskExecutor mapThreadPool() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(3);
executor.setMaxPoolSize(5);
executor.setQueueCapacity(20);
executor.setThreadNamePrefix("map-pool-");
executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
executor.initialize();
return executor;
}
}
/**
* 使用隔离线程池调用
*/
@Service
public class IsolatedThirdPartyService {
@Resource(name = "smsThreadPool")
private ThreadPoolTaskExecutor smsThreadPool;
@Resource(name = "logisticsThreadPool")
private ThreadPoolTaskExecutor logisticsThreadPool;
public CompletableFuture<SmsResponse> sendSmsAsync(SmsRequest request) {
return CompletableFuture.supplyAsync(() -> {
try (Entry entry = SphU.entry("thirdparty:sms")) {
return aliyunSmsClient.send(request);
} catch (BlockException e) {
return tencentSmsClient.send(request);
}
}, smsThreadPool);
}
public CompletableFuture<TrackResponse> queryLogisticsAsync(String expCode, String expNo) {
return CompletableFuture.supplyAsync(() -> {
try (Entry entry = SphU.entry("thirdparty:logistics")) {
return kdniaoClient.query(expCode, expNo);
} catch (BlockException e) {
return kuaidi100Client.query(expCode, expNo);
}
}, logisticsThreadPool);
}
}4.3 统一监控
4.3.1 指标采集
java
/**
* 第三方调用监控 - Micrometer 指标采集
*/
@Component
public class ThirdPartyMetricsCollector {
private final MeterRegistry meterRegistry;
private final Counter totalCounter;
private final Counter successCounter;
private final Counter failureCounter;
private final Timer responseTimer;
public ThirdPartyMetricsCollector(MeterRegistry meterRegistry) {
this.meterRegistry = meterRegistry;
this.totalCounter = Counter.builder("thirdparty.api.total")
.description("第三方接口总调用量")
.tag("type", "total")
.register(meterRegistry);
this.successCounter = Counter.builder("thirdparty.api.success")
.description("第三方接口成功调用量")
.register(meterRegistry);
this.failureCounter = Counter.builder("thirdparty.api.failure")
.description("第三方接口失败调用量")
.register(meterRegistry);
this.responseTimer = Timer.builder("thirdparty.api.response.time")
.description("第三方接口响应时间")
.publishPercentiles(0.5, 0.9, 0.95, 0.99)
.register(meterRegistry);
}
/**
* 记录调用指标
*/
public void record(String apiCode, boolean success, long durationMs) {
// 增加维度标签
List<Tag> tags = List.of(
Tag.of("api_code", apiCode),
Tag.of("success", String.valueOf(success))
);
totalCounter.increment();
if (success) {
successCounter.increment();
} else {
failureCounter.increment();
}
responseTimer.record(Duration.ofMillis(durationMs));
}
/**
* 指标采集切面
*/
@Around("@annotation(MonitorThirdParty)")
public Object monitor(ProceedingJoinPoint pjp) throws Throwable {
long start = System.currentTimeMillis();
String apiCode = getApiCode(pjp);
boolean success = true;
try {
return pjp.proceed();
} catch (Exception e) {
success = false;
throw e;
} finally {
long duration = System.currentTimeMillis() - start;
record(apiCode, success, duration);
}
}
}Prometheus 告警规则:
yaml
groups:
- name: thirdparty_alerts
rules:
# 成功率告警: 5 分钟内成功率低于 95%
- alert: ThirdPartyApiHighFailureRate
expr: |
rate(thirdparty_api_failure_total[5m])
/
rate(thirdparty_api_total[5m]) > 0.05
for: 3m
labels:
severity: critical
annotations:
summary: "第三方 API 成功率低于 95%"
description: "API {{ $labels.api_code }} 5 分钟成功率 {{ $value | humanizePercentage }}"
# 响应时间告警: P99 响应时间超过 5 秒
- alert: ThirdPartyApiSlowResponse
expr: |
histogram_quantile(0.99,
rate(thirdparty_api_response_time_seconds_bucket[5m])
) > 5
for: 2m
labels:
severity: warning
annotations:
summary: "第三方 API 响应缓慢"
description: "API {{ $labels.api_code }} P99 响应时间 {{ $value }}s"4.4 第三方密钥安全存储
4.4.1 数据库加密存储
java
/**
* 第三方凭证实体
*/
@Data
@TableName("third_party_credential")
public class ThirdPartyCredential {
@TableId(type = IdType.AUTO)
private Long id;
private String provider; // 服务商标识
private String envTag; // dev/staging/prod
@TableField(insertStrategy = FieldStrategy.NOT_NULL)
private String appKey; // 数据库密文存储
@TableField(insertStrategy = FieldStrategy.NOT_NULL)
private String appSecret; // 数据库密文存储
private String extraConfig; // 额外配置 JSON (密文)
private Integer status; // 0-禁用, 1-启用
private LocalDateTime rotateAt; // 上次轮换时间
private LocalDateTime expireAt; // 过期时间
private LocalDateTime createdAt;
private LocalDateTime updatedAt;
}
/**
* AES-256-GCM 加密解密工具
*/
@Component
public class CredentialEncryptor {
@Value("${credential.encrypt.key}")
private String secretKey;
private static final String ALGORITHM = "AES/GCM/NoPadding";
private static final int GCM_IV_LENGTH = 12;
private static final int GCM_TAG_LENGTH = 128;
/**
* 加密
*/
public String encrypt(String plainText) {
try {
byte[] iv = new byte[GCM_IV_LENGTH];
SecureRandom secureRandom = new SecureRandom();
secureRandom.nextBytes(iv);
Cipher cipher = Cipher.getInstance(ALGORITHM);
SecretKeySpec keySpec = new SecretKeySpec(
secretKey.getBytes(StandardCharsets.UTF_8), "AES");
GCMParameterSpec gcmSpec = new GCMParameterSpec(GCM_TAG_LENGTH, iv);
cipher.init(Cipher.ENCRYPT_MODE, keySpec, gcmSpec);
byte[] encrypted = cipher.doFinal(plainText.getBytes(StandardCharsets.UTF_8));
// 将 IV 和密文合并存储
byte[] combined = new byte[GCM_IV_LENGTH + encrypted.length];
System.arraycopy(iv, 0, combined, 0, GCM_IV_LENGTH);
System.arraycopy(encrypted, 0, combined, GCM_IV_LENGTH, encrypted.length);
return Base64.getEncoder().encodeToString(combined);
} catch (Exception e) {
throw new SecurityException("凭证加密失败", e);
}
}
/**
* 解密
*/
public String decrypt(String cipherText) {
try {
byte[] combined = Base64.getDecoder().decode(cipherText);
byte[] iv = new byte[GCM_IV_LENGTH];
byte[] encrypted = new byte[combined.length - GCM_IV_LENGTH];
System.arraycopy(combined, 0, iv, 0, GCM_IV_LENGTH);
System.arraycopy(combined, GCM_IV_LENGTH, encrypted, 0, encrypted.length);
Cipher cipher = Cipher.getInstance(ALGORITHM);
SecretKeySpec keySpec = new SecretKeySpec(
secretKey.getBytes(StandardCharsets.UTF_8), "AES");
GCMParameterSpec gcmSpec = new GCMParameterSpec(GCM_TAG_LENGTH, iv);
cipher.init(Cipher.DECRYPT_MODE, keySpec, gcmSpec);
byte[] decrypted = cipher.doFinal(encrypted);
return new String(decrypted, StandardCharsets.UTF_8);
} catch (Exception e) {
throw new SecurityException("凭证解密失败", e);
}
}
}
/**
* 凭证管理器 - 加解密 + 内存缓存
*/
@Component
public class CredentialManager {
@Autowired
private ThirdPartyCredentialMapper credentialMapper;
@Autowired
private CredentialEncryptor encryptor;
private final Map<String, ThirdPartyCredential> cache = new ConcurrentHashMap<>();
@PostConstruct
public void init() {
loadCredentials();
}
@Scheduled(fixedRate = 60000) // 每分钟刷新一次
public void refresh() {
loadCredentials();
}
private void loadCredentials() {
List<ThirdPartyCredential> credentials = credentialMapper.selectList(
new LambdaQueryWrapper<ThirdPartyCredential>()
.eq(ThirdPartyCredential::getStatus, 1));
for (ThirdPartyCredential cred : credentials) {
// 解密
cred.setAppKey(encryptor.decrypt(cred.getAppKey()));
cred.setAppSecret(encryptor.decrypt(cred.getAppSecret()));
if (cred.getExtraConfig() != null) {
cred.setExtraConfig(encryptor.decrypt(cred.getExtraConfig()));
}
cache.put(buildCacheKey(cred.getProvider(), cred.getEnvTag()), cred);
}
}
public ThirdPartyCredential getCredential(String provider) {
String env = System.getProperty("spring.profiles.active", "dev");
ThirdPartyCredential cred = cache.get(buildCacheKey(provider, env));
if (cred == null) {
throw new BusinessException("未找到第三方凭证: " + provider + "[" + env + "]");
}
return cred;
}
private String buildCacheKey(String provider, String env) {
return provider + ":" + env;
}
}4.4.2 凭据轮换
java
/**
* 凭据轮换定时任务
*/
@Component
public class CredentialRotationScheduler {
@Autowired
private ThirdPartyCredentialMapper credentialMapper;
@Autowired
private CredentialEncryptor encryptor;
/**
* 每天凌晨 2 点检查即将到期的凭证
*/
@Scheduled(cron = "0 0 2 * * ?")
@Transactional
public void rotateExpiringCredentials() {
// 查询 7 天内即将到期的凭证
List<ThirdPartyCredential> expiringList = credentialMapper.selectList(
new LambdaQueryWrapper<ThirdPartyCredential>()
.eq(ThirdPartyCredential::getStatus, 1)
.lt(ThirdPartyCredential::getExpireAt,
LocalDateTime.now().plusDays(7)));
for (ThirdPartyCredential cred : expiringList) {
try {
rotateCredential(cred);
} catch (Exception e) {
log.error("凭据轮换失败: provider={}, id={}",
cred.getProvider(), cred.getId(), e);
alertService.sendAlert("凭据轮换失败",
"provider: " + cred.getProvider());
}
}
}
private void rotateCredential(ThirdPartyCredential oldCred) {
// 1. 通过第三方 API 生成新密钥
RotatedCredential newKeys = rotateViaProviderApi(oldCred);
// 2. 保存新凭证
ThirdPartyCredential newCred = new ThirdPartyCredential();
newCred.setProvider(oldCred.getProvider());
newCred.setEnvTag(oldCred.getEnvTag());
newCred.setAppKey(encryptor.encrypt(newKeys.getAppKey()));
newCred.setAppSecret(encryptor.encrypt(newKeys.getAppSecret()));
newCred.setExtraConfig(oldCred.getExtraConfig());
newCred.setStatus(1);
newCred.setRotateAt(LocalDateTime.now());
newCred.setExpireAt(LocalDateTime.now().plusMonths(3));
credentialMapper.insert(newCred);
// 3. 旧凭证保留 24 小时用于过渡,然后禁用
oldCred.setStatus(0);
credentialMapper.updateById(oldCred);
}
}4.4.3 Vault 集成
yaml
# application.yml - HashiCorp Vault 配置
spring:
cloud:
vault:
host: ${VAULT_HOST:localhost}
port: ${VAULT_PORT:8200}
scheme: https
authentication: TOKEN
token: ${VAULT_TOKEN}
kv:
enabled: true
backend: secret
default-context: saas-platform
application-name: thirdparty-credentials
# Vault 中存储的密钥结构:
# secret/saas-platform/thirdparty-credentials/dev
# {
# "aliyun_sms_access_key": "LTAI...",
# "aliyun_sms_access_secret": "...",
# "amap_api_key": "...",
# "kdniao_api_key": "...",
# "kdniao_api_secret": "..."
# }java
/**
* Vault 密钥加载器
*/
@Component
public class VaultCredentialLoader {
@Autowired
private VaultTemplate vaultTemplate;
private Map<String, Object> credentialCache;
@PostConstruct
public void loadFromVault() {
String path = "secret/saas-platform/thirdparty-credentials/" + getEnv();
VaultResponseSupport<Map<String, Object>> response =
vaultTemplate.read(path, new ParameterizedTypeReference<>() {});
if (response != null && response.getData() != null) {
this.credentialCache = response.getData();
}
}
public String getSecret(String key) {
if (credentialCache == null || !credentialCache.containsKey(key)) {
throw new BusinessException("Vault 中未找到密钥: " + key);
}
return (String) credentialCache.get(key);
}
private String getEnv() {
return System.getProperty("spring.profiles.active", "dev");
}
}5. 开发平台接入规范
5.1 开发者入驻流程
开发者入驻开放平台的标准流程:
text
注册账号 -> 实名认证 -> 创建应用 -> 申请 API 权限 -> 获取 AK/SK -> 开始开发java
/**
* 开发者注册服务
*/
@Service
public class DeveloperRegistrationService {
@Transactional
public Developer registerDeveloper(RegisterRequest request) {
// 1. 校验邮箱/手机唯一性
if (developerMapper.existsByEmail(request.getEmail())) {
throw new BusinessException("该邮箱已被注册");
}
// 2. 创建开发者账号
Developer developer = new Developer();
developer.setDeveloperId(generateDeveloperId());
developer.setEmail(request.getEmail());
developer.setPassword(passwordEncoder.encode(request.getPassword()));
developer.setCompanyName(request.getCompanyName());
developer.setContactName(request.getContactName());
developer.setContactPhone(request.getContactPhone());
developer.setStatus(0); // 待认证
developer.setRegisterTime(LocalDateTime.now());
developerMapper.insert(developer);
// 3. 发送实名认证通知
notificationService.sendVerificationEmail(developer.getEmail());
return developer;
}
/**
* 实名认证审核
*/
@Transactional
public void verifyIdentity(String developerId, IdentityVerifyRequest request) {
Developer developer = developerMapper.selectById(developerId);
developer.setIdentityType(request.getIdentityType()); // COMPANY / PERSONAL
developer.setIdentityNo(encryptor.encrypt(request.getIdentityNo()));
developer.setIdentityFileUrls(String.join(",", request.getFileUrls()));
developer.setStatus(1); // 已认证
developer.setVerifyTime(LocalDateTime.now());
developerMapper.updateById(developer);
// 创建默认应用
createDefaultApp(developer);
}
private void createDefaultApp(Developer developer) {
App app = new App();
app.setAppId(generateAppId());
app.setAppName(developer.getCompanyName() + " 默认应用");
app.setDeveloperId(developer.getDeveloperId());
app.setStatus(AppStatus.DRAFT);
app.setAppType(ISV);
appMapper.insert(app);
}
}
/**
* 应用信息
*/
@Data
@TableName("open_app")
public class App {
@TableId
private String appId;
private String appName;
private String appDescription;
private String appIconUrl;
private String developerId;
private String appType; // ISV / SELF
private String appStatus; // DRAFT
private String apiPermissions; // API 权限列表
private String ipWhitelist; // IP 白名单
private String callbackUrl; // 回调地址
private LocalDateTime createdAt;
private LocalDateTime updatedAt;
}5.2 沙箱环境说明
沙箱环境为开发者提供独立、隔离的测试环境,包含模拟数据和 Mock 接口。
5.2.1 沙箱环境配置
yaml
# application-sandbox.yml
server:
port: 8081
spring:
profiles: sandbox
datasource:
url: jdbc:mysql://sandbox-db:3306/openapi_sandbox
username: sandbox_user
password: ${SANDBOX_DB_PASSWORD}
openapi:
sandbox:
enabled: true
mock-response: true
simulate-delay: true
min-delay-ms: 50
max-delay-ms: 500
# 模拟数据配置
mock-data:
orders:
count: 1000
amount-range: [0.01, 9999.99]
users:
count: 500
payments:
success-rate: 0.9 # 90% 成功率5.2.2 Mock 接口实现
java
/**
* 沙箱 Mock 控制器
* 在沙箱环境中自动启用,返回模拟数据
*/
@RestController
@Profile("sandbox")
@RequestMapping("/v1")
public class SandboxMockController {
@PostMapping("/orders")
public Response<CreateOrderResponse> createOrder(
@Valid @RequestBody CreateOrderRequest request) {
// 模拟延迟
simulateDelay();
// 模拟随机失败
if (randomFailure(0.05)) {
return Response.error("SYSTEM_BUSY", "系统繁忙,请稍后重试");
}
// 返回模拟响应
CreateOrderResponse response = new CreateOrderResponse();
response.setOrderId(generateMockOrderId());
response.setOrderNo("SANDBOX_" + System.currentTimeMillis());
response.setAmount(request.getAmount());
response.setStatus("CREATED");
response.setCreatedAt(LocalDateTime.now());
return Response.success(response);
}
@GetMapping("/orders/{orderId}")
public Response<OrderVO> getOrder(@PathVariable String orderId) {
simulateDelay();
// 返回模拟订单数据
OrderVO order = new OrderVO();
order.setOrderId(orderId);
order.setOrderNo("SANDBOX_" + orderId);
order.setAmount(new BigDecimal("99.99"));
order.setStatus("PAID");
order.setCreatedAt(LocalDateTime.now().minusHours(2));
return Response.success(order);
}
private void simulateDelay() {
try {
int delay = ThreadLocalRandom.current().nextInt(50, 300);
Thread.sleep(delay);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
private boolean randomFailure(double rate) {
return ThreadLocalRandom.current().nextDouble() < rate;
}
}5.2.3 沙箱测试账号
| 账号类型 | 账号 | 说明 |
|---|---|---|
| 开发者账号 | sandbox_dev@example.com | 预注册的开发者账号,已实名认证 |
| 商家账号 | sandbox_merchant@example.com | 预置商家账号,包含店铺和商品数据 |
| 管理员账号 | sandbox_admin@example.com | 预置管理员账号,用于审核测试 |
5.3 API 调用签名流程
完整的 API 调用签名流程:
text
客户端 服务端
| |
| 1. 构造请求参数 (body) |
| 2. 获取当前时间戳 timestamp |
| 3. 计算签名: |
| signature = HMAC-SHA256( |
| secret, timestamp + "|" + JSON(body) |
| ) |
| 4. 设置 Header: |
| X-Access-Key: AK |
| X-Signature: signature |
| X-Timestamp: timestamp |
| X-Nonce: uuid (可选, 防重放) |
|----------------------------------------------->
| |
| 5. 校验时间戳 (5 分钟内有效) |
| 6. 查询 AK 对应的 SK |
| 7. 重新计算签名并比对 |
| 8. 校验 Nonce 是否已使用 (可选) |
| 9. 校验 API 权限 |
| 10. 执行业务逻辑 |
|<-----------------------------------------------|
| 11. 返回响应数据 |5.4 SDK 与示例代码
5.4.1 Java SDK
java
/**
* 开放平台 Java SDK 使用示例
*/
public class OpenApiSdkExample {
public static void main(String[] args) {
// 初始化客户端
OpenApiClient client = OpenApiClient.newBuilder()
.accessKey("your-access-key")
.accessSecret("your-access-secret")
.baseUrl("https://api.example.com")
.connectTimeout(10, TimeUnit.SECONDS)
.readTimeout(30, TimeUnit.SECONDS)
.build();
// 创建订单
CreateOrderRequest request = new CreateOrderRequest();
request.setOutOrderNo("OUT_" + System.currentTimeMillis());
request.setProductId("PROD_001");
request.setQuantity(2);
request.setAmount(new BigDecimal("199.98"));
request.setNotifyUrl("https://your-domain.com/callback");
CreateOrderResponse response = client.call(
HttpMethod.POST, "/v1/orders", request, CreateOrderResponse.class);
System.out.println("订单创建成功: " + response.getOrderId());
// 查询订单
OrderVO order = client.call(
HttpMethod.GET, "/v1/orders/" + response.getOrderId(), null, OrderVO.class);
System.out.println("订单状态: " + order.getStatus());
}
}5.4.2 Python SDK
python
"""
开放平台 Python SDK 使用示例
"""
import hashlib
import hmac
import json
import time
import requests
class OpenApiClient:
def __init__(self, access_key, access_secret, base_url):
self.access_key = access_key
self.access_secret = access_secret
self.base_url = base_url.rstrip('/')
def _sign(self, timestamp, body):
data = f"{timestamp}|{json.dumps(body, ensure_ascii=False)}"
sign = hmac.new(
self.access_secret.encode('utf-8'),
data.encode('utf-8'),
hashlib.sha256
).digest()
return base64.b64encode(sign).decode('utf-8')
def call(self, method, path, body=None):
timestamp = int(time.time())
signature = self._sign(timestamp, body)
headers = {
'X-Access-Key': self.access_key,
'X-Signature': signature,
'X-Timestamp': str(timestamp),
'Content-Type': 'application/json'
}
url = self.base_url + path
if method == 'GET':
resp = requests.get(url, headers=headers)
else:
resp = requests.post(url, headers=headers, json=body)
if resp.status_code != 200:
raise Exception(f"API 调用失败: {resp.status_code} {resp.text}")
return resp.json()
# 使用示例
client = OpenApiClient(
access_key='your-access-key',
access_secret='your-access-secret',
base_url='https://api.example.com'
)
# 创建订单
order = client.call('POST', '/v1/orders', {
'out_order_no': 'OUT_' + str(int(time.time())),
'product_id': 'PROD_001',
'quantity': 2,
'amount': 199.98,
'notify_url': 'https://your-domain.com/callback'
})
print(f"订单创建成功: {order['data']['order_id']}")5.4.3 JavaScript SDK
javascript
/**
* 开放平台 JavaScript SDK
*/
class OpenApiClient {
constructor({ accessKey, accessSecret, baseUrl }) {
this.accessKey = accessKey;
this.accessSecret = accessSecret;
this.baseUrl = baseUrl.replace(/\/$/, '');
}
async _sign(timestamp, body) {
const data = `${timestamp}|${JSON.stringify(body || {})}`;
const encoder = new TextEncoder();
const key = await crypto.subtle.importKey(
'raw', encoder.encode(this.accessSecret),
{ name: 'HMAC', hash: 'SHA-256' },
false, ['sign']
);
const signature = await crypto.subtle.sign(
'HMAC', key, encoder.encode(data)
);
return btoa(String.fromCharCode(...new Uint8Array(signature)));
}
async call(method, path, body = null) {
const timestamp = Math.floor(Date.now() / 1000);
const signature = await this._sign(timestamp, body);
const headers = {
'X-Access-Key': this.accessKey,
'X-Signature': signature,
'X-Timestamp': timestamp.toString(),
'Content-Type': 'application/json'
};
const url = this.baseUrl + path;
const options = { method, headers };
if (body && method !== 'GET') {
options.body = JSON.stringify(body);
}
const response = await fetch(url, options);
if (!response.ok) {
throw new Error(`API 调用失败: ${response.status}`);
}
return response.json();
}
}
// 使用示例
const client = new OpenApiClient({
accessKey: 'your-access-key',
accessSecret: 'your-access-secret',
baseUrl: 'https://api.example.com'
});
const order = await client.call('POST', '/v1/orders', {
out_order_no: 'OUT_' + Date.now(),
product_id: 'PROD_001',
quantity: 2,
amount: 199.98,
notify_url: 'https://your-domain.com/callback'
});
console.log('订单创建成功:', order.data.order_id);5.5 OpenAPI 规范文档模板
yaml
# openapi.yaml - SaaS 开放平台 API 规范模板
openapi: 3.0.0
info:
title: SaaS 开放平台 API
version: v1
description: |
SaaS 开放平台提供标准化的 API 接口,支持第三方开发者集成。
所有 API 使用 AK/SK HMAC-SHA256 签名鉴权,请求时需携带以下 Header:
- X-Access-Key: 开发者 Access Key
- X-Signature: 请求签名
- X-Timestamp: 请求时间戳(Unix 秒级)
contact:
name: SaaS 平台技术支持
email: dev-support@example.com
termsOfService: https://example.com/terms
servers:
- url: https://api.example.com
description: 生产环境
- url: https://sandbox-api.example.com
description: 沙箱环境
components:
securitySchemes:
AccessKeyAuth:
type: apiKey
in: header
name: X-Access-Key
description: AK/SK 鉴权方式,需同时提供 X-Signature 和 X-Timestamp
JWTAuth:
type: http
scheme: bearer
bearerFormat: JWT
description: JWT Token 鉴权,适用于服务端到服务端
schemas:
ErrorResponse:
type: object
properties:
code:
type: string
description: 错误码
message:
type: string
description: 错误描述
request_id:
type: string
description: 请求 ID(用于追踪)
required:
- code
- message
PaginationRequest:
type: object
properties:
page_no:
type: integer
default: 1
minimum: 1
page_size:
type: integer
default: 20
maximum: 100
minimum: 1
PaginationResponse:
type: object
properties:
total:
type: integer
description: 总记录数
page_no:
type: integer
description: 当前页码
page_size:
type: integer
description: 每页条数
total_pages:
type: integer
description: 总页数
responses:
UnauthorizedError:
description: 鉴权失败
content:
application/json:
schema:
$ref: '#/components/schemas/ErrorResponse'
RateLimitError:
description: 请求频率超限
content:
application/json:
schema:
$ref: '#/components/schemas/ErrorResponse'
ServerError:
description: 服务器内部错误
content:
application/json:
schema:
$ref: '#/components/schemas/ErrorResponse'
paths:
/v1/orders:
post:
summary: 创建订单
description: 第三方开发者通过此接口创建交易订单
tags:
- 订单管理
security:
- AccessKeyAuth: []
parameters:
- name: X-Signature
in: header
required: true
schema:
type: string
description: 请求签名
- name: X-Timestamp
in: header
required: true
schema:
type: integer
description: 请求时间戳
- name: X-Nonce
in: header
required: false
schema:
type: string
description: 随机字符串,用于防重放
requestBody:
required: true
content:
application/json:
schema:
type: object
required:
- out_order_no
- product_id
- quantity
- amount
properties:
out_order_no:
type: string
description: 外部订单号(开发者侧唯一)
maxLength: 64
product_id:
type: string
description: 商品 ID
quantity:
type: integer
description: 购买数量
minimum: 1
maximum: 999
amount:
type: number
format: double
description: 订单金额
notify_url:
type: string
format: uri
description: 异步通知回调地址
responses:
'200':
description: 创建成功
content:
application/json:
schema:
type: object
properties:
code:
type: string
example: SUCCESS
data:
type: object
properties:
order_id:
type: string
description: 平台订单 ID
order_no:
type: string
description: 平台订单号
amount:
type: number
format: double
status:
type: string
enum: [CREATED, PAID, SHIPPED, COMPLETED, CLOSED]
created_at:
type: string
format: date-time
'400':
$ref: '#/components/responses/UnauthorizedError'
'429':
$ref: '#/components/responses/RateLimitError'
'500':
$ref: '#/components/responses/ServerError'
/v1/orders/{order_id}:
get:
summary: 查询订单详情
tags:
- 订单管理
security:
- AccessKeyAuth: []
parameters:
- name: order_id
in: path
required: true
schema:
type: string
responses:
'200':
description: 查询成功
content:
application/json:
schema:
type: object
properties:
code:
type: string
data:
$ref: '#/components/schemas/Order'
/v1/orders/list:
get:
summary: 查询订单列表
tags:
- 订单管理
security:
- AccessKeyAuth: []
parameters:
- name: status
in: query
required: false
schema:
type: string
enum: [CREATED, PAID, SHIPPED, COMPLETED, CLOSED]
- name: start_time
in: query
required: false
schema:
type: string
format: date-time
- name: end_time
in: query
required: false
schema:
type: string
format: date-time
- $ref: '#/components/schemas/PaginationRequest'
responses:
'200':
description: 查询成功
content:
application/json:
schema:
type: object
properties:
code:
type: string
data:
type: object
properties:
items:
type: array
items:
$ref: '#/components/schemas/Order'
pagination:
$ref: '#/components/schemas/PaginationResponse'
---
## 6. 计量计费
### 6.1 API 调用次数计量
#### 6.1.1 计量采集
```java
/**
* API 调用计量采集
*/
@Component
public class ApiMeteringCollector {
@Autowired
private RedisTemplate<String, String> redisTemplate;
/**
* 记录 API 调用
*/
public void recordCall(String appId, String apiCode, boolean success) {
String date = LocalDate.now().toString();
// 1. 实时计数(用于仪表盘)
String realtimeKey = "metering:realtime:" + appId + ":" + apiCode;
redisTemplate.opsForValue().increment(realtimeKey);
// 2. 日累计(用于计费结算)
String dailyKey = "metering:daily:" + date + ":" + appId + ":" + apiCode;
redisTemplate.opsForValue().increment(dailyKey);
redisTemplate.expire(dailyKey, Duration.ofDays(90));
// 3. 成功/失败分别计数
String statusKey = "metering:status:" + date + ":" + appId + ":" + (success ? "success" : "failure");
redisTemplate.opsForValue().increment(statusKey);
redisTemplate.expire(statusKey, Duration.ofDays(90));
}
/**
* 批量持久化(定时任务,每 5 分钟执行一次)
*/
@Scheduled(fixedRate = 300000)
@Transactional
public void batchPersist() {
// 从 Redis 读取计量数据,批量写入数据库
Set<String> keys = redisTemplate.keys("metering:daily:*");
List<ApiMeteringRecord> records = new ArrayList<>();
for (String key : keys) {
String count = redisTemplate.opsForValue().get(key);
if (count == null || "0".equals(count)) continue;
// 解析 key: metering:daily:2024-01-01:app-123:order.create
String[] parts = key.split(":");
ApiMeteringRecord record = new ApiMeteringRecord();
record.setRecordDate(LocalDate.parse(parts[2]));
record.setAppId(parts[3]);
record.setApiCode(parts[4]);
record.setCallCount(Long.parseLong(count));
records.add(record);
}
if (!records.isEmpty()) {
meteringRecordMapper.batchInsertOrUpdate(records);
// 清除已持久化的 Redis 计数
redisTemplate.delete(keys);
}
}
}
/**
* API 调用计量记录
*/
@Data
@TableName("open_api_metering")
public class ApiMeteringRecord {
@TableId(type = IdType.AUTO)
private Long id;
private String appId; // 应用 ID
private String apiCode; // API 编码
private LocalDate recordDate; // 记录日期
private Long callCount; // 调用次数
private Long successCount; // 成功次数
private Long failureCount; // 失败次数
private LocalDateTime createdAt;
}6.1.2 用量查询
java
/**
* 用量查询服务
*/
@Service
public class UsageQueryService {
/**
* 查询应用当前周期用量
*/
public UsageVO getCurrentUsage(String appId) {
String date = LocalDate.now().toString();
String planCode = appPlanMapper.getPlanCode(appId);
// 从 Redis 查询实时用量
long totalCalls = 0;
Set<String> keys = redisTemplate.keys("metering:daily:" + date + ":" + appId + ":*");
for (String key : keys) {
String count = redisTemplate.opsForValue().get(key);
totalCalls += Long.parseLong(count);
}
// 查询套餐限额
PlanLimit planLimit = planLimitMapper.getByPlanCode(planCode);
UsageVO usage = new UsageVO();
usage.setAppId(appId);
usage.setPlanCode(planCode);
usage.setTotalCalls(totalCalls);
usage.setApiLimit(planLimit.getMaxApiCalls());
usage.setUsagePercent((double) totalCalls / planLimit.getMaxApiCalls() * 100);
usage.setResetDate(calculateResetDate(planCode));
return usage;
}
}6.2 套餐管理
6.2.1 套餐定义
| 套餐类型 | 月 API 调用上限 | 单 API QPS | 数据隔离级别 | 技术支持 | 价格参考 |
|---|---|---|---|---|---|
| 免费版 | 10,000 次 | 10 | 共享表 | 社区支持 | 免费 |
| 专业版 | 100,000 次 | 100 | Schema 隔离 | 工单支持 | 999 元/月 |
| 企业版 | 1,000,000 次 | 1000 | 独立数据库 | 专属技术支持 | 9999 元/月 |
java
/**
* 套餐定义
*/
@Data
@TableName("open_plan_definition")
public class PlanDefinition {
@TableId
private String planCode; // free / professional / enterprise
private String planName;
private String description;
private BigDecimal monthlyPrice;
private Long maxApiCalls; // 月 API 调用上限
private Integer maxQps; // 单应用 QPS
private Integer maxApps; // 最大应用数
private Integer isolationLevel; // 数据隔离级别
private Integer maxStorageGb; // 最大存储
private String supportLevel; // 技术支持级别
private Integer status; // 0-下架, 1-上架
}
/**
* 套餐订购服务
*/
@Service
public class PlanSubscriptionService {
@Transactional
public Subscription subscribePlan(String appId, String planCode) {
// 1. 校验套餐可用性
PlanDefinition plan = planMapper.selectById(planCode);
if (plan == null || plan.getStatus() != 1) {
throw new BusinessException("套餐不可用");
}
// 2. 创建订购记录
Subscription subscription = new Subscription();
subscription.setSubscriptionNo(generateSubscriptionNo());
subscription.setAppId(appId);
subscription.setPlanCode(planCode);
subscription.setStatus(SubscriptionStatus.ACTIVE);
subscription.setStartTime(LocalDateTime.now());
subscription.setEndTime(LocalDateTime.now().plusMonths(1));
subscription.setAutoRenew(true);
subscriptionMapper.insert(subscription);
// 3. 加载对应的限流规则
rateLimitRuleService.loadRateLimitRules(appId, planCode);
// 4. 触发套餐变更通知
notificationService.notifyAppOwner(appId, "套餐变更",
"您的应用已升级至 " + plan.getPlanName() + " 套餐");
return subscription;
}
/**
* 自动续费检查
*/
@Scheduled(cron = "0 0 10 * * ?") // 每天上午 10 点
public void checkAutoRenew() {
List<Subscription> expiringSubs = subscriptionMapper.selectExpiringSoon(
LocalDateTime.now(), LocalDateTime.now().plusDays(3));
for (Subscription sub : expiringSubs) {
if (sub.getAutoRenew()) {
try {
paymentService.charge(sub);
sub.setStartTime(sub.getEndTime());
sub.setEndTime(sub.getEndTime().plusMonths(1));
subscriptionMapper.updateById(sub);
} catch (Exception e) {
log.error("自动续费失败: {}", sub.getSubscriptionNo(), e);
notificationService.notifyAppOwner(sub.getAppId(),
"续费失败", "自动续费失败,请手动续费");
}
}
}
}
}6.3 限流策略与套餐绑定
java
/**
* 套餐限流规则配置
*/
@Component
public class PlanRateLimitBinder {
/**
* 获取套餐对应的限流配置
*/
public RateLimitConfig getRateLimitConfig(String planCode) {
switch (planCode) {
case "free":
return RateLimitConfig.builder()
.appQps(10)
.dailyCallLimit(10000)
.concurrentLimit(2)
.build();
case "professional":
return RateLimitConfig.builder()
.appQps(100)
.dailyCallLimit(100000)
.concurrentLimit(20)
.build();
case "enterprise":
return RateLimitConfig.builder()
.appQps(1000)
.dailyCallLimit(1000000)
.concurrentLimit(100)
.build();
default:
throw new BusinessException("未知套餐: " + planCode);
}
}
/**
* 同步套餐限流规则到网关
*/
public void syncRateLimitToGateway(String appId, String planCode) {
RateLimitConfig config = getRateLimitConfig(planCode);
// 通过配置中心推送限流规则到 API 网关
String ruleJson = JSON.toJSONString(Map.of(
"appId", appId,
"qps", config.getAppQps(),
"dailyLimit", config.getDailyCallLimit(),
"concurrentLimit", config.getConcurrentLimit()
));
configCenter.publish("gateway:ratelimit:" + appId, ruleJson);
}
}
@Data
@Builder
public class RateLimitConfig {
private int appQps;
private long dailyCallLimit;
private int concurrentLimit;
}6.4 用量预警与自动停服
6.4.1 用量预警
java
/**
* 用量预警服务
*/
@Component
public class UsageAlertService {
@Autowired
private NotificationService notificationService;
/**
* 每小时检查一次用量情况
*/
@Scheduled(cron = "0 0 * * * ?")
public void checkUsageAndAlert() {
List<Subscription> activeSubs = subscriptionMapper.selectActive();
for (Subscription sub : activeSubs) {
PlanDefinition plan = planMapper.selectById(sub.getPlanCode());
long usedCalls = meteringRecordMapper.sumCallsByAppId(
sub.getAppId(), getCurrentCycleStart(sub));
double usagePercent = (double) usedCalls / plan.getMaxApiCalls() * 100;
// 80% 预警
if (usagePercent >= 80 && !hasAlerted(sub.getAppId(), 80)) {
notificationService.notifyAppOwner(sub.getAppId(),
"API 用量预警",
String.format("您的 API 用量已达 %.0f%%,请及时升级套餐或等待下个周期重置", usagePercent));
markAlerted(sub.getAppId(), 80);
}
// 95% 紧急预警
if (usagePercent >= 95 && !hasAlerted(sub.getAppId(), 95)) {
notificationService.notifyAppOwner(sub.getAppId(),
"API 用量紧急预警",
String.format("您的 API 用量已达 %.0f%%,即将触发停服", usagePercent));
markAlerted(sub.getAppId(), 95);
}
// 100% 触发停服
if (usedCalls >= plan.getMaxApiCalls()) {
suspendService(sub.getAppId());
}
}
}
/**
* 停服处理
*/
@Transactional
public void suspendService(String appId) {
// 1. 更新应用状态
App app = appMapper.selectById(appId);
app.setStatus("SUSPENDED");
app.setSuspendReason("API 调用量超限");
app.setSuspendTime(LocalDateTime.now());
appMapper.updateById(app);
// 2. 推送网关规则,拒绝该应用的所有请求
rateLimitRuleService.loadRateLimitRules(appId, "blocked");
// 3. 通知开发者
notificationService.notifyAppOwner(appId,
"服务已暂停",
"您的应用因 API 调用量超限已被暂停服务,请升级套餐或等待下个周期自动恢复");
// 4. 记录审计日志
auditLogService.record("app:suspend", appId,
"API 调用量超限自动停服");
}
/**
* 周期重置时自动恢复
*/
@Scheduled(cron = "0 0 0 1 * ?") // 每月 1 日凌晨
public void autoRestore() {
List<App> suspendedApps = appMapper.selectSuspended();
for (App app : suspendedApps) {
app.setStatus("ENABLED");
app.setSuspendReason(null);
app.setSuspendTime(null);
appMapper.updateById(app);
// 恢复限流规则
String planCode = subscriptionMapper.getPlanCode(app.getAppId());
rateLimitRuleService.loadRateLimitRules(app.getAppId(), planCode);
notificationService.notifyAppOwner(app.getAppId(),
"服务已恢复", "您的应用 API 服务已自动恢复");
}
}
}6.4.2 预警通知渠道
| 通知渠道 | 适用场景 | 优先级 | 配置方式 |
|---|---|---|---|
| 站内信 | 所有开发者 | 默认 | 自动开通 |
| 邮件 | 所有开发者 | 中 | 注册时绑定 |
| 短信 | 紧急预警(95% 以上) | 高 | 开发者中心配置 |
| Webhook | 集成自有监控系统 | 低 | 开发者中心配置回调 URL |
java
/**
* 多通道通知实现
*/
@Component
public class MultiChannelNotifier {
@Autowired
private EmailService emailService;
@Autowired
private SmsService smsService;
@Autowired
private WebhookService webhookService;
public void notify(AlertEvent event) {
// 发送站内信(强制)
sendInternalMessage(event);
// 根据预警级别发送额外通知
switch (event.getLevel()) {
case WARNING:
emailService.sendEmail(event.getEmail(), event.getTitle(), event.getContent());
break;
case CRITICAL:
emailService.sendEmail(event.getEmail(), event.getTitle(), event.getContent());
smsService.sendSms(event.getPhone(), event.getContent());
break;
case INFO:
// 仅站内信
break;
}
// 如果配置了 Webhook,则推送
if (event.getWebhookUrl() != null) {
webhookService.postJson(event.getWebhookUrl(), Map.of(
"event", event.getType(),
"title", event.getTitle(),
"content", event.getContent(),
"timestamp", System.currentTimeMillis()
));
}
}
}