当前位置: 首页 > news >正文

你的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调优,这些基本功比任何炫酷的技术都重要。

异步是扩展的关键:能异步的绝不同步,能批处理的绝不单条。消息队列、事件驱动、反应式编程,让系统松耦合、高内聚。

监控是系统的眼睛:没有监控的系统就是在裸奔。全链路追踪、指标监控、日志聚合、智能告警,缺一不可。

容错比优化更重要:系统一定会出问题,关键是如何快速发现、快速恢复、快速止损。熔断、降级、限流、重试,这些机制是系统的保险绳。

http://www.cnnetsun.cn/news/1261095.html

相关文章:

  • 想用 Claude Code 做 AI 编程,很多人其实卡在了接入这一步
  • 做 AI 测试用例系统时,Prompt、MCP、Agent、Skills、OpenClaw 到底分别是什么?
  • 鸿蒙常见问题分析三十三:如何解决Column子组件超出容器边界
  • JUC 并发编程:对可见性、有序性与 volatile的理解
  • RAG检索瓶颈突破实战指南(非常详细),Multi-HyDE与Adaptive HyDE从入门到精通,收藏这一篇就够了!
  • 实现一下简单的聊天功能,包括聊天消息自适应大小
  • KohakuRAG:层次化RAG的新范式
  • 关于DBeaver的一些配置
  • 性能优化-前端性能优化相关
  • asp毕业设计——基于asp+access的网上选题系统设计与实现(毕业论文+程序源码)——网上选题系统
  • 【OS】操作系统分类及用户界面
  • 探索 Stanford Alpaca: 一个强大的深度学习框架
  • 探索超轻量级人脸识别神器:Ultra Light Fast Generic Face Detector 1MB
  • 【亲测免费】 推荐一款强大的数据可视化库——Plotly.js
  • MAPPO动作类型改进(二)——MAPPO+连续环境
  • IPED元数据可视化案例:创建取证报告中的元数据图表
  • Apache Airflow 项目教程
  • Carbon 开源项目指南
  • YTKNetwork批量请求终极指南:YTKBatchRequest高效应用实战
  • OpenCart 开源电商系统推荐
  • 新手到大卖都在用!亚马逊运营全流程工具盘点
  • 开源项目推荐:ccv
  • YTKNetwork性能优化终极指南:10个提升iOS网络请求效率的实用方法
  • 深入理解粤语编程编译器:从Python转换到LLVM执行
  • LCD屏幕调光方式全解析:DC调光 vs PWM调光,哪种更适合你的眼睛?
  • DS1302实时时钟模块从入门到精通:手把手教你搭建可调时钟系统
  • 用Multisim仿真8种经典运放电路:手把手教你搭建比例/微分/积分放大器
  • Android11音频架构剖析:RK3568平台如何通过RK809实现智能功放切换
  • 告别DHCP!CentOS服务器静态IP设置避坑指南(附常见问题解决方案)
  • 2026年广东省职业院校技能大赛(高职组)移动应用设计与开发赛项样题(一)