拓十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

O2O四端微服务架构:SpringCloud Alibaba多端隔离与实时聚合

O2O四端微服务架构:SpringCloud Alibaba多端隔离与实时聚合

简介:本资源是一套基于SpringCloud微服务架构的O2O外卖系统后端数据库设计实现,面向Java后端开发者、微服务学习者及电商平台项目实践者,聚焦饿了么类业务场景下的多端协同与数据建模问题。压缩包为ZIP格式,共3个SQL文件(65KB),分别对应订单调度(db_dis_order.sql)、客户订单(db_cus_order.sql)和商家中心(shangcheng.sql)三大核心模块,完整覆盖用户、商家、骑手、订单等关键实体关系与事务逻辑,具备直接导入MySQL运行的基础结构与初始化数据。已有485人学习下载,适合用于微服务项目数据库层参考、SpringCloud多模块联调的数据准备,以及理解O2O平台中高并发订单与分布式事务下的表结构设计思路。

1. 这不是“仿饿了么”Demo,而是一套真实可压测、可拆分、可演进的O2O多端协同后端骨架:SpringCloud + 四端隔离 + 总后台聚合

你在网上搜“饿了么外卖系统 SpringCloud”,90% 的结果是单体 SpringBoot + 前端 Vue 的教学 Demo,订单表硬编码在 user_service 里,配送员和商家共用一个登录接口,连 Redis 缓存都只配了一个 localhost:6379。但真实 O2O 系统的后端根本不是这样——它必须支撑客户端(C端用户)日均千万级请求、商家端(B端)高并发上架/下架、配送端(Rider)毫秒级位置上报与调度响应、订单端(Order)强一致性事务与状态机驱动、总后台(Admin)跨域数据聚合与实时监控。这套架构不是为“跑通功能”设计的,而是为“扛住峰值、快速迭代、独立演进”设计的。它用 SpringCloud Alibaba(Nacos + Sentinel + Seata)做服务治理底座,四端业务逻辑物理隔离、API 网关层路由分流、数据库按端垂直分库+水平分表(订单库按月分表+用户ID哈希分片),总后台不直连业务库,只通过 Feign 调用聚合服务或 Kafka 消费埋点数据。如果你正要从单体迁移到微服务,或正在设计新 O2O 平台的后端基线,这篇笔记就是你跳过踩坑、直接复用的最小可行骨架。


2. 用 SpringCloud Alibaba 搭建四端服务注册与发现:Nacos 配置中心 + 多环境隔离 + 服务健康检查

2.1 为什么选 Nacos 而非 Eureka?——真实生产中“注册中心不可用”比“服务宕机”更致命

Eureka 的自我保护模式在大规模节点抖动时会盲目保留大量失效实例,导致流量打到已下线服务;ZooKeeper 的 CP 特性在脑裂场景下直接拒绝写入,网关无法动态更新路由。Nacos 同时支持 AP(服务发现)和 CP(配置管理),且提供控制台可视化心跳探测、权重灰度、元数据标签路由——这正是四端协同的核心需求:比如配送端 RiderService 必须优先调用同机房的 DispatchService,商家端需按城市 ID 标签路由到对应区域的 MerchantService。我们不用@EnableEurekaClient,而是统一依赖spring-cloud-starter-alibaba-nacos-discovery,每个模块的bootstrap.yml中强制指定 namespace(命名空间)隔离环境:

# client-service/bootstrap.yml spring: cloud: nacos: discovery: server-addr: 192.168.10.50:8848 namespace: 7a8b9c0d-1e2f-3a4b-5c6d-7e8f9a0b1c2d # dev 环境 namespace ID group: CLIENT_GROUP cluster-name: SHANGHAI config: server-addr: 192.168.10.50:8848 namespace: 7a8b9c0d-1e2f-3a4b-5c6d-7e8f9a0b1c2d group: CLIENT_GROUP file-extension: yaml

