你的ChatBI(问数)准确率不到%?带你深度拆解%准确率的高德ChatBI案例
亩兴哪傲数据库是系统的核心,也是千万级系统最难扩展的部分。
读写分离
@Configuration
public class DataSourceConfig {
@Bean
@Primary
public DataSource dataSource() {
// 主从数据源配置
Map targetDataSources = new HashMap<>();
// 主库
DataSource masterDataSource = createDataSource(
"jdbc:mysql://master-db:3306/db?useSSL=false",
"master_user",
"master_password"
);
targetDataSources.put(DataSourceType.MASTER, masterDataSource);
// 从库1
DataSource slave1DataSource = createDataSource(
"jdbc:mysql://slave1-db:3306/db?useSSL=false",
"slave_user",
"slave_password"
);
targetDataSources.put("slave1", slave1DataSource);
// 从库2
DataSource slave2DataSource = createDataSource(
"jdbc:mysql://slave2-db:3306/db?useSSL=false",
"slave_user",
"slave_password"
);
targetDataSources.put("slave2", slave2DataSource);
// 动态数据源
DynamicDataSource dynamicDataSource = new DynamicDataSource();
dynamicDataSource.setDefaultTargetDataSource(masterDataSource);
dynamicDataSource.setTargetDataSources(targetDataSources);
return dynamicDataSource;
}
@Bean
public AbstractRoutingDataSource routingDataSource() {
return new AbstractRoutingDataSource() {
@Override
protected Object determineCurrentLookupKey() {
return DynamicDataSourceContextHolder.getDataSourceType();
}
};
}
@Aspect
@Component
public class DataSourceAspect {
@Before("@annotation(com.example.annotation.Master)")
public void setMasterDataSource() {
DynamicDataSourceContextHolder.setDataSourceType(DataSourceType.MASTER);
}
@Before("@annotation(com.example.annotation.Slave)")
public void setSlaveDataSource() {
// 随机选择从库
String[] slaves = {"slave1", "slave2"};
String selectedSlave = slaves[ThreadLocalRandom.current().nextInt(slaves.length)];
DynamicDataSourceContextHolder.setDataSourceType(selectedSlave);
}
@After("@annotation(com.example.annotation.Master) || " +
"@annotation(com.example.annotation.Slave)")
public void clearDataSource() {
DynamicDataSourceContextHolder.clearDataSourceType();
}
}
}
// 使用示例
@Service
public class OrderService {
@Master // 使用主库
public void createOrder(Order order) {
orderRepository.save(order); // 写操作
}
@Slave // 使用从库
public Order getOrder(Long orderId) {
return orderRepository.findById(orderId).orElse(null); // 读操作
}
}
分库分表
// 分片策略:用户ID取模分片
@Component
public class UserShardingStrategy implements PreciseShardingAlgorithm {
@Override
public String doSharding(Collection availableTargetNames,
PreciseShardingValue shardingValue) {
long userId = shardingValue.getValue();
// 分片数
int shardCount = availableTargetNames.size();
// 简单取模分片
long shardIndex = userId % shardCount;
for (String tableName : availableTargetNames) {
if (tableName.endsWith("_" + shardIndex)) {
return tableName;
}
}
throw new UnsupportedOperationException("无法找到对应的分片表");
}
}
// 分库分表配置
@Configuration
public class ShardingConfig {
@Bean
public DataSource shardingDataSource() throws SQLException {
// 数据源映射
Map dataSourceMap = new HashMap<>();
dataSourceMap.put("ds0", createDataSource("jdbc:mysql://db0:3306/db"));
dataSourceMap.put("ds1", createDataSource("jdbc:mysql://db1:3306/db"));
dataSourceMap.put("ds2", createDataSource("jdbc:mysql://db2:3306/db"));
// 用户表分片规则
ShardingRuleConfiguration shardingRuleConfig = new ShardingRuleConfiguration();
// 用户表分表规则
TableRuleConfiguration userTableRuleConfig = new TableRuleConfiguration("user", "ds${0..2}.user_${0..7}");
userTableRuleConfig.setDatabaseShardingStrategyConfig(
new InlineShardingStrategyConfiguration("id", "ds${id % 3}")
);
userTableRuleConfig.setTableShardingStrategyConfig(
new StandardShardingStrategyConfiguration("id", new UserShardingStrategy())
);
// 订单表分表规则
TableRuleConfiguration orderTableRuleConfig = new TableRuleConfiguration("order", "ds${0..2}.order_${0..15}");
orderTableRuleConfig.setDatabaseShardingStrategyConfig(
new InlineShardingStrategyConfiguration("user_id", "ds${user_id % 3}")
);
orderTableRuleConfig.setTableShardingStrategyConfig(
new InlineShardingStrategyConfiguration("order_id", "order_${order_id % 16}")
);
shardingRuleConfig.getTableRuleConfigs().add(userTableRuleConfig);
shardingRuleConfig.getTableRuleConfigs().add(orderTableRuleConfig);
// 创建ShardingSphere数据源
return ShardingSphereDataSourceFactory.createDataSource(
dataSourceMap, Collections.singleton(shardingRuleConfig), new Properties()
);
}
}
数据库优化实战
-- 1. 索引优化示例
-- 错误的索引设计
CREATE INDEX idx_user_email ON user(email); -- 过长的索引
CREATE INDEX idx_user_status ON user(status); -- 低区分度的索引
-- 正确的索引设计
-- 联合索引,注意字段顺序
CREATE INDEX idx_user_created_status ON user(created_at, status);
-- 前缀索引,适合长字段
CREATE INDEX idx_user_email_prefix ON user(email(20));
-- 2. 慢查询优化示例
-- 优化前:全表扫描
EXPLAIN SELECT * FROM order WHERE YEAR(created_at) = 2024;
-- 优化后:使用索引
EXPLAIN SELECT * FROM order
WHERE created_at >= '2024-01-01'
AND created_at < '2025-01-01';
-- 3. 分页优化
-- 传统分页(数据量大时慢)
SELECT * FROM user ORDER BY id LIMIT 1000000, 20;
-- 优化分页(使用覆盖索引)
SELECT * FROM user
WHERE id > (SELECT id FROM user ORDER BY id LIMIT 1000000, 1)
ORDER BY id LIMIT 20;
-- 4. 批量操作优化
-- 批量插入
INSERT INTO user (name, email) VALUES
('user1', 'user1@example.com'),
('user2', 'user2@example.com'),
-- ... 1000条
('user1000', 'user1000@example.com');
-- 批量更新(避免在循环中单条更新)
UPDATE user SET status = 'active'
WHERE id IN (1, 2, 3, ..., 1000);
06 异步处理与消息队列
对于千万级系统,同步处理所有请求是不现实的。异步化是提高系统吞吐量的关键。
消息队列设计
@Configuration
public class KafkaConfig {
@Bean
public ProducerFactory producerFactory() {
Map configProps = new HashMap<>();
configProps.put(
ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,
"kafka1:9092,kafka2:9092,kafka3:9092"
);
configProps.put(
ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,
StringSerializer.class
);
configProps.put(
ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
StringSerializer.class
);
// 高吞吐量配置
configProps.put(ProducerConfig.LINGER_MS_CONFIG, 20); // 批量发送延迟
configProps.put(ProducerConfig.BATCH_SIZE_CONFIG, 32 * 1024); // 批量大小
configProps.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "lz4"); // 压缩
// 高可靠性配置
configProps.put(ProducerConfig.ACKS_CONFIG, "all"); // 所有副本确认
configProps.put(ProducerConfig.RETRIES_CONFIG, 3); // 重试次数
configProps.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true); // 幂等性
return new DefaultKafkaProducerFactory<>(configProps);
}
@Bean
public KafkaTemplate kafkaTemplate() {
return new KafkaTemplate<>(producerFactory());
}
}
@Service
public class OrderMessageService {
@Autowired
private KafkaTemplate kafkaTemplate;
// 顺序消息发送
public void sendOrderEvent(OrderEvent event) {
// 使用订单ID作为key,保证同一订单的消息顺序
String key = String.valueOf(event.getOrderId());
String topic = "order-events";
kafkaTemplate.send(topic, key, JsonUtils.toJson(event))
.addCallback(
result -> log.info("消息发送成功: {}", event),
ex -> {
log.error("消息发送失败: {}", event, ex);
// 失败处理:记录到数据库,定时重试
saveFailedMessage(event, ex.getMessage());
}
);
}
// 批量消息处理
@KafkaListener(topics = "order-events",
containerFactory = "batchFactory")
public void processOrderEvents(List> records) {
List events = new ArrayList<>();
for (ConsumerRecord record : records) {
try {
OrderEvent event = JsonUtils.fromJson(record.value(), OrderEvent.class);
events.add(event);
} catch (Exception e) {
log.error("消息解析失败: {}", record.value(), e);
}
}
if (!events.isEmpty()) {
// 批量处理
batchProcessOrders(events);
}
}
@Bean
public ConcurrentKafkaListenerContainerFactory batchFactory() {
ConcurrentKafkaListenerContainerFactory factory =
new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(consumerFactory());
factory.setBatchListener(true); // 开启批量监听
factory.getContainerProperties().setPollTimeout(3000);
return factory;
}
}
异步处理模式
// 1. CompletableFuture异步处理
@Service
public class AsyncOrderService {
@Autowired
private ThreadPoolTaskExecutor taskExecutor;
public CompletableFuture createOrderAsync(OrderRequest request) {
return CompletableFuture.supplyAsync(() -> {
// 第一步:校验
validateOrder(request);
return request;
}, taskExecutor).thenComposeAsync(validatedRequest -> {
// 第二步:扣减库存
return reduceInventoryAsync(validatedRequest);
}, taskExecutor).thenComposeAsync(inventoryResult -> {
// 第三步:生成订单
return generateOrderAsync(inventoryResult);
}, taskExecutor).thenApplyAsync(order -> {
// 第四步:发送通知
sendNotification(order);
return new OrderResult(order, "SUCCESS");
}, taskExecutor).exceptionally(ex -> {
// 异常处理
log.error("订单创建失败", ex);
return new OrderResult(null, "FAILED: " + ex.getMessage());
});
}
// 2. 事件驱动架构
@Component
public class OrderEventPublisher {
@Autowired
private ApplicationEventPublisher eventPublisher;
public void createOrder(Order order) {
// 保存订单
orderRepository.save(order);
// 发布领域事件
eventPublisher.publishEvent(new OrderCreatedEvent(this, order));
}
}
@Component
public class OrderEventListener {
@EventListener
@Async
public void handleOrderCreated(OrderCreatedEvent event) {
Order order = event.getOrder();
// 异步发送邮件
emailService.sendOrderConfirmation(order);
// 异步更新统计数据
statisticsService.updateOrderStats(order);
// 异步通知库存系统
inventoryService.notifyOrderCreated(order);
}
}
}
07 监控与治理:系统的眼睛和大脑
没有监控的系统就像盲人开车。千万级系统需要完善的监控体系。
全链路监控
// 1. 链路追踪(基于SkyWalking)
@Configuration
public class TracingConfig {
@Bean
public Tracer skywalkingTracer() {
// SkyWalking自动注入,只需添加依赖
return new SkywalkingTracer();
}
}
// 在业务代码中自动追踪
@Service
public class ProductService {
@Trace(operationName = "ProductService.getProduct") // SkyWalking注解
@Tag(key = "productId", value = "arg[0]") // 添加标签
public Product getProduct(Long productId) {
// 业务逻辑
return productRepository.findById(productId).orElse(null);
}
}
// 2. 指标监控(Micrometer + Prometheus)
@Configuration
public class MetricsConfig {
@Bean
public MeterRegistry meterRegistry() {
return new PrometheusMeterRegistry(PrometheusConfig.DEFAULT);
}
@Bean
public TimedAspect timedAspect(MeterRegistry registry) {
return new TimedAspect(registry);
}
}
@Service
public class OrderService {
private final Counter orderCounter;
private final Timer orderTimer;
public OrderService(MeterRegistry registry) {
orderCounter = Counter.builder("orders.created")
.description("创建的订单数量")
.tag("type", "normal")
.register(registry);
orderTimer = Timer.builder("orders.process.time")
.description("订单处理时间")
.publishPercentiles(0.5, 0.95, 0.99) // 50%, 95%, 99%分位数
.register(registry);
}
@Timed(value = "orders.create", extraTags = {"priority", "high"})
public Order createOrder(OrderRequest request) {
return orderTimer.record(() -> {
Order order = processOrder(request);
orderCounter.increment();
return order;
});
}
}
// 3. 健康检查
@Component
public class CustomHealthIndicator implements HealthIndicator {
@Autowired
private DataSource dataSource;
@Autowired
private RedisConnectionFactory redisConnectionFactory;
@Override
public Health health() {
// 检查数据库连接
boolean dbHealthy = checkDatabase();
// 检查Redis连接
boolean redisHealthy = checkRedis();
// 检查磁盘空间
boolean diskHealthy = checkDiskSpace();
if (dbHealthy && redisHealthy && diskHealthy) {
return Health.up()
.withDetail("database", "connected")
.withDetail("redis", "connected")
.withDetail("disk", "sufficient")
.build();
} else {
Map details = new HashMap<>();
details.put("database", dbHealthy ? "connected" : "disconnected");
details.put("redis", redisHealthy ? "connected" : "disconnected");
details.put("disk", diskHealthy ? "sufficient" : "insufficient");
return Health.down().withDetails(details).build();
}
}
private boolean checkDatabase() {
try {
return dataSource.getConnection().isValid(5);
} catch (SQLException e) {
return false;
}
}
}
智能告警
# Prometheus告警规则
groups:
- name: application_alerts
rules:
# 错误率告警
- alert: HighErrorRate
expr: rate(http_requests_total{status=~"5.."}[5m]) / rate(http_requests_total[5m]) > 0.01
for: 2m
labels:
severity: critical
annotations:
summary: "高错误率告警"
description: "错误率超过1% (当前值: {{ $value }})"
# 延迟告警
- alert: HighLatency
expr: histogram_quantile(0.95, rate(http_request_duration_seconds_bucket[5m])) > 1
for: 5m
labels:
severity: warning
annotations:
summary: "高延迟告警"
description: "95分位延迟超过1秒 (当前值: {{ $value }}s)"
# 服务宕机告警
- alert: ServiceDown
expr: up == 0
for: 1m
labels:
severity: critical
annotations:
summary: "服务宕机告警"
description: "{{ $labels.job }} 服务已宕机"
08 实战案例:秒杀系统设计
让我们用一个具体的例子来综合运用上述技术。
@Service
public class SeckillService {
// 1. 缓存预热
@PostConstruct
public void warmUpCache() {
List hotProducts = findHotProducts();
hotProducts.parallelStream().forEach(product -> {
String stockKey = "seckill:stock:" + product.getId();
redisTemplate.opsForValue().set(stockKey, product.getStock());
// 设置随机过期时间,防止雪崩
setWithRandomExpire(stockKey, product.getStock(), 10, TimeUnit.MINUTES);
});
}
// 2. 秒杀入口(防刷+限流)
@RateLimiter(key = "'seckill:' + #userId",
limit = 1, period = 10, unit = TimeUnit.SECONDS)
public SeckillResult seckill(Long userId, Long productId) {
// 2.1 资格校验(Redis实现)
String userKey = "seckill:user:" + productId + ":" + userId;
if (Boolean.TRUE.equals(redisTemplate.hasKey(userKey))) {
return SeckillResult.alreadyParticipated();
}
// 2.2 库存预扣减(Lua脚本保证原子性)
String stockKey = "seckill:stock:" + productId;
Long remaining = redisTemplate.execute(
stockDecrScript,
Collections.singletonList(stockKey),
"1"
);
if (remaining < 0) {
// 库存不足,回滚
redisTemplate.opsForValue().increment(stockKey, 1);
return SeckillResult.outOfStock();
}
// 2.3 记录用户购买资格
redisTemplate.opsForValue().set(userKey, "1", 1, TimeUnit.HOURS);
// 2.4 发送消息到队列,异步创建订单
String orderKey = "seckill:order:" + UUID.randomUUID();
SeckillOrderMessage message = new SeckillOrderMessage(userId, productId, orderKey);
kafkaTemplate.send("seckill-orders",
String.valueOf(productId),
JsonUtils.toJson(message)
);
return SeckillResult.success(orderKey);
}
// 3. 异步订单处理
@KafkaListener(topics = "seckill-orders",
groupId = "seckill-order-processor")
public void processSeckillOrder(String message) {
SeckillOrderMessage orderMessage = JsonUtils.fromJson(message, SeckillOrderMessage.class);
try {
// 3.1 创建订单(数据库操作)
Order order = createOrderInDB(orderMessage);
// 3.2 更新缓存状态
String resultKey = "seckill:result:" + orderMessage.getOrderKey();
redisTemplate.opsForValue().set(resultKey, order.getId(), 10, TimeUnit.MINUTES);
// 3.3 发送成功通知
notifyUser(orderMessage.getUserId(), order);
} catch (Exception e) {
log.error("处理秒杀订单失败", e);
// 失败处理:恢复库存
String stockKey = "seckill:stock:" + orderMessage.getProductId();
redisTemplate.opsForValue().increment(stockKey, 1);
// 移除用户资格记录
String userKey = "seckill:user:" + orderMessage.getProductId() + ":" + orderMessage.getUserId();
redisTemplate.delete(userKey);
}
}
// 4. 库存扣减Lua脚本
private final DefaultRedisScript stockDecrScript = new DefaultRedisScript<>(
"local key = KEYS[1]\n" +
"local decrement = tonumber(ARGV[1])\n" +
"local current = redis.call('get', key)\n" +
"if not current then\n" +
" return -1\n" +
"end\n" +
"local currentNum = tonumber(current)\n" +
"if currentNum < decrement then\n" +
" return -1\n" +
"end\n" +
"return redis.call('decrby', key, decrement)",
Long.class
);
}
总结
设计一个扛住千万级流量的系统,不是简单堆砌技术组件,而是构建一个有机的、能自我调节的生态系统。
通过本文的分享,我想你已经看到了一个完整的高并发系统架构图景。
让我最后总结一下关键点:
架构是演进而来的:不要一开始就追求完美架构,而是随着业务增长不断演进。从小而美的单体开始,逐步拆分、优化、扩展。
缓存是性能的银弹:但要用得聪明。多级缓存、智能淘汰、防止穿透/击穿/雪崩,每个细节都影响巨大。
数据库是瓶颈所在:读写分离、分库分表、索引优化、SQL调优,这些基本功比任何炫酷的技术都重要。
异步是扩展的关键:能异步的绝不同步,能批处理的绝不单条。消息队列、事件驱动、反应式编程,让系统松耦合、高内聚。
监控是系统的眼睛:没有监控的系统就是在裸奔。全链路追踪、指标监控、日志聚合、智能告警,缺一不可。
容错比优化更重要:系统一定会出问题,关键是如何快速发现、快速恢复、快速止损。熔断、降级、限流、重试,这些机制是系统的保险绳。
