Spring Cloud Alibaba微服务实战:从注册中心到分布式事务
·
引言
微服务架构已成为现代分布式系统的主流选择,而Spring Cloud Alibaba为Java开发者提供了一站式的微服务解决方案。从服务注册发现到分布式事务,从流量控制到消息驱动,本文将带你深入实战,掌握微服务架构的核心技术。
微服务基础架构搭建
1.1 环境准备与依赖配置
<!-- 父POM依赖管理 -->
<dependencyManagement>
<dependencies>
<dependency>
<groupId>com.alibaba.cloud</groupId>
<artifactId>spring-cloud-alibaba-dependencies</artifactId>
<version>2022.0.0.0</version>
<type>pom</type>
<scope>import</scope>
</dependency>
</dependencies>
</dependencyManagement>
<!-- 通用微服务依赖 -->
<dependencies>
<!-- Spring Cloud Alibaba Nacos 服务发现 -->
<dependency>
<groupId>com.alibaba.cloud</groupId>
<artifactId>spring-cloud-starter-alibaba-nacos-discovery</artifactId>
</dependency>
<!-- Spring Cloud Alibaba Nacos 配置中心 -->
<dependency>
<groupId>com.alibaba.cloud</groupId>
<artifactId>spring-cloud-starter-alibaba-nacos-config</artifactId>
</dependency>
<!-- Spring Cloud LoadBalancer -->
<dependency>
<groupId>org.springframework.cloud</groupId>
<artifactId>spring-cloud-starter-loadbalancer</artifactId>
</dependency>
<!-- Spring Boot Actuator 监控 -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-actuator</artifactId>
</dependency>
<!-- Spring Cloud OpenFeign 声明式HTTP客户端 -->
<dependency>
<groupId>org.springframework.cloud</groupId>
<artifactId>spring-cloud-starter-openfeign</artifactId>
</dependency>
</dependencies>
1.2 服务注册与发现
// 用户服务 - 服务提供者
@SpringBootApplication
@EnableDiscoveryClient
public class UserServiceApplication {
public static void main(String[] args) {
SpringApplication.run(UserServiceApplication.class, args);
}
}
// 订单服务 - 服务消费者
@SpringBootApplication
@EnableDiscoveryClient
@EnableFeignClients
public class OrderServiceApplication {
public static void main(String[] args) {
SpringApplication.run(OrderServiceApplication.class, args);
}
}
// 网关服务
@SpringBootApplication
@EnableDiscoveryClient
public class GatewayApplication {
public static void main(String[] args) {
SpringApplication.run(GatewayApplication.class, args);
}
}
// 应用配置文件
@Configuration
public class NacosConfig {
// bootstrap.yml 配置
/**
spring:
application:
name: user-service
cloud:
nacos:
discovery:
server-addr: 192.168.1.100:8848
namespace: dev
group: DEFAULT_GROUP
config:
server-addr: 192.168.1.100:8848
file-extension: yaml
namespace: dev
group: DEFAULT_GROUP
*/
}
Nacos深度实战
2.1 服务注册与发现进阶
// 用户服务实现
@RestController
@Slf4j
@RequestMapping("/api/users")
public class UserController {
@Autowired
private UserService userService;
@GetMapping("/{userId}")
public ResponseEntity<UserDTO> getUser(@PathVariable Long userId) {
log.info("查询用户信息: {}", userId);
UserDTO user = userService.getUserById(userId);
return ResponseEntity.ok(user);
}
@PostMapping
public ResponseEntity<UserDTO> createUser(@RequestBody @Valid CreateUserRequest request) {
log.info("创建用户: {}", request.getUsername());
UserDTO user = userService.createUser(request);
return ResponseEntity.status(HttpStatus.CREATED).body(user);
}
// 健康检查接口
@GetMapping("/health")
public ResponseEntity<Map<String, String>> health() {
Map<String, String> health = new HashMap<>();
health.put("status", "UP");
health.put("timestamp", Instant.now().toString());
return ResponseEntity.ok(health);
}
}
// 服务元数据配置
@Configuration
public class NacosMetadataConfig {
@Bean
@ConditionalOnMissingBean
public NacosDiscoveryProperties nacosDiscoveryProperties() {
NacosDiscoveryProperties properties = new NacosDiscoveryProperties();
// 配置服务元数据
Map<String, String> metadata = new HashMap<>();
metadata.put("version", "1.0.0");
metadata.put("environment", "production");
metadata.put("region", "cn-east-1");
metadata.put("weight", "1.0");
properties.setMetadata(metadata);
properties.setClusterName("DEFAULT_CLUSTER");
properties.setGroup("USER_SERVICE_GROUP");
return properties;
}
// 自定义服务实例配置
@Bean
public NacosServiceManager nacosServiceManager() {
return new NacosServiceManager();
}
}
// 服务发现客户端封装
@Component
@Slf4j
public class ServiceDiscoveryClient {
@Autowired
private DiscoveryClient discoveryClient;
@Autowired
private NacosServiceManager nacosServiceManager;
/**
* 获取指定服务的所有实例
*/
public List<ServiceInstance> getServiceInstances(String serviceId) {
try {
return discoveryClient.getInstances(serviceId);
} catch (Exception e) {
log.error("获取服务实例失败: {}", serviceId, e);
return Collections.emptyList();
}
}
/**
* 获取健康的服务实例
*/
public List<ServiceInstance> getHealthyInstances(String serviceId) {
return getServiceInstances(serviceId).stream()
.filter(instance -> {
// 检查实例的健康状态
Map<String, String> metadata = instance.getMetadata();
return !"false".equals(metadata.get("healthy"));
})
.collect(Collectors.toList());
}
/**
* 根据负载均衡策略选择实例
*/
public ServiceInstance selectInstance(String serviceId, String strategy) {
List<ServiceInstance> instances = getHealthyInstances(serviceId);
if (instances.isEmpty()) {
throw new ServiceNotFoundException("没有可用的服务实例: " + serviceId);
}
switch (strategy) {
case "random":
return instances.get(ThreadLocalRandom.current().nextInt(instances.size()));
case "round-robin":
return roundRobinSelect(serviceId, instances);
case "weighted":
return weightedSelect(instances);
default:
return instances.get(0);
}
}
private ServiceInstance roundRobinSelect(String serviceId, List<ServiceInstance> instances) {
// 简单的轮询实现
AtomicLong counter = roundRobinCounters.computeIfAbsent(serviceId, k -> new AtomicLong(0));
int index = (int) (counter.getAndIncrement() % instances.size());
return instances.get(index);
}
private ServiceInstance weightedSelect(List<ServiceInstance> instances) {
// 基于权重的选择
double totalWeight = instances.stream()
.mapToDouble(instance -> Double.parseDouble(
instance.getMetadata().getOrDefault("weight", "1.0")))
.sum();
double random = ThreadLocalRandom.current().nextDouble(totalWeight);
double current = 0;
for (ServiceInstance instance : instances) {
double weight = Double.parseDouble(
instance.getMetadata().getOrDefault("weight", "1.0"));
current += weight;
if (random <= current) {
return instance;
}
}
return instances.get(0);
}
private final Map<String, AtomicLong> roundRobinCounters = new ConcurrentHashMap<>();
}
// 服务注册生命周期管理
@Component
@Slf4j
public class ServiceLifecycleManager implements ApplicationListener<WebServerInitializedEvent> {
@Autowired
private NacosDiscoveryProperties discoveryProperties;
@Autowired
private NacosServiceManager nacosServiceManager;
@Override
public void onApplicationEvent(WebServerInitializedEvent event) {
// 服务启动完成后的处理
log.info("服务启动完成,注册到Nacos: {}", discoveryProperties.getService());
// 可以在这里执行服务注册后的初始化操作
initializeService();
}
@EventListener
public void onApplicationShutdown(ContextClosedEvent event) {
// 服务关闭时的处理
log.info("服务关闭,从Nacos注销: {}", discoveryProperties.getService());
// 执行清理操作
cleanupService();
}
private void initializeService() {
// 服务初始化逻辑
log.info("初始化服务资源...");
}
private void cleanupService() {
// 服务清理逻辑
log.info("清理服务资源...");
}
}
2.2 配置中心高级功能
// 动态配置管理
@RefreshScope
@Configuration
@Slf4j
public class DynamicConfiguration {
@Value("${app.feature.toggle.new-payment:false}")
private boolean newPaymentFeature;
@Value("${app.rate.limit.max-requests:100}")
private int rateLimitMaxRequests;
@Value("${app.cache.ttl:300}")
private long cacheTtl;
// 配置变更监听
@EventListener
public void onConfigChange(RefreshScopeRefreshedEvent event) {
log.info("配置已刷新: newPaymentFeature={}, rateLimit={}, cacheTtl={}",
newPaymentFeature, rateLimitMaxRequests, cacheTtl);
// 执行配置变更后的逻辑
onConfigurationUpdated();
}
private void onConfigurationUpdated() {
// 重新初始化相关组件
log.info("根据新配置重新初始化组件...");
}
// 配置验证
@PostConstruct
public void validateConfiguration() {
if (rateLimitMaxRequests <= 0) {
throw new IllegalStateException("rateLimitMaxRequests必须大于0");
}
if (cacheTtl <= 0) {
throw new IllegalStateException("cacheTtl必须大于0");
}
}
}
// 配置管理服务
@Service
@Slf4j
public class ConfigManagementService {
@Autowired
private ConfigService configService;
/**
* 发布配置到Nacos
*/
public boolean publishConfig(String dataId, String group, String content) {
try {
boolean result = configService.publishConfig(dataId, group, content);
log.info("发布配置成功: dataId={}, group={}", dataId, group);
return result;
} catch (NacosException e) {
log.error("发布配置失败: dataId={}, group={}", dataId, group, e);
return false;
}
}
/**
* 获取配置内容
*/
public String getConfig(String dataId, String group) {
try {
return configService.getConfig(dataId, group, 5000);
} catch (NacosException e) {
log.error("获取配置失败: dataId={}, group={}", dataId, group, e);
return null;
}
}
/**
* 监听配置变化
*/
public void addConfigListener(String dataId, String group, Listener listener) {
try {
configService.addListener(dataId, group, listener);
log.info("添加配置监听器: dataId={}, group={}", dataId, group);
} catch (NacosException e) {
log.error("添加配置监听器失败: dataId={}, group={}", dataId, group, e);
}
}
/**
* 删除配置
*/
public boolean removeConfig(String dataId, String group) {
try {
boolean result = configService.removeConfig(dataId, group);
log.info("删除配置: dataId={}, group={}, result={}", dataId, group, result);
return result;
} catch (NacosException e) {
log.error("删除配置失败: dataId={}, group={}", dataId, group, e);
return false;
}
}
}
// 配置变更监听器
@Component
@Slf4j
public class CustomConfigListener implements Listener {
private final String dataId;
private final String group;
public CustomConfigListener(String dataId, String group) {
this.dataId = dataId;
this.group = group;
}
@Override
public void receiveConfigInfo(String configInfo) {
log.info("配置发生变化 - dataId: {}, group: {}", dataId, group);
log.info("新配置内容: {}", configInfo);
try {
// 解析配置内容
ObjectMapper mapper = new ObjectMapper();
Map<String, Object> config = mapper.readValue(configInfo, Map.class);
// 处理配置变更
handleConfigChange(config);
} catch (Exception e) {
log.error("处理配置变更失败", e);
}
}
private void handleConfigChange(Map<String, Object> config) {
// 根据配置变更执行相应的操作
if (config.containsKey("featureToggles")) {
@SuppressWarnings("unchecked")
Map<String, Boolean> features = (Map<String, Boolean>) config.get("featureToggles");
updateFeatureToggles(features);
}
if (config.containsKey("rateLimits")) {
@SuppressWarnings("unchecked")
Map<String, Integer> rateLimits = (Map<String, Integer>) config.get("rateLimits");
updateRateLimits(rateLimits);
}
}
private void updateFeatureToggles(Map<String, Boolean> features) {
features.forEach((feature, enabled) -> {
log.info("特性开关更新: {} = {}", feature, enabled);
});
}
private void updateRateLimits(Map<String, Integer> rateLimits) {
rateLimits.forEach((resource, limit) -> {
log.info("限流配置更新: {} = {}", resource, limit);
});
}
@Override
public Executor getExecutor() {
return Executors.newSingleThreadExecutor(r -> {
Thread t = new Thread(r, "config-listener-" + dataId);
t.setDaemon(true);
return t;
});
}
}
OpenFeign声明式HTTP客户端
3.1 Feign客户端高级配置
// 用户服务Feign客户端
@FeignClient(
name = "user-service",
path = "/api/users",
configuration = UserFeignConfig.class,
fallbackFactory = UserServiceFallbackFactory.class
)
public interface UserServiceClient {
@GetMapping("/{userId}")
ResponseEntity<UserDTO> getUserById(@PathVariable("userId") Long userId);
@PostMapping
ResponseEntity<UserDTO> createUser(@RequestBody CreateUserRequest request);
@PutMapping("/{userId}")
ResponseEntity<UserDTO> updateUser(
@PathVariable("userId") Long userId,
@RequestBody UpdateUserRequest request
);
@GetMapping("/search")
ResponseEntity<List<UserDTO>> searchUsers(
@RequestParam("keyword") String keyword,
@RequestParam(value = "page", defaultValue = "0") int page,
@RequestParam(value = "size", defaultValue = "20") int size
);
}
// Feign配置类
@Configuration
@Slf4j
public class UserFeignConfig {
@Bean
public Logger.Level feignLoggerLevel() {
return Logger.Level.FULL;
}
@Bean
public RequestInterceptor authRequestInterceptor() {
return template -> {
// 添加认证头
template.header("Authorization", "Bearer " + getCurrentToken());
// 添加请求ID用于链路追踪
template.header("X-Request-ID", generateRequestId());
};
}
@Bean
public Retryer feignRetryer() {
// 自定义重试策略
return new Retryer.Default(100, 1000, 3);
}
@Bean
public ErrorDecoder feignErrorDecoder() {
return (methodKey, response) -> {
log.error("Feign调用失败: {} {}, 状态码: {}",
response.request().httpMethod(),
response.request().url(),
response.status());
if (response.status() == 404) {
return new UserNotFoundException("用户不存在");
} else if (response.status() == 429) {
return new RateLimitException("请求过于频繁");
} else if (response.status() >= 500) {
return new ServiceUnavailableException("服务暂时不可用");
} else {
return new FeignException.FeignClientException(
response.status(),
response.reason()
);
}
};
}
private String getCurrentToken() {
// 从安全上下文中获取token
return "mock-token";
}
private String generateRequestId() {
return UUID.randomUUID().toString();
}
}
// 降级工厂实现
@Component
@Slf4j
public class UserServiceFallbackFactory implements FallbackFactory<UserServiceClient> {
@Override
public UserServiceClient create(Throwable cause) {
log.warn("UserService服务降级,原因: {}", cause.getMessage());
return new UserServiceClientFallback(cause);
}
private static class UserServiceClientFallback implements UserServiceClient {
private final Throwable cause;
public UserServiceClientFallback(Throwable cause) {
this.cause = cause;
}
@Override
public ResponseEntity<UserDTO> getUserById(Long userId) {
log.warn(" getUserById降级,返回默认用户");
// 返回默认用户或缓存数据
UserDTO defaultUser = new UserDTO();
defaultUser.setId(userId);
defaultUser.setUsername("default-user");
defaultUser.setEmail("default@example.com");
return ResponseEntity.ok(defaultUser);
}
@Override
public ResponseEntity<UserDTO> createUser(CreateUserRequest request) {
log.error("createUser服务不可用", cause);
throw new ServiceUnavailableException("用户服务暂时不可用", cause);
}
@Override
public ResponseEntity<UserDTO> updateUser(Long userId, UpdateUserRequest request) {
log.error("updateUser服务不可用", cause);
throw new ServiceUnavailableException("用户服务暂时不可用", cause);
}
@Override
public ResponseEntity<List<UserDTO>> searchUsers(String keyword, int page, int size) {
log.warn("searchUsers降级,返回空列表");
return ResponseEntity.ok(Collections.emptyList());
}
}
}
// Feign调用封装服务
@Service
@Slf4j
public class FeignInvocationService {
@Autowired
private UserServiceClient userServiceClient;
/**
* 安全的Feign调用封装
*/
public <T> T executeWithFallback(Supplier<T> feignCall, Supplier<T> fallback) {
try {
return feignCall.get();
} catch (FeignException e) {
log.error("Feign调用失败: {}", e.getMessage());
return fallback.get();
} catch (Exception e) {
log.error("Feign调用异常", e);
return fallback.get();
}
}
/**
* 带重试的Feign调用
*/
public <T> T executeWithRetry(Supplier<T> feignCall, int maxRetries) {
int attempts = 0;
while (attempts <= maxRetries) {
try {
return feignCall.get();
} catch (FeignException e) {
attempts++;
if (attempts > maxRetries) {
throw e;
}
log.warn("Feign调用失败,进行第{}次重试", attempts);
exponentialBackoff(attempts);
}
}
throw new RuntimeException("重试次数耗尽");
}
private void exponentialBackoff(int attempt) {
try {
long delay = Math.min(1000 * (1L << (attempt - 1)), 30000L);
Thread.sleep(delay);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new RuntimeException("重试被中断", e);
}
}
/**
* 批量Feign调用
*/
public <T> List<T> executeBatch(List<Supplier<T>> feignCalls) {
return feignCalls.parallelStream()
.map(call -> executeWithFallback(call, () -> null))
.filter(Objects::nonNull)
.collect(Collectors.toList());
}
}
Sentinel流量控制与熔断降级
4.1 流量控制实战
// Sentinel配置类
@Configuration
@Slf4j
public class SentinelConfiguration {
@PostConstruct
public void init() {
// 初始化Sentinel配置
initFlowRules();
initDegradeRules();
initSystemRules();
initAuthorityRules();
}
private void initFlowRules() {
List<FlowRule> rules = new ArrayList<>();
// 用户查询接口限流
FlowRule userQueryRule = new FlowRule("GET:/api/users/{userId}");
userQueryRule.setCount(100); // 每秒100个请求
userQueryRule.setGrade(RuleConstant.FLOW_GRADE_QPS);
userQueryRule.setLimitApp(RuleConstant.LIMIT_APP_DEFAULT);
rules.add(userQueryRule);
// 用户创建接口限流
FlowRule userCreateRule = new FlowRule("POST:/api/users");
userCreateRule.setCount(50); // 每秒50个请求
userCreateRule.setGrade(RuleConstant.FLOW_GRADE_QPS);
userCreateRule.setLimitApp(RuleConstant.LIMIT_APP_DEFAULT);
rules.add(userCreateRule);
FlowRuleManager.loadRules(rules);
log.info("初始化流量控制规则: {}条", rules.size());
}
private void initDegradeRules() {
List<DegradeRule> rules = new ArrayList<>();
// 用户服务降级规则
DegradeRule userServiceRule = new DegradeRule("userService");
userServiceRule.setGrade(RuleConstant.DEGRADE_GRADE_EXCEPTION_COUNT);
userServiceRule.setCount(5); // 5个异常
userServiceRule.setTimeWindow(10); // 10秒
rules.add(userServiceRule);
DegradeRuleManager.loadRules(rules);
log.info("初始化熔断降级规则: {}条", rules.size());
}
private void initSystemRules() {
List<SystemRule> rules = new ArrayList<>();
// 系统保护规则
SystemRule loadRule = new SystemRule();
loadRule.setHighestSystemLoad(4.0); // 最大系统负载
rules.add(loadRule);
SystemRuleManager.loadRules(rules);
log.info("初始化系统保护规则: {}条", rules.size());
}
private void initAuthorityRules() {
List<AuthorityRule> rules = new ArrayList<>();
// 授权规则
AuthorityRule authRule = new AuthorityRule();
authRule.setResource("adminApi");
authRule.setLimitApp("192.168.1.100");
authRule.setStrategy(RuleConstant.AUTHORITY_WHITE);
rules.add(authRule);
AuthorityRuleManager.loadRules(rules);
log.info("初始化授权规则: {}条", rules.size());
}
}
// Sentinel资源注解使用
@Service
@Slf4j
public class UserServiceWithSentinel {
/**
* 使用@SentinelResource注解定义资源
*/
@SentinelResource(
value = "getUserById",
blockHandler = "handleBlockException",
fallback = "getUserFallback",
blockHandlerClass = {UserServiceBlockHandler.class},
exceptionsToIgnore = {IllegalArgumentException.class}
)
public UserDTO getUserById(Long userId) {
log.info("查询用户: {}", userId);
if (userId == null || userId <= 0) {
throw new IllegalArgumentException("用户ID无效");
}
// 模拟业务逻辑
if (userId == 999L) {
throw new RuntimeException("模拟业务异常");
}
UserDTO user = new UserDTO();
user.setId(userId);
user.setUsername("user-" + userId);
user.setEmail("user" + userId + "@example.com");
return user;
}
/**
* 热点参数限流
*/
@SentinelResource(
value = "searchUsers",
blockHandler = "searchUsersBlockHandler"
)
public List<UserDTO> searchUsers(String keyword, Integer page, Integer size) {
log.info("搜索用户: keyword={}, page={}, size={}", keyword, page, size);
// 模拟搜索逻辑
List<UserDTO> users = new ArrayList<>();
for (int i = 0; i < 10; i++) {
UserDTO user = new UserDTO();
user.setId((long) i);
user.setUsername(keyword + "-user-" + i);
user.setEmail(keyword + i + "@example.com");
users.add(user);
}
return users;
}
// Fallback方法
public UserDTO getUserFallback(Long userId, Throwable ex) {
log.warn("getUserFallback被调用, userId: {}, 异常: {}", userId, ex.getMessage());
UserDTO fallbackUser = new UserDTO();
fallbackUser.setId(userId);
fallbackUser.setUsername("fallback-user");
fallbackUser.setEmail("fallback@example.com");
fallbackUser.setFallback(true);
return fallbackUser;
}
}
// 阻塞异常处理器
@Slf4j
public class UserServiceBlockHandler {
/**
* 流控阻塞处理
*/
public static UserDTO handleBlockException(Long userId, BlockException ex) {
log.warn("用户服务流控阻塞: userId={}, rule={}", userId, ex.getRule());
UserDTO blockedUser = new UserDTO();
blockedUser.setId(userId);
blockedUser.setUsername("blocked-user");
blockedUser.setEmail("blocked@example.com");
blockedUser.setBlocked(true);
return blockedUser;
}
/**
* 搜索用户流控处理
*/
public static List<UserDTO> searchUsersBlockHandler(
String keyword, Integer page, Integer size, BlockException ex) {
log.warn("搜索用户流控阻塞: keyword={}, rule={}", keyword, ex.getRule());
return Collections.emptyList();
}
}
// Sentinel监控数据收集
@Component
@Slf4j
public class SentinelMetricsCollector {
@Autowired
private MeterRegistry meterRegistry;
@PostConstruct
public void init() {
startMetricsCollection();
}
private void startMetricsCollection() {
ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor(
r -> new Thread(r, "sentinel-metrics-collector"));
scheduler.scheduleAtFixedRate(this::collectMetrics, 0, 10, TimeUnit.SECONDS);
}
private void collectMetrics() {
try {
// 收集流控规则指标
collectFlowMetrics();
// 收集熔断规则指标
collectDegradeMetrics();
// 收集系统指标
collectSystemMetrics();
} catch (Exception e) {
log.error("收集Sentinel指标失败", e);
}
}
private void collectFlowMetrics() {
List<FlowRule> rules = FlowRuleManager.getRules();
Gauge.builder("sentinel.flow.rules.count")
.description("流控规则数量")
.register(meterRegistry, rules::size);
// 可以收集更详细的指标,如每个规则的通过QPS、阻塞QPS等
}
private void collectDegradeMetrics() {
List<DegradeRule> rules = DegradeRuleManager.getRules();
Gauge.builder("sentinel.degrade.rules.count")
.description("熔断规则数量")
.register(meterRegistry, rules::size);
}
private void collectSystemMetrics() {
double systemLoad = ManagementFactory.getOperatingSystemMXBean().getSystemLoadAverage();
Gauge.builder("system.load.average")
.description("系统负载")
.register(meterRegistry, () -> systemLoad);
}
}
Seata分布式事务实战
5.1 AT模式分布式事务
// 订单服务 - 分布式事务发起方
@Service
@Slf4j
public class OrderService {
@Autowired
private OrderMapper orderMapper;
@Autowired
private AccountServiceClient accountServiceClient;
@Autowired
private InventoryServiceClient inventoryServiceClient;
/**
* 创建订单 - 分布式事务
*/
@GlobalTransactional(timeoutMills = 300000, name = "create-order-tx")
public OrderDTO createOrder(CreateOrderRequest request) {
log.info("开始创建订单分布式事务");
// 1. 创建本地订单
OrderDTO order = createLocalOrder(request);
// 2. 扣减库存 - 远程调用
deductInventory(order);
// 3. 扣减余额 - 远程调用
deductBalance(order);
// 4. 更新订单状态
updateOrderStatus(order.getId(), OrderStatus.COMPLETED);
log.info("订单创建成功: {}", order.getId());
return order;
}
/**
* 创建本地订单
*/
private OrderDTO createLocalOrder(CreateOrderRequest request) {
Order order = new Order();
order.setUserId(request.getUserId());
order.setProductId(request.getProductId());
order.setQuantity(request.getQuantity());
order.setAmount(calculateAmount(request));
order.setStatus(OrderStatus.CREATED);
order.setCreateTime(new Date());
orderMapper.insert(order);
log.info("创建本地订单: {}", order.getId());
return convertToDTO(order);
}
/**
* 扣减库存
*/
private void deductInventory(OrderDTO order) {
try {
DeductInventoryRequest deductRequest = new DeductInventoryRequest();
deductRequest.setProductId(order.getProductId());
deductRequest.setQuantity(order.getQuantity());
ResponseEntity<Void> response = inventoryServiceClient.deductInventory(deductRequest);
if (!response.getStatusCode().is2xxSuccessful()) {
throw new RuntimeException("扣减库存失败");
}
log.info("扣减库存成功");
} catch (Exception e) {
log.error("扣减库存异常", e);
throw new RuntimeException("扣减库存服务调用失败", e);
}
}
/**
* 扣减余额
*/
private void deductBalance(OrderDTO order) {
try {
DeductBalanceRequest deductRequest = new DeductBalanceRequest();
deductRequest.setUserId(order.getUserId());
deductRequest.setAmount(order.getAmount());
ResponseEntity<Void> response = accountServiceClient.deductBalance(deductRequest);
if (!response.getStatusCode().is2xxSuccessful()) {
throw new RuntimeException("扣减余额失败");
}
log.info("扣减余额成功");
} catch (Exception e) {
log.error("扣减余额异常", e);
throw new RuntimeException("扣减余额服务调用失败", e);
}
}
/**
* 更新订单状态
*/
private void updateOrderStatus(Long orderId, OrderStatus status) {
Order order = new Order();
order.setId(orderId);
order.setStatus(status);
order.setUpdateTime(new Date());
orderMapper.updateById(order);
log.info("更新订单状态: {} -> {}", orderId, status);
}
private BigDecimal calculateAmount(CreateOrderRequest request) {
// 模拟金额计算
return BigDecimal.valueOf(request.getQuantity() * 100);
}
private OrderDTO convertToDTO(Order order) {
OrderDTO dto = new OrderDTO();
dto.setId(order.getId());
dto.setUserId(order.getUserId());
dto.setProductId(order.getProductId());
dto.setQuantity(order.getQuantity());
dto.setAmount(order.getAmount());
dto.setStatus(order.getStatus());
return dto;
}
}
// 库存服务 - 分布式事务参与者
@Service
@Slf4j
public class InventoryService {
@Autowired
private InventoryMapper inventoryMapper;
/**
* 扣减库存 - 分布式事务参与者
*/
@Transactional
@GlobalTransactional
public void deductInventory(DeductInventoryRequest request) {
log.info("开始扣减库存: {}", request);
// 查询当前库存
Inventory inventory = inventoryMapper.selectByProductId(request.getProductId());
if (inventory == null) {
throw new RuntimeException("产品不存在: " + request.getProductId());
}
// 检查库存是否充足
if (inventory.getStock() < request.getQuantity()) {
throw new RuntimeException("库存不足");
}
// 扣减库存
int updated = inventoryMapper.deductStock(
request.getProductId(),
request.getQuantity()
);
if (updated == 0) {
throw new RuntimeException("扣减库存失败");
}
log.info("扣减库存成功: 产品={}, 数量={}",
request.getProductId(), request.getQuantity());
}
/**
* 补偿操作 - 增加库存
*/
@Compensable
public void compensateDeductInventory(DeductInventoryRequest request) {
log.info("执行库存补偿操作: {}", request);
try {
int updated = inventoryMapper.addStock(
request.getProductId(),
request.getQuantity()
);
if (updated > 0) {
log.info("库存补偿成功");
} else {
log.error("库存补偿失败");
}
} catch (Exception e) {
log.error("库存补偿异常", e);
// 补偿操作失败需要人工干预
alertCompensationFailure("inventory", request, e);
}
}
private void alertCompensationFailure(String service, Object request, Exception e) {
// 发送告警通知
log.error("补偿操作失败告警: service={}, request={}", service, request);
// 可以集成到监控系统或发送邮件/短信
}
}
// Seata配置类
@Configuration
@Slf4j
public class SeataConfiguration {
@Bean
public GlobalTransactionScanner globalTransactionScanner() {
return new GlobalTransactionScanner(
"order-service",
"my_test_tx_group"
);
}
@Bean
@ConfigurationProperties(prefix = "spring.datasource")
public DruidDataSource druidDataSource() {
return new DruidDataSource();
}
@Bean
public DataSourceProxy dataSourceProxy(DataSource dataSource) {
return new DataSourceProxy(dataSource);
}
@Bean
public SqlSessionFactory sqlSessionFactoryBean(DataSourceProxy dataSourceProxy) throws Exception {
SqlSessionFactoryBean sqlSessionFactoryBean = new SqlSessionFactoryBean();
sqlSessionFactoryBean.setDataSource(dataSourceProxy);
return sqlSessionFactoryBean.getObject();
}
}
// 分布式事务监控
@Component
@Slf4j
public class DistributedTransactionMonitor {
@EventListener
public void onGlobalTransactionEvent(GlobalTransactionEvent event) {
log.info("分布式事务事件: {}", event);
switch (event.getStatus()) {
case Begin:
log.info("分布式事务开始: xid={}", event.getXid());
break;
case Commit:
log.info("分布式事务提交: xid={}", event.getXid());
break;
case Rollback:
log.info("分布式事务回滚: xid={}", event.getXid());
break;
case Timeout:
log.warn("分布式事务超时: xid={}", event.getXid());
break;
}
// 记录事务指标
recordTransactionMetrics(event);
}
private void recordTransactionMetrics(GlobalTransactionEvent event) {
// 记录事务相关的监控指标
Counter.builder("distributed.transaction.total")
.tag("status", event.getStatus().name())
.register(Metrics.globalRegistry)
.increment();
}
// 定期检查悬挂事务
@Scheduled(fixedRate = 60000) // 每分钟检查一次
public void checkHangingTransactions() {
log.info("检查悬挂的分布式事务...");
// 实现悬挂事务检查逻辑
}
}
微服务监控与治理
6.1 全链路监控
// 链路追踪配置
@Configuration
@Slf4j
public class TracingConfiguration {
@Bean
public Tracer tracer() {
return new OpenTracingTracer();
}
@Bean
public Brave brave() {
return newBrave("user-service");
}
private Brave newBrave(String serviceName) {
return new Brave.Builder(serviceName)
.traceSampler(Sampler.ALWAYS_SAMPLE)
.build();
}
@Bean
public BraveServletFilter braveServletFilter(Brave brave) {
return new BraveServletFilter(
brave.serverRequestInterceptor(),
brave.serverResponseInterceptor(),
new DefaultSpanNameProvider()
);
}
}
// 自定义链路追踪拦截器
@Component
@Slf4j
public class CustomTracingInterceptor implements HandlerInterceptor {
@Autowired
private Tracer tracer;
private static final String SPAN_NAME_PREFIX = "HTTP_";
@Override
public boolean preHandle(HttpServletRequest request, HttpServletResponse response, Object handler) {
String spanName = SPAN_NAME_PREFIX + request.getMethod() + ":" + request.getRequestURI();
Span span = tracer.buildSpan(spanName)
.withTag("http.method", request.getMethod())
.withTag("http.url", request.getRequestURL().toString())
.withTag("http.headers", getHeadersAsString(request))
.start();
tracer.activateSpan(span);
request.setAttribute("span", span);
log.info("开始处理请求: {} {}", request.getMethod(), request.getRequestURI());
return true;
}
@Override
public void afterCompletion(HttpServletRequest request, HttpServletResponse response, Object handler, Exception ex) {
Span span = (Span) request.getAttribute("span");
if (span != null) {
span.setTag("http.status_code", response.getStatus());
if (ex != null) {
span.setTag("error", true);
span.log(Map.of("error.message", ex.getMessage()));
}
span.finish();
}
log.info("请求处理完成: {} {} - {}",
request.getMethod(), request.getRequestURI(), response.getStatus());
}
private String getHeadersAsString(HttpServletRequest request) {
Enumeration<String> headerNames = request.getHeaderNames();
StringBuilder headers = new StringBuilder();
while (headerNames.hasMoreElements()) {
String headerName = headerNames.nextElement();
String headerValue = request.getHeader(headerName);
headers.append(headerName).append(": ").append(headerValue).append("; ");
}
return headers.toString();
}
}
// 微服务健康检查
@Component
@Slf4j
public class MicroserviceHealthIndicator implements HealthIndicator {
@Autowired
private DataSource dataSource;
@Autowired
private DiscoveryClient discoveryClient;
@Value("${spring.application.name}")
private String applicationName;
@Override
public Health health() {
Map<String, Object> details = new HashMap<>();
// 检查数据库连接
boolean dbHealthy = checkDatabaseHealth();
details.put("database", dbHealthy ? "UP" : "DOWN");
// 检查服务注册状态
boolean registrationHealthy = checkServiceRegistration();
details.put("serviceRegistration", registrationHealthy ? "UP" : "DOWN");
// 检查依赖服务状态
Map<String, String> dependencies = checkDependencies();
details.put("dependencies", dependencies);
// 检查系统资源
details.putAll(checkSystemResources());
if (dbHealthy && registrationHealthy) {
return Health.up().withDetails(details).build();
} else {
return Health.down().withDetails(details).build();
}
}
private boolean checkDatabaseHealth() {
try (Connection conn = dataSource.getConnection()) {
return conn.isValid(5);
} catch (SQLException e) {
log.error("数据库健康检查失败", e);
return false;
}
}
private boolean checkServiceRegistration() {
try {
List<ServiceInstance> instances = discoveryClient.getInstances(applicationName);
return !instances.isEmpty();
} catch (Exception e) {
log.error("服务注册状态检查失败", e);
return false;
}
}
private Map<String, String> checkDependencies() {
Map<String, String> dependencies = new HashMap<>();
// 检查关键依赖服务
String[] criticalServices = {"nacos-server", "config-service", "gateway-service"};
for (String service : criticalServices) {
try {
List<ServiceInstance> instances = discoveryClient.getInstances(service);
dependencies.put(service, instances.isEmpty() ? "DOWN" : "UP");
} catch (Exception e) {
dependencies.put(service, "UNKNOWN");
}
}
return dependencies;
}
private Map<String, Object> checkSystemResources() {
Map<String, Object> resources = new HashMap<>();
Runtime runtime = Runtime.getRuntime();
resources.put("memory.used", runtime.totalMemory() - runtime.freeMemory());
resources.put("memory.max", runtime.maxMemory());
resources.put("processors", runtime.availableProcessors());
try {
OperatingSystemMXBean osBean = ManagementFactory.getOperatingSystemMXBean();
resources.put("system.load", osBean.getSystemLoadAverage());
} catch (Exception e) {
log.warn("获取系统负载失败", e);
}
return resources;
}
}
总结
🎯 Spring Cloud Alibaba核心组件
-
Nacos - 服务注册发现与配置中心
- 动态服务发现与健康检查
- 配置集中管理与实时刷新
- 命名空间与分组管理
-
Sentinel - 流量控制与系统保护
- 细粒度流量控制
- 熔断降级与系统自适应保护
- 实时监控与规则管理
-
Seata - 分布式事务解决方案
- AT模式无侵入分布式事务
- TCC模式高性能事务
- Saga模式长事务支持
-
OpenFeign - 声明式HTTP客户端
- 服务间调用简化
- 负载均衡与熔断降级
- 可扩展的配置体系
🚀 微服务最佳实践
-
服务设计原则
- 单一职责与界限上下文
- 接口版本管理与兼容性
- 优雅停机与健康检查
-
分布式系统考量
- 分布式事务与数据一致性
- 服务容错与降级策略
- 链路追踪与问题定位
-
性能与可观测性
- 监控指标收集与告警
- 日志聚合与分析
- 性能调优与容量规划
📋 生产环境检查表
- 服务注册发现配置正确
- 配置中心数据备份
- 流量控制规则验证
- 分布式事务测试
- 监控告警配置完善
- 安全认证机制到位
- 日志链路追踪完整
- 容灾备份方案准备
Spring Cloud Alibaba为微服务架构提供了完整的解决方案,但在生产环境中需要根据具体业务场景进行适当的调整和优化。
下一篇预告:《Docker & Kubernetes Java应用部署实战》
如果觉得本文对你有帮助,请点赞、收藏、关注!欢迎在评论区分享你的微服务实践经验和问题。
更多推荐


所有评论(0)