提示:namespace 是 Nacos 最关键的隔离单元,不是用profile区分环境,而是用 namespace ID。dev/test/prod 各自独立 namespace,避免测试环境误删生产配置。group 用于同一环境下按端划分(CLIENT_GROUP / MERCHANT_GROUP / RIDER_GROUP / ORDER_GROUP),便于权限管控。

2.2 四端服务命名规范与元数据注入:让网关和熔断器“看懂”你是谁

单纯靠服务名client-service不足以表达语义。我们在每个服务启动时注入元数据,供 SpringCloud Gateway 动态路由和 Sentinel 流控规则识别:

// client-service/src/main/java/com/eleme/client/config/NacosMetadataConfig.java @Configuration public class NacosMetadataConfig { @Bean @ConditionalOnMissingBean public Registration registration() { return new NacosRegistration() { @Override public Map<String, String> getMetadata() { Map<String, String> metadata = new HashMap<>(); metadata.put("service-type", "client"); // 标识 C端 metadata.put("biz-domain", "user"); // 用户域 metadata.put("version", "v2.3.1"); // 接口版本 metadata.put("qps-threshold", "5000"); // 基准 QPS return metadata; } }; } }

这样,Gateway 的 Predicate 就能写成:

spring: cloud: gateway: routes: - id: client-api uri: lb://client-service predicates: - Header[X-Client-Type], client # 只转发带此 header 的请求 - Metadata[service-type], client # 或直接匹配元数据

而 Sentinel 控制台就能基于service-type=client统一配置 C端全局流控规则,无需为每个接口单独设限。

2.3 服务健康检查的“真存活”验证:不只是 ping,而是查 DB + Redis + MQ 连通性

Nacos 默认只检测 HTTP/actuator/health端点,但该端点可能返回 UP 却实际连不上 MySQL。我们重写HealthIndicator,组合校验:

