引言

微服务架构已成为现代分布式系统的主流选择,而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核心组件

  1. Nacos - 服务注册发现与配置中心

    • 动态服务发现与健康检查
    • 配置集中管理与实时刷新
    • 命名空间与分组管理
  2. Sentinel - 流量控制与系统保护

    • 细粒度流量控制
    • 熔断降级与系统自适应保护
    • 实时监控与规则管理
  3. Seata - 分布式事务解决方案

    • AT模式无侵入分布式事务
    • TCC模式高性能事务
    • Saga模式长事务支持
  4. OpenFeign - 声明式HTTP客户端

    • 服务间调用简化
    • 负载均衡与熔断降级
    • 可扩展的配置体系

🚀 微服务最佳实践

  1. 服务设计原则

    • 单一职责与界限上下文
    • 接口版本管理与兼容性
    • 优雅停机与健康检查
  2. 分布式系统考量

    • 分布式事务与数据一致性
    • 服务容错与降级策略
    • 链路追踪与问题定位
  3. 性能与可观测性

    • 监控指标收集与告警
    • 日志聚合与分析
    • 性能调优与容量规划

📋 生产环境检查表

  • 服务注册发现配置正确
  • 配置中心数据备份
  • 流量控制规则验证
  • 分布式事务测试
  • 监控告警配置完善
  • 安全认证机制到位
  • 日志链路追踪完整
  • 容灾备份方案准备

Spring Cloud Alibaba为微服务架构提供了完整的解决方案,但在生产环境中需要根据具体业务场景进行适当的调整和优化。


下一篇预告:《Docker & Kubernetes Java应用部署实战》

如果觉得本文对你有帮助,请点赞、收藏、关注!欢迎在评论区分享你的微服务实践经验和问题。

更多推荐