10-微服务架构设计与实践
·
微服务架构设计与实践
引言微服务架构作为现代软件架构的重要模式,通过将大型应用拆分为多个小型、独立的服务来提高系统的可维护性、可扩展性和技术多样性。本文将深入探讨微服务架构的设计原则、实现方式和最佳实践。## 微服务架构概述### 什么是微服务微服务是一种架构风格,它将应用程序构建为一组小型服务,每个服务:- 运行在自己的进程中- 通过轻量级机制(通常是HTTP API)进行通信- 围绕业务能力构建- 可以独立部署- 由小型团队维护### 微服务 vs 单体架构mermaidgraph TB subgraph "单体架构" A[用户界面] --> B[业务逻辑层] B --> C[数据访问层] C --> D[数据库] end subgraph "微服务架构" E[API网关] --> F[用户服务] E --> G[订单服务] E --> H[支付服务] E --> I[库存服务] F --> J[用户数据库] G --> K[订单数据库] H --> L[支付数据库] I --> M[库存数据库] end## 微服务设计原则### 1. 单一职责原则每个微服务应该专注于一个业务领域:java// 用户服务 - 只处理用户相关业务@RestController@RequestMapping("/api/users")public class UserController { @Autowired private UserService userService; @PostMapping public ResponseEntity<User> createUser(@RequestBody CreateUserRequest request) { User user = userService.createUser(request); return ResponseEntity.ok(user); } @GetMapping("/{userId}") public ResponseEntity<User> getUser(@PathVariable Long userId) { User user = userService.findById(userId); return ResponseEntity.ok(user); } @PutMapping("/{userId}") public ResponseEntity<User> updateUser( @PathVariable Long userId, @RequestBody UpdateUserRequest request) { User user = userService.updateUser(userId, request); return ResponseEntity.ok(user); }}// 订单服务 - 只处理订单相关业务@RestController@RequestMapping("/api/orders")public class OrderController { @Autowired private OrderService orderService; @PostMapping public ResponseEntity<Order> createOrder(@RequestBody CreateOrderRequest request) { Order order = orderService.createOrder(request); return ResponseEntity.ok(order); } @GetMapping("/{orderId}") public ResponseEntity<Order> getOrder(@PathVariable Long orderId) { Order order = orderService.findById(orderId); return ResponseEntity.ok(order); }}### 2. 数据库独立性每个微服务拥有自己的数据库:yaml# docker-compose.ymlversion: '3.8'services: user-service: build: ./user-service environment: - DB_HOST=user-db - DB_NAME=userdb depends_on: - user-db user-db: image: postgres:13 environment: - POSTGRES_DB=userdb - POSTGRES_USER=userservice - POSTGRES_PASSWORD=password volumes: - user_data:/var/lib/postgresql/data order-service: build: ./order-service environment: - DB_HOST=order-db - DB_NAME=orderdb depends_on: - order-db order-db: image: postgres:13 environment: - POSTGRES_DB=orderdb - POSTGRES_USER=orderservice - POSTGRES_PASSWORD=password volumes: - order_data:/var/lib/postgresql/datavolumes: user_data: order_data:### 3. API优先设计定义清晰的API契约:yaml# user-service-api.yamlopenapi: 3.0.0info: title: User Service API version: 1.0.0 description: 用户服务APIpaths: /api/users: post: summary: 创建用户 requestBody: required: true content: application/json: schema: $ref: '#/components/schemas/CreateUserRequest' responses: '201': description: 用户创建成功 content: application/json: schema: $ref: '#/components/schemas/User' '400': description: 请求参数错误 '409': description: 用户已存在 /api/users/{userId}: get: summary: 获取用户信息 parameters: - name: userId in: path required: true schema: type: integer format: int64 responses: '200': description: 成功获取用户信息 content: application/json: schema: $ref: '#/components/schemas/User' '404': description: 用户不存在components: schemas: User: type: object properties: id: type: integer format: int64 username: type: string email: type: string createdAt: type: string format: date-time CreateUserRequest: type: object required: - username - email - password properties: username: type: string minLength: 3 maxLength: 50 email: type: string format: email password: type: string minLength: 8## 服务间通信### 1. 同步通信 - REST APIjava// 订单服务调用用户服务@Servicepublic class OrderService { @Autowired private UserServiceClient userServiceClient; @Autowired private OrderRepository orderRepository; public Order createOrder(CreateOrderRequest request) { // 验证用户是否存在 User user = userServiceClient.getUser(request.getUserId()); if (user == null) { throw new UserNotFoundException("用户不存在: " + request.getUserId()); } // 创建订单 Order order = new Order(); order.setUserId(request.getUserId()); order.setProductId(request.getProductId()); order.setQuantity(request.getQuantity()); order.setStatus(OrderStatus.PENDING); order.setCreatedAt(LocalDateTime.now()); return orderRepository.save(order); }}// Feign客户端@FeignClient(name = "user-service", url = "${services.user-service.url}")public interface UserServiceClient { @GetMapping("/api/users/{userId}") User getUser(@PathVariable("userId") Long userId); @PostMapping("/api/users") User createUser(@RequestBody CreateUserRequest request);}// 配置@Configuration@EnableFeignClientspublic class FeignConfig { @Bean public RequestInterceptor requestInterceptor() { return requestTemplate -> { // 添加认证头 String token = SecurityContextHolder.getContext() .getAuthentication().getCredentials().toString(); requestTemplate.header("Authorization", "Bearer " + token); }; } @Bean public ErrorDecoder errorDecoder() { return new CustomErrorDecoder(); }}### 2. 异步通信 - 消息队列java// 事件发布@Servicepublic class OrderService { @Autowired private RabbitTemplate rabbitTemplate; public Order createOrder(CreateOrderRequest request) { Order order = // ... 创建订单逻辑 // 发布订单创建事件 OrderCreatedEvent event = new OrderCreatedEvent( order.getId(), order.getUserId(), order.getProductId(), order.getQuantity(), order.getTotalAmount() ); rabbitTemplate.convertAndSend( "order.exchange", "order.created", event ); return order; }}// 事件监听@Componentpublic class InventoryEventListener { @Autowired private InventoryService inventoryService; @RabbitListener(queues = "inventory.order.created") public void handleOrderCreated(OrderCreatedEvent event) { try { // 减少库存 inventoryService.decreaseStock( event.getProductId(), event.getQuantity() ); log.info("库存已更新,订单ID: {}", event.getOrderId()); } catch (InsufficientStockException e) { // 发布库存不足事件 publishInsufficientStockEvent(event); } } private void publishInsufficientStockEvent(OrderCreatedEvent event) { InsufficientStockEvent stockEvent = new InsufficientStockEvent( event.getOrderId(), event.getProductId(), event.getQuantity() ); rabbitTemplate.convertAndSend( "inventory.exchange", "stock.insufficient", stockEvent ); }}// 消息队列配置@Configuration@EnableRabbitpublic class RabbitConfig { // 订单交换机 @Bean public TopicExchange orderExchange() { return new TopicExchange("order.exchange"); } // 库存队列 @Bean public Queue inventoryOrderQueue() { return QueueBuilder.durable("inventory.order.created").build(); } // 绑定 @Bean public Binding inventoryOrderBinding() { return BindingBuilder .bind(inventoryOrderQueue()) .to(orderExchange()) .with("order.created"); } // 库存交换机 @Bean public TopicExchange inventoryExchange() { return new TopicExchange("inventory.exchange"); } // 库存不足队列 @Bean public Queue stockInsufficientQueue() { return QueueBuilder.durable("order.stock.insufficient").build(); } @Bean public Binding stockInsufficientBinding() { return BindingBuilder .bind(stockInsufficientQueue()) .to(inventoryExchange()) .with("stock.insufficient"); }}## 服务发现与注册### 1. Eureka服务注册中心java// Eureka Server@SpringBootApplication@EnableEurekaServerpublic class EurekaServerApplication { public static void main(String[] args) { SpringApplication.run(EurekaServerApplication.class, args); }}// application.ymlserver: port: 8761eureka: instance: hostname: localhost client: register-with-eureka: false fetch-registry: false service-url: defaultZone: http://${eureka.instance.hostname}:${server.port}/eureka/ server: enable-self-preservation: false``````java// 微服务客户端@SpringBootApplication@EnableEurekaClient@EnableFeignClientspublic class UserServiceApplication { public static void main(String[] args) { SpringApplication.run(UserServiceApplication.class, args); }}// application.ymlspring: application: name: user-serviceserver: port: 8081eureka: client: service-url: defaultZone: http://localhost:8761/eureka/ instance: prefer-ip-address: true lease-renewal-interval-in-seconds: 10 lease-expiration-duration-in-seconds: 30### 2. Consul服务发现java// Consul配置@Configurationpublic class ConsulConfig { @Bean public ConsulClient consulClient() { return new ConsulClient("localhost", 8500); } @Bean public ServiceRegistry serviceRegistry() { return new ConsulServiceRegistry(consulClient()); }}// 服务注册@Componentpublic class ServiceRegistration { @Autowired private ConsulClient consulClient; @Value("${spring.application.name}") private String serviceName; @Value("${server.port}") private int serverPort; @PostConstruct public void registerService() { NewService service = new NewService(); service.setId(serviceName + "-" + UUID.randomUUID().toString()); service.setName(serviceName); service.setAddress("localhost"); service.setPort(serverPort); // 健康检查 NewService.Check check = new NewService.Check(); check.setHttp("http://localhost:" + serverPort + "/actuator/health"); check.setInterval("10s"); service.setCheck(check); consulClient.agentServiceRegister(service); }}## API网关### 1. Spring Cloud Gatewayjava@SpringBootApplicationpublic class GatewayApplication { public static void main(String[] args) { SpringApplication.run(GatewayApplication.class, args); }}// 路由配置@Configurationpublic class GatewayConfig { @Bean public RouteLocator customRouteLocator(RouteLocatorBuilder builder) { return builder.routes() // 用户服务路由 .route("user-service", r -> r .path("/api/users/**") .filters(f -> f .addRequestHeader("X-Gateway", "Spring-Cloud-Gateway") .circuitBreaker(config -> config .setName("user-service-cb") .setFallbackUri("forward:/fallback/users")) ) .uri("lb://user-service")) // 订单服务路由 .route("order-service", r -> r .path("/api/orders/**") .filters(f -> f .addRequestHeader("X-Gateway", "Spring-Cloud-Gateway") .retry(config -> config .setRetries(3) .setMethods(HttpMethod.GET) .setBackoff(Duration.ofMillis(100), Duration.ofMillis(1000), 2, true)) ) .uri("lb://order-service")) // 支付服务路由 .route("payment-service", r -> r .path("/api/payments/**") .filters(f -> f .requestRateLimiter(config -> config .setRateLimiter(redisRateLimiter()) .setKeyResolver(userKeyResolver())) ) .uri("lb://payment-service")) .build(); } @Bean public RedisRateLimiter redisRateLimiter() { return new RedisRateLimiter(10, 20, 1); } @Bean public KeyResolver userKeyResolver() { return exchange -> exchange.getRequest().getHeaders() .getFirst("X-User-Id"); }}// 全局过滤器@Componentpublic class AuthenticationFilter implements GlobalFilter, Ordered { @Autowired private JwtTokenUtil jwtTokenUtil; @Override public Mono<Void> filter(ServerWebExchange exchange, GatewayFilterChain chain) { ServerHttpRequest request = exchange.getRequest(); // 跳过认证的路径 if (isExcludedPath(request.getPath().value())) { return chain.filter(exchange); } String token = extractToken(request); if (token == null || !jwtTokenUtil.validateToken(token)) { return unauthorized(exchange); } // 添加用户信息到请求头 String userId = jwtTokenUtil.getUserIdFromToken(token); ServerHttpRequest modifiedRequest = request.mutate() .header("X-User-Id", userId) .build(); return chain.filter(exchange.mutate().request(modifiedRequest).build()); } private boolean isExcludedPath(String path) { return path.startsWith("/api/auth/") || path.startsWith("/actuator/") || path.equals("/api/health"); } private String extractToken(ServerHttpRequest request) { String authHeader = request.getHeaders().getFirst("Authorization"); if (authHeader != null && authHeader.startsWith("Bearer ")) { return authHeader.substring(7); } return null; } private Mono<Void> unauthorized(ServerWebExchange exchange) { ServerHttpResponse response = exchange.getResponse(); response.setStatusCode(HttpStatus.UNAUTHORIZED); return response.setComplete(); } @Override public int getOrder() { return -100; }}### 2. 熔断器模式java// 熔断器配置@Configurationpublic class CircuitBreakerConfig { @Bean public Customizer<ReactiveResilience4JCircuitBreakerFactory> defaultCustomizer() { return factory -> factory.configureDefault(id -> new Resilience4JConfigBuilder(id) .circuitBreakerConfig(CircuitBreakerConfig.custom() .slidingWindowSize(10) .minimumNumberOfCalls(5) .failureRateThreshold(50.0f) .waitDurationInOpenState(Duration.ofSeconds(30)) .slowCallRateThreshold(50.0f) .slowCallDurationThreshold(Duration.ofSeconds(2)) .build()) .timeLimiterConfig(TimeLimiterConfig.custom() .timeoutDuration(Duration.ofSeconds(3)) .build()) .build()); }}// 降级处理@RestControllerpublic class FallbackController { @RequestMapping("/fallback/users") public ResponseEntity<Map<String, Object>> userFallback() { Map<String, Object> response = new HashMap<>(); response.put("message", "用户服务暂时不可用,请稍后重试"); response.put("status", "SERVICE_UNAVAILABLE"); return ResponseEntity.status(HttpStatus.SERVICE_UNAVAILABLE).body(response); } @RequestMapping("/fallback/orders") public ResponseEntity<Map<String, Object>> orderFallback() { Map<String, Object> response = new HashMap<>(); response.put("message", "订单服务暂时不可用,请稍后重试"); response.put("status", "SERVICE_UNAVAILABLE"); return ResponseEntity.status(HttpStatus.SERVICE_UNAVAILABLE).body(response); }}## 分布式事务### 1. Saga模式java// 订单Saga编排器@Componentpublic class OrderSagaOrchestrator { @Autowired private PaymentService paymentService; @Autowired private InventoryService inventoryService; @Autowired private OrderService orderService; @Autowired private NotificationService notificationService; public void processOrder(CreateOrderRequest request) { SagaTransaction saga = SagaTransaction.builder() .step("createOrder", () -> orderService.createOrder(request)) .compensate("cancelOrder", orderId -> orderService.cancelOrder(orderId)) .step("reserveInventory", orderId -> inventoryService.reserveStock( request.getProductId(), request.getQuantity())) .compensate("releaseInventory", orderId -> inventoryService.releaseStock( request.getProductId(), request.getQuantity())) .step("processPayment", orderId -> paymentService.processPayment( request.getUserId(), request.getTotalAmount())) .compensate("refundPayment", paymentId -> paymentService.refund(paymentId)) .step("confirmOrder", orderId -> orderService.confirmOrder(orderId)) .step("sendNotification", orderId -> notificationService.sendOrderConfirmation(orderId)) .build(); try { saga.execute(); } catch (SagaExecutionException e) { log.error("订单处理失败,执行补偿操作", e); saga.compensate(); } }}// Saga事务框架public class SagaTransaction { private List<SagaStep> steps = new ArrayList<>(); private List<Object> results = new ArrayList<>(); public static SagaBuilder builder() { return new SagaBuilder(); } public void execute() throws SagaExecutionException { for (int i = 0; i < steps.size(); i++) { try { SagaStep step = steps.get(i); Object result = step.execute(getLastResult()); results.add(result); log.info("步骤 {} 执行成功", step.getName()); } catch (Exception e) { log.error("步骤 {} 执行失败", steps.get(i).getName(), e); compensate(i - 1); throw new SagaExecutionException("Saga执行失败", e); } } } public void compensate() { compensate(results.size() - 1); } private void compensate(int fromIndex) { for (int i = fromIndex; i >= 0; i--) { try { SagaStep step = steps.get(i); if (step.getCompensateAction() != null) { step.compensate(results.get(i)); log.info("步骤 {} 补偿成功", step.getName()); } } catch (Exception e) { log.error("步骤 {} 补偿失败", steps.get(i).getName(), e); } } } private Object getLastResult() { return results.isEmpty() ? null : results.get(results.size() - 1); }}### 2. 分布式事务消息java// 事务消息发送@Servicepublic class OrderTransactionService { @Autowired private OrderRepository orderRepository; @Autowired private TransactionMessageProducer messageProducer; @Transactional public Order createOrderWithMessage(CreateOrderRequest request) { // 1. 创建订单 Order order = new Order(); order.setUserId(request.getUserId()); order.setProductId(request.getProductId()); order.setQuantity(request.getQuantity()); order.setStatus(OrderStatus.PENDING); order = orderRepository.save(order); // 2. 发送事务消息 OrderCreatedEvent event = new OrderCreatedEvent( order.getId(), order.getUserId(), order.getProductId(), order.getQuantity() ); messageProducer.sendTransactionMessage( "order.created", event, order.getId().toString() ); return order; }}// 事务消息生产者@Componentpublic class TransactionMessageProducer { @Autowired private RabbitTemplate rabbitTemplate; @Autowired private TransactionMessageRepository messageRepository; public void sendTransactionMessage(String routingKey, Object message, String businessKey) { // 1. 保存消息到本地事务表 TransactionMessage txMessage = new TransactionMessage(); txMessage.setId(UUID.randomUUID().toString()); txMessage.setRoutingKey(routingKey); txMessage.setPayload(JsonUtils.toJson(message)); txMessage.setBusinessKey(businessKey); txMessage.setStatus(MessageStatus.PENDING); txMessage.setCreatedAt(LocalDateTime.now()); messageRepository.save(txMessage); // 2. 发送消息到MQ try { rabbitTemplate.convertAndSend("order.exchange", routingKey, message); // 3. 更新消息状态为已发送 txMessage.setStatus(MessageStatus.SENT); txMessage.setSentAt(LocalDateTime.now()); messageRepository.save(txMessage); } catch (Exception e) { log.error("发送消息失败", e); txMessage.setStatus(MessageStatus.FAILED); txMessage.setErrorMessage(e.getMessage()); messageRepository.save(txMessage); throw e; } }}// 消息补偿任务@Componentpublic class MessageCompensationTask { @Autowired private TransactionMessageRepository messageRepository; @Autowired private RabbitTemplate rabbitTemplate; @Scheduled(fixedDelay = 60000) // 每分钟执行一次 public void compensateFailedMessages() { LocalDateTime cutoff = LocalDateTime.now().minusMinutes(5); List<TransactionMessage> failedMessages = messageRepository .findByStatusAndCreatedAtBefore(MessageStatus.PENDING, cutoff); for (TransactionMessage message : failedMessages) { try { // 重新发送消息 Object payload = JsonUtils.fromJson(message.getPayload(), Object.class); rabbitTemplate.convertAndSend( "order.exchange", message.getRoutingKey(), payload ); message.setStatus(MessageStatus.SENT); message.setSentAt(LocalDateTime.now()); messageRepository.save(message); log.info("补偿发送消息成功: {}", message.getId()); } catch (Exception e) { log.error("补偿发送消息失败: {}", message.getId(), e); message.setRetryCount(message.getRetryCount() + 1); if (message.getRetryCount() >= 3) { message.setStatus(MessageStatus.FAILED); } messageRepository.save(message); } } }}## 配置管理### 1. Spring Cloud Configjava// Config Server@SpringBootApplication@EnableConfigServerpublic class ConfigServerApplication { public static void main(String[] args) { SpringApplication.run(ConfigServerApplication.class, args); }}// application.ymlserver: port: 8888spring: cloud: config: server: git: uri: https://github.com/your-org/config-repo search-paths: '{application}' default-label: main encrypt: enabled: true encrypt: key: mySecretKey``````yaml# 配置文件结构config-repo/├── user-service/│ ├── user-service.yml│ ├── user-service-dev.yml│ ├── user-service-prod.yml└── order-service/ ├── order-service.yml ├── order-service-dev.yml └── order-service-prod.yml``````yaml# user-service-prod.ymlspring: datasource: url: jdbc:postgresql://prod-db:5432/userdb username: '{cipher}AQA...' # 加密的用户名 password: '{cipher}AQB...' # 加密的密码 redis: host: prod-redis port: 6379 password: '{cipher}AQC...'logging: level: com.example.userservice: INFO pattern: console: '%d{yyyy-MM-dd HH:mm:ss} - %msg%n'management: endpoints: web: exposure: include: health,info,metrics### 2. 动态配置刷新java// 配置类@Component@RefreshScope@ConfigurationProperties(prefix = "app")public class AppConfig { private String name; private int maxConnections; private Duration timeout; private List<String> allowedOrigins; // getters and setters}// 使用配置@RestController@RefreshScopepublic class ConfigController { @Autowired private AppConfig appConfig; @Value("${app.feature.enabled:false}") private boolean featureEnabled; @GetMapping("/config") public Map<String, Object> getConfig() { Map<String, Object> config = new HashMap<>(); config.put("appName", appConfig.getName()); config.put("maxConnections", appConfig.getMaxConnections()); config.put("featureEnabled", featureEnabled); return config; }}// 配置变更监听@Componentpublic class ConfigChangeListener { @EventListener public void handleRefreshEvent(RefreshRemoteApplicationEvent event) { log.info("配置刷新事件: {}", event.getDestinationService()); // 执行配置变更后的逻辑 }}## 监控与日志### 1. 分布式链路追踪java// Sleuth配置@Configurationpublic class TracingConfig { @Bean public Sender sender() { return OkHttpSender.create("http://zipkin:9411/api/v2/spans"); } @Bean public AsyncReporter<Span> spanReporter() { return AsyncReporter.create(sender()); } @Bean public Tracing tracing() { return Tracing.newBuilder() .localServiceName("user-service") .spanReporter(spanReporter()) .sampler(Sampler.create(1.0f)) // 100%采样 .build(); }}// 自定义Span@Servicepublic class UserService { private final Tracer tracer; public UserService(Tracing tracing) { this.tracer = tracing.tracer(); } public User findById(Long userId) { Span span = tracer.nextSpan() .name("user-service.find-by-id") .tag("user.id", userId.toString()) .start(); try (Tracer.SpanInScope ws = tracer.withSpanInScope(span)) { // 业务逻辑 User user = userRepository.findById(userId); if (user != null) { span.tag("user.found", "true"); } else { span.tag("user.found", "false"); } return user; } catch (Exception e) { span.tag("error", e.getMessage()); throw e; } finally { span.end(); } }}### 2. 集中化日志yaml# logback-spring.xml<configuration> <springProfile name="!local"> <appender name="STDOUT" class="ch.qos.logback.core.ConsoleAppender"> <encoder class="net.logstash.logback.encoder.LoggingEventCompositeJsonEncoder"> <providers> <timestamp/> <logLevel/> <loggerName/> <mdc/> <pattern> <pattern> { "traceId": "%X{traceId:-}", "spanId": "%X{spanId:-}", "service": "user-service", "message": "%message" } </pattern> </pattern> </providers> </encoder> </appender> </springProfile> <springProfile name="local"> <appender name="STDOUT" class="ch.qos.logback.core.ConsoleAppender"> <encoder> <pattern>%d{HH:mm:ss.SSS} [%thread] %-5level [%X{traceId:-},%X{spanId:-}] %logger{36} - %msg%n</pattern> </encoder> </appender> </springProfile> <root level="INFO"> <appender-ref ref="STDOUT"/> </root></configuration>### 3. 应用监控java// 自定义指标@Componentpublic class UserServiceMetrics { private final Counter userCreatedCounter; private final Timer userQueryTimer; private final Gauge activeUsersGauge; public UserServiceMetrics(MeterRegistry meterRegistry) { this.userCreatedCounter = Counter.builder("user.created") .description("用户创建数量") .register(meterRegistry); this.userQueryTimer = Timer.builder("user.query.duration") .description("用户查询耗时") .register(meterRegistry); this.activeUsersGauge = Gauge.builder("user.active") .description("活跃用户数") .register(meterRegistry, this, UserServiceMetrics::getActiveUserCount); } public void incrementUserCreated() { userCreatedCounter.increment(); } public Timer.Sample startUserQueryTimer() { return Timer.start(); } public void recordUserQueryTime(Timer.Sample sample) { sample.stop(userQueryTimer); } private double getActiveUserCount() { // 实现获取活跃用户数的逻辑 return 0.0; }}// 健康检查@Componentpublic class UserServiceHealthIndicator implements HealthIndicator { @Autowired private UserRepository userRepository; @Autowired private RedisTemplate<String, String> redisTemplate; @Override public Health health() { Health.Builder builder = Health.up(); try { // 检查数据库连接 long userCount = userRepository.count(); builder.withDetail("database", "UP") .withDetail("userCount", userCount); } catch (Exception e) { builder.down() .withDetail("database", "DOWN") .withDetail("error", e.getMessage()); } try { // 检查Redis连接 redisTemplate.opsForValue().get("health-check"); builder.withDetail("redis", "UP"); } catch (Exception e) { builder.withDetail("redis", "DOWN") .withDetail("error", e.getMessage()); } return builder.build(); }}## 安全性### 1. JWT认证java// JWT工具类@Componentpublic class JwtTokenUtil { private String secret = "mySecretKey"; private int expiration = 86400; // 24小时 public String generateToken(UserDetails userDetails) { Map<String, Object> claims = new HashMap<>(); claims.put("sub", userDetails.getUsername()); claims.put("roles", userDetails.getAuthorities().stream() .map(GrantedAuthority::getAuthority) .collect(Collectors.toList())); return createToken(claims, userDetails.getUsername()); } private String createToken(Map<String, Object> claims, String subject) { return Jwts.builder() .setClaims(claims) .setSubject(subject) .setIssuedAt(new Date(System.currentTimeMillis())) .setExpiration(new Date(System.currentTimeMillis() + expiration * 1000)) .signWith(SignatureAlgorithm.HS512, secret) .compact(); } public Boolean validateToken(String token, UserDetails userDetails) { final String username = getUsernameFromToken(token); return (username.equals(userDetails.getUsername()) && !isTokenExpired(token)); } public String getUsernameFromToken(String token) { return getClaimFromToken(token, Claims::getSubject); } public Date getExpirationDateFromToken(String token) { return getClaimFromToken(token, Claims::getExpiration); } public <T> T getClaimFromToken(String token, Function<Claims, T> claimsResolver) { final Claims claims = getAllClaimsFromToken(token); return claimsResolver.apply(claims); } private Claims getAllClaimsFromToken(String token) { return Jwts.parser().setSigningKey(secret).parseClaimsJws(token).getBody(); } private Boolean isTokenExpired(String token) { final Date expiration = getExpirationDateFromToken(token); return expiration.before(new Date()); }}// JWT过滤器@Componentpublic class JwtAuthenticationFilter extends OncePerRequestFilter { @Autowired private UserDetailsService userDetailsService; @Autowired private JwtTokenUtil jwtTokenUtil; @Override protected void doFilterInternal(HttpServletRequest request, HttpServletResponse response, FilterChain chain) throws ServletException, IOException { final String requestTokenHeader = request.getHeader("Authorization"); String username = null; String jwtToken = null; if (requestTokenHeader != null && requestTokenHeader.startsWith("Bearer ")) { jwtToken = requestTokenHeader.substring(7); try { username = jwtTokenUtil.getUsernameFromToken(jwtToken); } catch (IllegalArgumentException e) { logger.error("无法获取JWT Token", e); } catch (ExpiredJwtException e) { logger.error("JWT Token已过期", e); } } if (username != null && SecurityContextHolder.getContext().getAuthentication() == null) { UserDetails userDetails = this.userDetailsService.loadUserByUsername(username); if (jwtTokenUtil.validateToken(jwtToken, userDetails)) { UsernamePasswordAuthenticationToken authToken = new UsernamePasswordAuthenticationToken( userDetails, null, userDetails.getAuthorities()); authToken.setDetails(new WebAuthenticationDetailsSource().buildDetails(request)); SecurityContextHolder.getContext().setAuthentication(authToken); } } chain.doFilter(request, response); }}### 2. OAuth2集成java// OAuth2配置@Configuration@EnableWebSecurity@EnableGlobalMethodSecurity(prePostEnabled = true)public class SecurityConfig { @Bean public SecurityFilterChain filterChain(HttpSecurity http) throws Exception { http.csrf().disable() .sessionManagement().sessionCreationPolicy(SessionCreationPolicy.STATELESS) .and() .authorizeHttpRequests(authz -> authz .requestMatchers("/api/auth/**").permitAll() .requestMatchers("/actuator/health").permitAll() .requestMatchers(HttpMethod.GET, "/api/users/**").hasRole("USER") .requestMatchers(HttpMethod.POST, "/api/users/**").hasRole("ADMIN") .anyRequest().authenticated() ) .oauth2ResourceServer(oauth2 -> oauth2 .jwt(jwt -> jwt .decoder(jwtDecoder()) .jwtAuthenticationConverter(jwtAuthenticationConverter()) ) ); return http.build(); } @Bean public JwtDecoder jwtDecoder() { return NimbusJwtDecoder.withJwkSetUri("http://auth-server:8080/oauth2/jwks") .build(); } @Bean public JwtAuthenticationConverter jwtAuthenticationConverter() { JwtGrantedAuthoritiesConverter authoritiesConverter = new JwtGrantedAuthoritiesConverter(); authoritiesConverter.setAuthorityPrefix("ROLE_"); authoritiesConverter.setAuthoritiesClaimName("roles"); JwtAuthenticationConverter converter = new JwtAuthenticationConverter(); converter.setJwtGrantedAuthoritiesConverter(authoritiesConverter); return converter; }}## 部署与运维### 1. Docker容器化dockerfile# 多阶段构建FROM openjdk:17-jdk-slim as builderWORKDIR /appCOPY pom.xml .COPY src ./src# 构建应用RUN ./mvnw clean package -DskipTests# 运行时镜像FROM openjdk:17-jre-slim# 创建非root用户RUN groupadd -r appuser && useradd -r -g appuser appuserWORKDIR /app# 复制jar文件COPY --from=builder /app/target/*.jar app.jar# 设置文件权限RUN chown appuser:appuser app.jar# 切换到非root用户USER appuser# 健康检查HEALTHCHECK --interval=30s --timeout=3s --start-period=60s --retries=3 \ CMD curl -f http://localhost:8080/actuator/health || exit 1# 启动应用ENTRYPOINT ["java", "-jar", "app.jar"]### 2. Kubernetes部署yaml# deployment.yamlapiVersion: apps/v1kind: Deploymentmetadata: name: user-service labels: app: user-servicespec: replicas: 3 selector: matchLabels: app: user-service template: metadata: labels: app: user-service spec: containers: - name: user-service image: user-service:1.0.0 ports: - containerPort: 8080 env: - name: SPRING_PROFILES_ACTIVE value: "kubernetes" - name: DB_HOST valueFrom: secretKeyRef: name: user-service-secret key: db-host - name: DB_PASSWORD valueFrom: secretKeyRef: name: user-service-secret key: db-password resources: requests: memory: "512Mi" cpu: "250m" limits: memory: "1Gi" cpu: "500m" livenessProbe: httpGet: path: /actuator/health/liveness port: 8080 initialDelaySeconds: 60 periodSeconds: 30 readinessProbe: httpGet: path: /actuator/health/readiness port: 8080 initialDelaySeconds: 30 periodSeconds: 10---apiVersion: v1kind: Servicemetadata: name: user-servicespec: selector: app: user-service ports: - port: 80 targetPort: 8080 type: ClusterIP---apiVersion: autoscaling/v2kind: HorizontalPodAutoscalermetadata: name: user-service-hpaspec: scaleTargetRef: apiVersion: apps/v1 kind: Deployment name: user-service minReplicas: 3 maxReplicas: 10 metrics: - type: Resource resource: name: cpu target: type: Utilization averageUtilization: 70 - type: Resource resource: name: memory target: type: Utilization averageUtilization: 80## 最佳实践### 1. 服务拆分原则java// 按业务能力拆分// 用户管理服务@RestController@RequestMapping("/api/users")public class UserController { // 用户CRUD操作 // 用户认证 // 用户权限管理}// 订单管理服务@RestController@RequestMapping("/api/orders")public class OrderController { // 订单创建 // 订单状态管理 // 订单查询}// 支付服务@RestController@RequestMapping("/api/payments")public class PaymentController { // 支付处理 // 退款处理 // 支付状态查询}// 库存服务@RestController@RequestMapping("/api/inventory")public class InventoryController { // 库存查询 // 库存预留 // 库存释放}### 2. 数据一致性策略java// 最终一致性实现@Servicepublic class OrderConsistencyService { @Autowired private OrderRepository orderRepository; @Autowired private EventPublisher eventPublisher; // 订单状态变更 @Transactional public void updateOrderStatus(Long orderId, OrderStatus newStatus) { Order order = orderRepository.findById(orderId) .orElseThrow(() -> new OrderNotFoundException(orderId)); OrderStatus oldStatus = order.getStatus(); order.setStatus(newStatus); order.setUpdatedAt(LocalDateTime.now()); orderRepository.save(order); // 发布状态变更事件 OrderStatusChangedEvent event = new OrderStatusChangedEvent( orderId, oldStatus, newStatus, LocalDateTime.now() ); eventPublisher.publish(event); } // 补偿机制 @EventListener public void handlePaymentFailed(PaymentFailedEvent event) { try { updateOrderStatus(event.getOrderId(), OrderStatus.PAYMENT_FAILED); } catch (Exception e) { log.error("处理支付失败事件异常", e); // 重试机制或人工介入 } }}### 3. 性能优化java// 缓存策略@Servicepublic class UserService { @Autowired private UserRepository userRepository; @Autowired private RedisTemplate<String, User> redisTemplate; private static final String USER_CACHE_KEY = "user:"; private static final Duration CACHE_TTL = Duration.ofHours(1); public User findById(Long userId) { // 1. 查询缓存 String cacheKey = USER_CACHE_KEY + userId; User cachedUser = redisTemplate.opsForValue().get(cacheKey); if (cachedUser != null) { return cachedUser; } // 2. 查询数据库 User user = userRepository.findById(userId) .orElseThrow(() -> new UserNotFoundException(userId)); // 3. 更新缓存 redisTemplate.opsForValue().set(cacheKey, user, CACHE_TTL); return user; } @CacheEvict(value = "users", key = "#userId") public void updateUser(Long userId, UpdateUserRequest request) { User user = findById(userId); // 更新用户信息 userRepository.save(user); // 清除缓存 String cacheKey = USER_CACHE_KEY + userId; redisTemplate.delete(cacheKey); }}// 连接池优化@Configurationpublic class DatabaseConfig { @Bean @Primary public DataSource dataSource() { HikariConfig config = new HikariConfig(); config.setJdbcUrl("jdbc:postgresql://localhost:5432/userdb"); config.setUsername("username"); config.setPassword("password"); // 连接池配置 config.setMaximumPoolSize(20); config.setMinimumIdle(5); config.setConnectionTimeout(30000); config.setIdleTimeout(600000); config.setMaxLifetime(1800000); config.setLeakDetectionThreshold(60000); return new HikariDataSource(config); }}## 总结微服务架构设计与实践是一个复杂的系统工程,需要综合考虑:### 核心要点1. 服务拆分:按业务能力进行合理拆分2. 通信机制:选择合适的同步/异步通信方式3. 服务治理:实现服务发现、负载均衡、熔断降级4. 数据管理:保证数据一致性和事务完整性5. 配置管理:集中化配置和动态更新6. 监控运维:完善的监控、日志和链路追踪7. 安全防护:认证授权和API安全8. 部署运维:容器化和自动化部署### 最佳实践- 渐进式迁移:从单体到微服务的平滑过渡- 领域驱动:基于业务领域进行服务划分- API设计:遵循RESTful规范和版本管理- 容错设计:实现熔断、重试、降级机制- 性能优化:合理使用缓存和连接池- 持续集成:自动化测试和部署流水线通过系统性的设计和实践,可以构建出高可用、高性能、易维护的微服务架构系统。大数据
更多推荐



所有评论(0)