@Component public class CompositeHealthIndicator implements HealthIndicator { private final JdbcTemplate jdbcTemplate; private final RedisTemplate redisTemplate; private final KafkaTemplate kafkaTemplate; @Override public Health health() { try { // 1. 检查 MySQL:执行轻量 SELECT jdbcTemplate.queryForObject("SELECT 1", Integer.class); // 2. 检查 Redis:SETNX key redisTemplate.opsForValue().setIfAbsent("health:check", "ok", Duration.ofSeconds(1)); // 3. 检查 Kafka:发送测试消息(异步,超时 2s) kafkaTemplate.send("health-test", "ping").get(2, TimeUnit.SECONDS); return Health.up().withDetail("db", "OK").withDetail("redis", "OK").withDetail("kafka", "OK").build(); } catch (Exception e) { return Health.down().withDetail("error", e.getMessage()).build(); } } }

Nacos 会每 5 秒调用此接口,连续 3 次失败则将实例标记为DOWN并从服务列表剔除。这是防止“假活”流量打到半死服务的关键防线。


3. 四端 API 网关统一入口:SpringCloud Gateway + JWT 鉴权 + 多租户路由 + 请求体缓存

3.1 网关层不做业务鉴权,只做“身份透传”与“端标识校验”

很多团队把 JWT 解析、RBAC 权限校验全堆在 Gateway,结果网关成为性能瓶颈和单点故障。我们的做法是:Gateway 只验证 token 签名有效性、过期时间、并提取client_type(client/merchant/rider/order)、tenant_id(商家 ID 或城市编码),然后以X-Client-Type、X-Tenant-ID等 Header 透传给下游服务。真正的权限校验由各端服务自己完成(如商家端校验tenant_id是否属于当前登录商家),网关只做“准入”而非“授权”。

// gateway/src/main/java/com/eleme/gateway/filter/JwtAuthFilter.java public class JwtAuthFilter implements GlobalFilter { @Override public Mono<Void> filter(ServerWebExchange exchange, GatewayFilterChain chain) { String authHeader = exchange.getRequest().getHeaders().getFirst("Authorization"); if (authHeader == null || !authHeader.startsWith("Bearer ")) { exchange.getResponse().setStatusCode(HttpStatus.UNAUTHORIZED); return exchange.getResponse().setComplete(); } String token = authHeader.substring(7); try { Jws<Claims> claimsJws = Jwts.parserBuilder() .setSigningKey(rsaPublicKey) .build() .parseClaimsJws(token); Claims claims = claimsJws.getBody(); // 提取并透传关键字段 exchange.getRequest().mutate() .headers(h -> { h.set("X-Client-Type", claims.get("client_type", String.class)); h.set("X-Tenant-ID", claims.get("tenant_id", String.class)); h.set("X-User-ID", claims.get("user_id", String.class)); }) .build(); } catch (JwtException e) { exchange.getResponse().setStatusCode(HttpStatus.UNAUTHORIZED); return exchange.getResponse().setComplete(); } return chain.filter(exchange); } }

注意:RSA 公钥必须预加载到内存,禁止每次解析都读文件;claims.get("xxx", String.class)比claims.get("xxx").toString()更安全,避免空指针。

3.2 多租户路由策略:同一商家多个子门店,流量按 tenant_id 哈希到不同实例组

商家端(merchant-service)需支持连锁品牌(如“肯德基上海徐汇店”、“肯德基上海静安店”)共享一套代码但数据物理隔离。我们不在数据库层做 tenant_id 字段过滤,而是在网关层做路由分发:

# gateway/src/main/resources/application.yml spring: cloud: gateway: routes: - id: merchant-service uri: lb://merchant-service predicates: - Header[X-Client-Type], merchant filters: - name: RequestRateLimiter args: redis-rate-limiter.replenishRate: 100 redis-rate-limiter.burstCapacity: 200 - name: DeduplicateResponseHeader args: name: Vary strategy: RETAIN_FIRST # 关键:按 X-Tenant-ID 哈希,路由到不同集群 metadata: tenant-hash: true

配合自定义LoadBalancer实现:

@Bean public ReactorLoadBalancer<ServiceInstance> reactorLoadBalancer( Environment environment, LoadBalancerClientFactory loadBalancerClientFactory) { String serviceId = environment.getProperty(LoadBalancerClientFactory.PROPERTY_NAME); return new TenantHashLoadBalancer(loadBalancerClientFactory.getLazyProvider(serviceId, ServiceInstanceListSupplier.class)); } // TenantHashLoadBalancer.java 内部逻辑: // 1. 从 exchange 获取 X-Tenant-ID // 2. 对 tenant_id 做 MurmurHash3,取模 instance 数量 // 3. 返回对应索引的 ServiceInstance // 保证同一 tenant_id 总打到同一组 merchant-service 实例,利于本地缓存命中

3.3 请求体缓存:解决 POST/PUT 请求体被 Gateway 读取后下游收不到的问题

SpringCloud Gateway 默认不缓存请求体,ServerWebExchange.getRequest().getBody()只能读一次。当需要记录日志、做风控校验、或转发到多个下游时,必须手动缓存:

@Component public class CacheRequestBodyGlobalFilter implements GlobalFilter { @Override public Mono<Void> filter(ServerWebExchange exchange, GatewayFilterChain chain) { ServerHttpRequest request = exchange.getRequest(); if (request.getMethod().equals(HttpMethod.POST) || request.getMethod().equals(HttpMethod.PUT)) { // 缓存 body 到 exchange 属性 return DataBufferUtils.join(request.getBody()) .map(dataBuffer -> { byte[] bytes = new byte[dataBuffer.readableByteCount()]; dataBuffer.read(bytes); DataBufferUtils.release(dataBuffer); return bytes; }) .flatMap(bytes -> { // 存入 exchange 属性,供后续 Filter 或 Controller 使用 exchange.getAttributes().put("cached-request-body", bytes); // 构造新的 cachedRequest Flux<DataBuffer> cachedFlux = Flux.just(exchange.getResponse().bufferFactory().wrap(bytes)); ServerHttpRequest cachedRequest = new ModifyRequestBodyServerHttpRequest(request, cachedFlux); return chain.filter(exchange.mutate().request(cachedRequest).build()); }); } return chain.filter(exchange); } }

下游服务可通过exchange.getAttribute("cached-request-body")获取原始 JSON,避免重复解析或丢失。


4. 订单端强一致性保障:Seata AT 模式 + TCC 补偿 + Saga 状态机驱动

4.1 为什么不用本地事务?——跨库操作天然存在分布式事务问题

下单流程涉及:
① 客户端扣减用户余额(client-db)
② 商家端冻结商品库存(merchant-db)
③ 订单端生成主订单 + 子订单(order-db)
④ 配送端预分配骑手(rider-db)

四个库分布在不同物理节点,MySQL 本地事务无法覆盖。若用消息最终一致性,会出现“用户付了钱但没生成订单”或“库存扣了但订单失败”的资金/库存损失。我们必须用分布式事务中间件。

4.2 Seata AT 模式落地细节:undo_log 表必须与业务表同库,且开启全局事务注解

AT 模式对业务代码侵入最小,但配置极易出错。关键点:

  • 每个业务库(client-db、merchant-db、order-db、rider-db)必须创建undo_log表,且与业务表在同一 schema 下;
  • undo_log表结构必须严格匹配 Seata 1.7+ 版本要求(含branch_id、xid、context字段);
  • 启动类添加@EnableAutoDataSourceProxy(Seata 1.7+ 自动代理数据源);
  • 分布式事务入口方法加@GlobalTransactional(rollbackFor = Exception.class)。
// order-service/src/main/java/com/eleme/order/service/OrderService.java @Service public class OrderServiceImpl implements OrderService { @Autowired private ClientAccountService clientAccountService; // Feign 调用 client-service @Autowired private MerchantInventoryService merchantInventoryService; // Feign 调用 merchant-service @GlobalTransactional(rollbackFor = Exception.class) @Override public Order createOrder(CreateOrderRequest req) { // 1. 扣用户余额(client-service) clientAccountService.deductBalance(req.getUserId(), req.getTotalAmount()); // 2. 冻结库存(merchant-service) merchantInventoryService.freezeStock(req.getMerchantId(), req.getItems()); // 3. 本地生成订单(order-db) Order order = orderMapper.insert(req); // 4. 预分配骑手(rider-service) riderService.preAssignRider(order.getId(), req.getDeliveryAddress()); return order; } }

注意:Feign 调用必须走@GlobalTransactional包裹的方法内部,否则 Seata 无法传播 XID。clientAccountService.deductBalance()内部需有@GlobalLock注解确保行锁。

4.3 TCC 模式兜底:针对无法自动回滚的操作(如调用微信支付)

AT 模式无法处理外部三方服务(微信支付、短信网关)。此时用 TCC:Try-Confirm-Cancel三阶段。

// 微信支付服务(wechat-pay-service) @DubboService public class WechatPayTccService implements TccAction { @Override @TwoPhaseBusinessAction(name = "wechatPayTry", commitMethod = "confirm", rollbackMethod = "cancel") public boolean prepare(BusinessActionContext actionContext, PayRequest req) { // Try 阶段:调用微信统一下单 API,获取 prepay_id,本地记录支付流水(status=TRYING) String prepayId = wechatApi.unifiedOrder(req); payMapper.insert(new PayRecord(req.getOrderId(), prepayId, "TRYING")); return true; } public boolean confirm(BusinessActionContext actionContext) { // Confirm 阶段:查询微信支付结果,更新 status=SUCCESS String orderId = actionContext.getAttachment("orderId"); PayRecord record = payMapper.selectByOrderId(orderId); String result = wechatApi.query(record.getPrepayId()); if ("SUCCESS".equals(result)) { payMapper.updateStatus(orderId, "SUCCESS"); return true; } return false; } public boolean cancel(BusinessActionContext actionContext) { // Cancel 阶段:调用微信关单 API,更新 status=CANCELLED String orderId = actionContext.getAttachment("orderId"); PayRecord record = payMapper.selectByOrderId(orderId); wechatApi.closeOrder(record.getPrepayId()); payMapper.updateStatus(orderId, "CANCELLED"); return true; } }

订单服务在createOrder中调用wechatPayTccService.prepare(...),Seata 自动协调 Confirm/Cancel。

4.4 Saga 状态机:复杂长流程(如售后退款)用 JSON DSL 定义,避免硬编码

退货退款涉及:取消订单 → 释放库存 → 退用户款 → 通知商家 → 关闭物流单。步骤多、失败点分散、需人工干预。我们用 Seata Saga 模式,定义refund.saga.json:

{ "name": "refund-saga", "states": [ { "name": "cancelOrder", "type": "ServiceTask", "serviceName": "order-service", "serviceMethod": "cancelOrder", "compensateServiceName": "order-service", "compensateServiceMethod": "restoreOrder" }, { "name": "releaseStock", "type": "ServiceTask", "serviceName": "merchant-service", "serviceMethod": "releaseStock", "compensateServiceName": "merchant-service", "compensateServiceMethod": "freezeStock" }, { "name": "refundMoney", "type": "ServiceTask", "serviceName": "client-service", "serviceMethod": "refundMoney", "compensateServiceName": "client-service", "compensateServiceMethod": "deductBalance" } ], "startState": "cancelOrder", "endState": "refundMoney" }

Saga 引擎自动执行、失败回滚、持久化状态,比手写状态机代码更可靠、易维护。


5. 避坑指南:SpringCloud O2O 四端架构中 5 个血泪经验总结

5.1 现象:Nacos 配置中心修改后,部分服务未实时刷新,重启才生效

原因:Nacos 客户端默认使用长轮询(Long Polling),但某些云厂商 SLB 或防火墙会中断 30s 以上连接,导致配置变更丢失;同时@RefreshScope注解未加在所有需要刷新的 Bean 上。
解决:
① 在application.yml中显式配置长轮询超时:

spring: cloud: nacos: config: timeout: 10000 # 缩短超时至 10s,避免连接挂起 max-retry: 3 # 失败重试次数

② 所有含@Value("${xxx}")的类必须加@RefreshScope,包括@ConfigurationProperties类;
③ 生产环境禁用@RefreshScope的@Bean方法,改用ApplicationContext.getBean()动态获取。

5.2 现象:Gateway 路由到 merchant-service 后,Feign 调用 client-service 失败,报No instances available for client-service

原因:Feign 默认使用 Ribbon 负载均衡,而 Ribbon 的服务发现依赖DiscoveryClient,但 Gateway 的lb://协议已绕过 Ribbon,直接走 ReactorLoadBalancer;若 merchant-service 未正确引入spring-cloud-starter-loadbalancer,则其内部 Feign 无法发现 client-service。
解决:
① 所有业务服务(client/merchant/rider/order)的pom.xml必须同时引入:

<dependency> <groupId>org.springframework.cloud</groupId> <artifactId>spring-cloud-starter-loadbalancer</artifactId> </dependency> <dependency> <groupId>org.springframework.cloud</groupId> <artifactId>spring-cloud-starter-openfeign</artifactId> </dependency>

② 禁用 Ribbon:spring.cloud.loadbalancer.ribbon.enabled=false;
③ FeignClient 接口必须标注@FeignClient(contextId = "clientService", value = "client-service"),避免 contextId 冲突。

5.3 现象:Seata 全局事务中,order-service 插入订单成功,但 client-service 扣余额失败,order 订单未回滚

原因:Seata AT 模式要求所有参与方数据源必须被 Seata 代理;若 client-service 使用了 Druid 多数据源(如主库+报表库),而只代理了主数据源,则扣余额操作未进入全局事务。
解决:
① 检查client-service的DataSourceProxy是否包装了所有业务数据源;
② 在DruidDataSource初始化后,必须用new DataSourceProxy(druidDataSource)包装;
③ 日志中搜索Branch Register,确认 client-service 是否向 TC 注册了 branch;若无,则说明数据源未被代理。

5.4 现象:Redis 缓存穿透,大量请求击穿到 DB,CPU 突增 90%

原因:商家端商品详情页用GET product:{id}查询,但恶意请求构造不存在的id(如product:999999999),缓存未命中,DB 直接查询返回 null,且未对该 key 设置空值缓存。
解决:
① 所有缓存查询必须做空值缓存(SET product:999999999 "null" EX 60);
② 使用布隆过滤器前置拦截:启动时加载所有有效商品 ID 到布隆过滤器(RedisBloom 模块),查询前先BF.EXISTS product_bf 999999999;
③ 对高频恶意 key(如product:-1、product:abc)做 IP 限流,用 Sentinel@SentinelResource拦截。

5.5 现象:总后台 Admin 查询“昨日订单 TOP10 商家”,响应超时 30s

原因:总后台直接 JOIN 四个库的表(order + merchant + client + rider),MySQL 无法跨库关联,实际执行的是笛卡尔积 + 应用层内存聚合,数据量达百万级。
解决:
① 总后台不直连业务库,改为消费 Kafka 订单埋点 Topic(order-created),用 Flink 实时计算 TOP10,结果写入 Elasticsearch;
② 或使用 ShardingSphere-Proxy 作为数据库中间件,配置broadcast表(merchant_info)和sharding表(order_info),让 SQL 在 Proxy 层下推执行;
③ 禁止总后台写任何SELECT * FROM ... JOIN ...,所有报表查询必须走预聚合宽表(每日凌晨 ETL 生成dws_order_merchant_daily)。


6. 总后台数据聚合实战:用 Kafka + Flink 实现实时订单看板,替代慢 SQL 和定时任务

6.1 为什么总后台不能直连业务库?——数据耦合与性能雪崩的双重陷阱

我见过太多项目,总后台一个“实时销量榜”接口,直接SELECT m.name, COUNT(*) FROM order o JOIN merchant m ON o.merchant_id = m.id WHERE o.create_time > '2024-06-01' GROUP BY m.name ORDER BY COUNT(*) DESC LIMIT 10,结果订单库 CPU 100%,影响 C端下单。根源在于:总后台和业务库共享同一套 MySQL 实例,慢查询拖垮整个 OLTP 链路。解耦唯一路径是:业务库只写,总后台只读;中间用消息队列做数据管道,流计算引擎做实时聚合。

6.2 四端埋点统一 Schema 设计:用 Avro 定义事件结构,避免 JSON 字段歧义

我们定义核心事件 SchemaOrderCreatedEvent.avsc:

{ "type": "record", "name": "OrderCreatedEvent", "namespace": "com.eleme.event", "fields": [ {"name": "event_id", "type": "string"}, {"name": "order_id", "type": "string"}, {"name": "merchant_id", "type": "string"}, {"name": "client_id", "type": "string"}, {"name": "rider_id", "type": ["string", "null"]}, {"name": "total_amount", "type": "double"}, {"name": "create_time", "type": "long"}, // Unix timestamp millis {"name": "city_code", "type": "string"}, {"name": "items_count", "type": "int"} ] }

所有端(client/merchant/order/rider)在订单创建成功后,发送此 Avro 序列化消息到 Kafka Topicorder-created。Avro 比 JSON 体积小 40%,且 Schema Registry 保证前后兼容——新增字段discount_amount不影响旧消费者。

6.3 Flink 实时聚合作业:每分钟统计各城市 TOP10 商家,写入 Elasticsearch

Flink 作业代码(Scala):

val env = StreamExecutionEnvironment.getExecutionEnvironment env.enableCheckpointing(60000) // 60s checkpoint val kafkaSource = KafkaSource.builder[String] .setBootstrapServers("kafka:9092") .setGroupId("flink-order-consumer") .setTopics("order-created") .setValueOnlyDeserializer(new SimpleStringSchema) .build() val orderStream = env.fromSource(kafkaSource, WatermarkStrategy.noWatermarks(), "Kafka Source") .map(json => { val parser = new JsonParser() val root = parser.parse(json).getAsJsonObject OrderCreatedEvent( root.get("event_id").getAsString, root.get("order_id").getAsString, root.get("merchant_id").getAsString, root.get("client_id").getAsString, Option(root.get("rider_id")).map(_.getAsString).orNull, root.get("total_amount").getAsDouble, root.get("create_time").getAsLong, root.get("city_code").getAsString, root.get("items_count").getAsInt ) }) // 按 city_code + merchant_id 滚动窗口聚合(1分钟) val top10Stream = orderStream .keyBy(_.cityCode, _.merchantId) .window(TumblingEventTimeWindows.of(Time.minutes(1))) .aggregate(new OrderAggregateFunction) // 写入 ES val esSink = new ElasticsearchSink.Builder[String]( Collections.singletonList(new HttpHost("es:9200", 9200, "http")), new ElasticsearchSinkFunction[String] { override def process(element: String, context: RuntimeContext, sinkFunction: RequestIndexer): Unit = { val doc = Map( "timestamp" -> System.currentTimeMillis(), "city_code" -> element.split("|")(0), "merchant_id" -> element.split("|")(1), "order_count" -> element.split("|")(2).toInt, "total_amount" -> element.split("|")(3).toDouble ).asJava val indexRequest = Requests.indexRequest() .index("order_top10_daily") .source(doc, XContentType.JSON) sinkFunction.add(indexRequest) } } ) top10Stream.map(t => s"${t._1}|${t._2}|${t._3}|${t._4}") .addSink(esSink)

OrderAggregateFunction内部实现add/getResult/merge,每分钟输出一条city|merchant|count|amount字符串。

6.4 总后台前端对接:用 Elasticsearch DSL 替代 SQL,实现亚秒级响应

总后台查询接口不再走 MySQL,而是调用 ES:

// admin-service/src/main/java/com/eleme/admin/controller/DashboardController.java @GetMapping("/dashboard/top-merchants") public List<TopMerchant> getTopMerchants(@RequestParam String cityCode) { SearchRequest searchRequest = new SearchRequest("order_top10_daily"); SearchSourceBuilder sourceBuilder = new SearchSourceBuilder(); sourceBuilder.query(QueryBuilders.termQuery("city_code", cityCode)); sourceBuilder.sort(SortBuilders.fieldSort("order_count").order(SortOrder.DESC)); sourceBuilder.size(10); searchRequest.source(sourceBuilder); SearchResponse response = restHighLevelClient.search(searchRequest, RequestOptions.DEFAULT); return Arrays.stream(response.getHits().getHits()) .map(hit -> { Map<String, Object> source = hit.getSourceAsMap(); return new TopMerchant( (String) source.get("merchant_id"), ((Number) source.get("order_count")).intValue(), ((Number) source.get("total_amount")).doubleValue() ); }) .collect(Collectors.toList()); }

实测:1000 万订单数据,ES 聚合响应 < 300ms,而 MySQL JOIN 需 25s+。这才是真正“实时”的含义。

我带过的三个 O2O 项目,前两个用慢 SQL 抗了半年,最后都重构为 Kafka+Flink+ES 架构。不是因为技术炫酷,而是当订单量突破 50 万/天,那个“实时销量榜”就不再是锦上添花的功能,而是运营决策的生命线。现在我的习惯是:任何总后台报表需求,第一反应不是写 SQL,而是画 Kafka Topic 和 Flink DAG 图。希望帮到你。

本文还有配套的精品资源,点击获取

返回列表