Skip to content

Java:业务项目进阶 #141

Description

@peng-yin

表结构设计:主子表

上图是实际功能的业务效果图,用户可以新增场景,场景下会包含基本信息、角色信息、剧本流程、评估维度。其中角色信息下方有会包含行为准则信息,剧本流程、评估维度同样会关联一些可以增删的信息。

对上述的功能的表结构设计时,采用了 主子表 的结构设计,子表有进一步包含了子子表。随着开发的进行,发现它是一种常用的模式,并且操作起来还有一些固定的套路。

是什么/为什么

主子表结构(Master-Detail)是数据库设计中最常见的一种模式,它本质上是建立了一种"一对多"(1:N)的关联关系。

在我们的业务场景中,要实现的是子表单独更新,相比较将子表字段用JSON存储,对于分拆主子表结构自然会更灵活和方便。同时由于对于像 陪练角色表,它的数据还会有版本的概念,也即更新之后,training_role中会存储一条新的数据,这样一来,行为准则表 就能关联不同的数据。

怎么用

创建的套路:层级递进保存

  1. 保存主表:保存 training_scenario,获取生成的 scenario_id。
  2. 保存第一层子表
    • 将 scenario_id 赋值给 basic_info、training_role、script_process 等。
    • 关键点:保存 training_role 等列表时,保存后必须立即获取它们各自生成的 ID(如 role_id)。
  3. 保存第二层子表
    • 遍历角色列表,将生成的 role_id 赋值给对应的 behavior_rules。
    • 同理,将 process_id 赋值给 script_process_detail。

查询的套路

对于查询功能,触发入口在列表中,前端传递是时主表 ID,因此查询子表的套路遵循:

  1. 用 主表 ID 批量查出所有 L1 子表(注意此时返回的ID是子表ID)
  2. 用 L1 的 ID 集合(如 role_ids)批量查出所有 L2 子表(behavior_rules)。
  3. 在 Service 中通过 Stream 或循环,根据 ID 将 L2 挂载到 L1 对象中,将 L1 挂载到主对象中。

更新L1子表的套路

如果需要更新:

  • 根据 scenario_id 删除旧的所有角色记录(由于级联关系,通常也会删掉 L2)。
  • 将前端传来的新角色列表重新插入。

如果无需更新:

  • 根据 role_id找到角色记录,并更新

更新L2子表的套路

比如更新training_role表的 behavior_rules,必须传入其归属的 role_id。

原则:不要直接提供一个"更新规则"的独立接口。应该在"更新角色"的业务流程中,一并处理该角色下的规则。

处理策略

  1. 先更新 training_role 的基本字段。
  2. 针对该 role_id 下的 behavior_rules 执行"先全删再全增"。

技巧/注意事项

清楚具体使用的套路,防止模型错误实现

上述列举主子表涉及到的套路,具体的场景不同,可能采用的策略又会不一样,不可直接照搬,亦不可直接在自己没清楚采用什么套路时,直接交模型写。

例如 更新L2子表的套路 中,开发者或模型构想的做法是,behavior_rules中的数据有id就更新,没有就插入新数据。而实际山用户更新 training_role表,传入的数据可能有 新增(无id)、修改(有id)、删除(并不体现在数据 behavior_rules),要能兼容这三种case,直接先全删除数据库关联的 behavior_rules再插入新数据无疑是最简单的做法。


初识异步调用(Asynchronous Call)

是什么/为什么

简单来说,异步调用是一种"发完指令就不管了,直接去干别的事"的通信方式。

1. 生活中的类比:去餐厅吃饭

  • 同步调用(Synchronous):你排队点餐,点完后站在收银台前死等,厨师不把饭做好,你就不走,后面的人也点不了餐。你的时间被完全锁死在了"等饭"这件事上。
  • 异步调用(Asynchronous):你点完餐,服务员给你一个取餐号(或者震动铃)。你直接去找位子坐下玩手机、聊天(干别的事)。等饭做好了,铃声响了,你再去取餐。

2. 技术层面的定义

  • 同步调用:调用方(Java)发出请求后,必须等待被调用方(Python)返回结果,才能继续执行后面的代码。程序会"卡"在这一行。
  • 异步调用:调用方(Java)发出请求后,立即执行下一行代码,不等待结果。被调用方(Python)会在处理完后,通过"回调"或者"消息"的方式通知调用方。

在实际业务开发中,用户直接跟java程序打交道,java程序再调用python程序,python会调用模型进行内容生成(耗时较长),如下结合此案例对为什么要使用异步调用进行了总结:

维度 详细解释 / 原理 在你的项目(Java+Python AI)中的体现 对比:如果用同步会怎样?
用户体验:防止界面卡死 用户的请求立即得到响应,不需要等待后台漫长的计算过程。 用户点击"生成角色",Java 立即返回"已提交",前端开始展示进度条。 页面一直转圈,用户无法操作,甚至以为浏览器或 App 挂了。
系统吞吐量:高并发支持 Java 的线程发完指令就释放,去处理下一个请求,而不是原地等结果。 即使有 100 个人同时生成角色,Java 也能轻松应付,因为它只负责"下订单"。 100 个请求把 Java 线程池占满,第 101 个用户连网页都打不开。
资源利用率:避免线程阻塞 消除"无谓的等待"。CPU 不会在等待 IO 或远程调用返回时闲置。 Java 不用"陪着" Python 跑模型,Java 线程可以去干别的工作(如查库、鉴权)。 Java 线程处于 Blocked 状态,白白消耗内存和 CPU 资源却没干活。
系统稳定性:故障隔离与解耦 调用方(Java)的性能不再受被调用方(Python)运行速度的直接牵制。 如果 Python 运行特别慢或临时宕机,Java 依然能正常运行,不会被拖垮。 Python 变慢会导致 Java 积压大量请求,引发雪崩效应,整个系统一起崩溃。
业务逻辑:处理长耗时任务 绕过网络传输的时间限制(Timeout)。 AI 模型运行可能需要几分钟,远超 HTTP 请求通常的 30 秒超时限制。 还没等模型跑完,浏览器就报"网关超时(504 Gateway Timeout)"错误。

怎么用

进程和线程

如果不理解进程和线程,你很难写出高性能、高可用的企业级应用。

学习它是必不可少的,主要体现在以下几个方面:

  • 优化用户体验:通过异步处理耗时任务(如文件下载),防止主线程阻塞导致界面卡死。
  • 支撑高并发:利用线程池同时处理海量用户请求,大幅提升服务器的吞吐能力。
  • 榨干多核性能:将大数据计算拆分为多任务并行执行,成倍缩短处理时间。
  • 保障数据准确:通过同步锁机制(如 synchronized),规避多线程竞争导致的数据冲突与业务 Bug。
  • 掌握前沿技术:理解底层原理才能驾驭 Java 21 虚拟线程等新技术,应对百万级高并发挑战。

学习建议路线:

  1. 入门:学会如何继承 Thread 类或实现 Runnable 接口。
  2. 进阶:掌握线程池 (ThreadPoolExecutor),这是实际开发中最常用的方式(不要手动 new Thread)。
  3. 核心:理解 Java 内存模型 (JMM) 和线程安全问题(死锁、原子性、可见性),掌握synchronized、Lock、volatile 等工具来保证数据准确。
  4. 架构:学习 CompletableFuture 或响应式编程,处理复杂的异步业务逻辑。

进程 (Process) —— 资源的分配单位

定义:进程是操作系统进行资源分配和调度的一个独立单位。简单来说,一个运行中的软件就是一个进程。

类比:想象一家工厂。工厂有独立的土地、电力、原材料仓库。不同工厂之间是物理隔离的,一家工厂倒闭了通常不会影响另一家。

在Java中:当你运行 java Main 时,操作系统就会启动一个 JVM 进程。

线程 (Thread) —— 执行的最小单位

定义:线程是进程内部的一条执行路径,是 CPU 调度的最小单位。一个进程可以包含多个线程,它们共享进程的资源。

类比:想象工厂里的工人。

  • 一个工厂(进程)里至少有一个工人(主线程)。
  • 工人们共享工厂的设备、原材料(共享内存)。
  • 多个工人同时干活,生产效率更高。

在Java中:Java 程序天生就是多线程的。即使你只写一个 main 方法,JVM 也会启动垃圾回收线程(GC)、信号处理线程等。

实践的场景

实践过程中用到异步调用的场景多是java和python进行交互的场景,如下是一例:

Java端提供了一个生成角色接口[1],接口会调用Python端的生成角色接口[2],[2] 会调用大模型耗时较长,等到执行完毕后,会再调用你Java端的接口[3]将角色信息回写。

如果[1] [2]接口都不是异步的,那么用户在浏览器端将会得到一个超时失败的接口,而通过将接口[2]实现为异步,用户可以很快完成角色的生成。

5种实现方法

基于模型总结了5种异步调用的实现方式,需要使时用于自查:

1. 最基础的做法:手动创建新线程 (Thread)

这是最直观的方法:在接口 A 中手动开启一个新线程来运行接口 B。

技术点:Thread 或 Runnable。

代码示例

public void interfaceA() {
    System.out.println("接口A开始执行");

    // 开启一个新线程来执行接口B
    new Thread(() -> {
        interfaceB(); // 这里的耗时操作不会阻塞主线程
    }).start();

    System.out.println("接口A已响应用户(无需等待B)");
}

缺点:每次请求都创建新线程,如果用户访问量大,系统会因为创建过多线程而崩溃。实际开发中不推荐直接用。

2. 标准的做法:线程池 (ExecutorService)

这是工业级的标准做法。预先准备好一批工人(线程池),有任务了就丢给其中一个工人,干完活工人不销毁,等下一个任务。

技术点:java.util.concurrent.ExecutorService。

代码示例

// 1. 创建一个固定大小的线程池(通常定义为全局静态变量)
private static final ExecutorService executor = Executors.newFixedThreadPool(10);

public void interfaceA() {
    System.out.println("接口A开始");

    // 2. 将耗时任务提交给线程池
    executor.submit(() -> {
        interfaceB();
    });

    System.out.println("接口A返回结果");
}

优点:资源可控,不会因为任务太多压垮服务器。

3. 现代 Java 的方式:CompletableFuture

如果你使用的是 Java 8 及以上版本,这是非常推荐的工具,它处理异步任务非常优雅。

技术点:java.util.concurrent.CompletableFuture。

代码示例

public void interfaceA() {
    // 异步执行接口B
    CompletableFuture.runAsync(() -> interfaceB());

    System.out.println("接口A立即返回");
}

优点:语法简洁,支持非常复杂的逻辑编排(比如 B 执行完后再自动执行 C)。

4. 最省心的方式:Spring 框架的 @async 注解

如果你正在学习或使用 Spring Boot(这是目前 Java 开发的主流框架),实现异步只需要一个注解。

技术点@async

实现步骤

  1. 在启动类上加 @EnableAsync
  2. 在接口 B 的方法上加上 @Async

代码示例

// 在 Service 类中
@Async
public void interfaceB() {
    // 模拟耗时操作
    Thread.sleep(5000);
    System.out.println("接口B终于干完活了");
}

// 在 Controller 中调用
public void interfaceA() {
    service.interfaceB(); // 它会自动进入异步线程运行
    return "成功";
}

优点:代码侵入性极低,阅读最清晰。

5. 进阶方案:消息队列 (MQ)

如果接口 B 极其耗时(比如要跑几个小时),或者非常重要(绝对不能丢失),我们会用到中间件。

技术点:RabbitMQ, Kafka, RocketMQ。

逻辑:接口 A 把任务发到一个"消息队列"里就完事了。另外一个独立的程序(消费者)从队列里取出任务慢慢执行。

技巧/注意事项

在自己的开发场景中,如果按照模型的建议,在Java和Python中要采取双端异步的做法。但在Java侧似乎并无必要,只需要将Python的方法改为异步支持即可。因此本节内容称之为初探,只是理解了它的概念和用途,所推荐的各种实现异步的方式并无实践。


HSF(High-Speed Service Framework)

是什么/为什么

HSF(High-Speed Service Framework) 是阿里巴巴自研的高性能、分布式 服务框架,主要用于支撑大规模微服务架构下的远程过程调用(RPC)。HSF 本身是闭源的,但阿里巴巴后来将 HSF 的核心思想和能力 开源并演进为 Dubbo, HSF ≈ Dubbo(尤其是 Dubbo 2.x / 3.x)。

远程过程调用(RPC,Remote Procedure Call)是一种让程序像调用本地函数一样调用远程服务器上的服务或方法的通信机制。它是构建分布式系统、微服务架构的核心技术之一。

完整的 HSF 使用应该包含 Provider 和 Consumer,如下是一个典型的 HSF 项目目录结构:

user-service-parent/                     # 父工程(pom 类型)
│
├── user-service-api/                    # ✅ 核心:API 模块(只包含接口和 DTO)
│   ├── src/main/java/
│   │   └── com/alibaba/sample/api/
│   │       ├── UserService.java         # 服务接口
│   │       └── User.java                # 数据传输对象(DTO),必须 Serializable
│   └── pom.xml
│
├── user-service-provider/               # ✅ 服务提供方(Provider)
│   ├── src/main/java/
│   │   └── com/alibaba/sample/provider/
│   │       ├── UserServiceImpl.java     # 接口实现
│   │       └── UserApplication.java     # 启动类(Pandora Boot)
│   ├── src/main/resources/
│   │   ├── application.properties       # 应用配置
│   │   └── hsf.properties               # HSF 专用配置(可选)
│   └── pom.xml
│
├── user-service-consumer/               # ✅ 服务消费方(Consumer,可选,常集成在其他业务系统)
│   ├── src/main/java/
│   │   └── com/alibaba/sample/consumer/
│   │       ├── OrderService.java        # 调用 UserService 的业务逻辑
│   │       └── OrderApplication.java
│   └── pom.xml
│
└── pom.xml                              # 父 POM,管理子模块依赖版本

怎么用(定义)

  1. 在api项目的api目录下,定义接口和dto
  2. 在api项目的pom.xml中定义这个服务
  3. 在service项目的provider中定义具体的实现
  4. 将服务注册到服务中心

怎么用(消费)

项目开发中,使用到另外一个hsf服务提供的方法:根据redisKey获取导入的userList。

实际项为例,要在项目中使用另外一个项目提供的hsf服务,如下是具体的代码:

如下在pom.xml中引用使用的hsf服务:

<dependency>
    <groupId>me.ele.uni</groupId>
    <!--  API 模块 通常只包含 接口定义、DTO 数据类、异常类等不包含具体实现逻辑  -->
    <artifactId>learning_manager_api</artifactId>
    <!-- 开发中的快照版本(-SNAPSHOT),说明该 API 仍在迭代,适用于开发/测试环境;生产环境通常用正式版(如 2.1.2)。-->
    <version>2.1.2-SNAPSHOT</version>
    <!-- 显式排除 learning_manager_api 所依赖的两个库。-->
    <exclusions>
        <exclusion>
            <artifactId>uni_common</artifactId>
            <groupId>me.ele.uni</groupId>
        </exclusion>
        <exclusion>
            <artifactId>servlet-api</artifactId>
            <groupId>javax.servlet</groupId>
        </exclusion>
    </exclusions>
</dependency>

注册要使用的service:

package me.ele.uni.hsf.consumer;

import com.alibaba.boot.hsf.annotation.HSFConsumer;
import me.ele.uni.learning_manager_api.api.CommonHsfService;
import me.ele.uni.annotation.FastHSFConsumer;

@Configuration
public class HsfConfig {
    @HSFConsumer
    @FastHSFConsumer
    CommonHsfService commonHsfService;
}

使用hsf服务:

package me.ele.uni.service.impl.kk.admin;

import me.ele.uni.learning_manager_api.api.CommonHsfService;

@Service
@Slf4j
public class AiPushServiceImpl implements AiPushService {

    @Lazy
    @Autowired
    private CommonHsfService commonHsfService;

    @Override
    @Transactional(rollbackFor = Exception.class)
    public PushAddRes addPush(PushAddVo pushAddVo, LoginUserDto loginUserDto) {
      List<String> importUserList = commonHsfService.getUserListByImportRedisKey(pushAddVo.getImportUserKey());
    }
}

技巧/注意事项

使用最新版版本的hsf服务

就像前端的npm一来一样,引用的如果是老的版本的服务,是无法看到新版本上的API的。

实际开发中,自己因为引用了老版本的hsf服务,发现该服务上缺少功能,实际开发发现功能是有的。


MyBatis

是什么/为什么

MyBatis 是一款优秀的持久层框架,它支持自定义 SQL、存储过程以及高级映射。

MyBatis 避免了几乎所有的 JDBC 代码和手动设置参数以及获取结果集的工作。

怎么用

在第一篇文章 前端老鸟的Java新手期:AI Coding 踩坑记录与收获 中,我们使用的是MyBatis-Plus,数据持久化的代码要简单很多。

而在业务项目中,很多项目是直接使用的MyBatis,它涉及了 Entity.java、Mapper.java、Mapper.xml、EntityExample.java 组件:

对应角色 协作流向
Mapper.java 接口 (Interface) 1. 发起指令:Service 调用 Mapper 接口方法
EntityExample.java 条件 (Criteria) 2. 封装条件:作为参数传给 Mapper 方法
Mapper.xml 实现 (SQL/XML) 3. 执行 SQL:解析 Example 生成 SQL,查询数据库
Entity.java 数据 (POJO) 4. 返回结果:MyBatis 将结果封装成 Entity 返回

假设你有一张表 User,对应的类是 User.java,辅助类是 UserExample.java。

第一步:在 Java 代码中定义数据操作

你想查询 id < 100 的用户:

UserExample example = new UserExample();
example.createCriteria().andIdLessThan(100L); // 使用 Example 构建条件

第二步:调用 Mapper 接口

// 调用接口方法,传入 example 对象
List<User> list = userMapper.selectByExample(example);

第三步:MyBatis 解析 Mapper.xml

MyBatis 找到 XML 中的这段代码:

<select id="selectByExample" parameterType="UserExample" resultMap="BaseResultMap">
  select * from user
  <if test="_parameter != null">
    <include refid="Example_Where_Clause" /> <!-- 关键:这里把 example 对象转成了 WHERE id < 100 -->
  </if>
</select>

第四步:封装返回

MyBatis 执行 SQL 后,根据 XML 里的 ResultMap 定义,自动把数据库每一行记录 new 成一个 User (Entity) 对象,最后返回一个 List<User> 给调用者。

技巧/注意事项

使用generate自动生成有关文件

<?xml version="1.0" encoding="UTF-8"?>
<!DOCTYPE generatorConfiguration
        PUBLIC "-//mybatis.org//DTD MyBatis Generator Configuration 1.0//EN"
        "http://mybatis.org/dtd/mybatis-generator-config_1_0.dtd">

<generatorConfiguration>
    <classPathEntry
            location="/Users/kougyu/Documents/mysql-connector-java-5.1.47.jar"/>

    <context id="Mysql" targetRuntime="MyBatis3" defaultModelType="flat">
        <!-- 自动识别数据库关键字,默认false,如果设置为true,根据SqlReservedWords中定义的关键字列表;
           一般保留默认值,遇到数据库关键字(Java关键字),使用columnOverride覆盖
        -->
        <property name="autoDelimitKeywords" value="false"/>
        <!-- 格式化java代码 -->
        <property name="javaFormatter" value="org.mybatis.generator.api.dom.DefaultJavaFormatter"/>
        <!-- 格式化XML代码 -->
        <property name="xmlFormatter" value="org.mybatis.generator.api.dom.DefaultXmlFormatter"/>
        <commentGenerator>
            <property name="suppressDate" value="true"/>
            <property name="suppressAllComments" value="true"/>
        </commentGenerator>
        <jdbcConnection driverClass="com.mysql.jdbc.Driver"
                        connectionURL="jdbc:mysql://rm-hp3k73805459c2aj8.mysql.huhehaote.rds.aliyuncs.com/elearprod?useSSL=false&amp;allowMultiQueries=true" userId="eleardev"
                        password="dldewredkR7wk"/>
        <!-- 生成的entity目录,默认模块.base,base表示基础的,迁移之前的都是ext扩展的,后面全部过度到base,扩展的继承base,base尽量做到重用 -->
        <javaModelGenerator targetPackage="me.ele.uni.entity"
                            targetProject="./src/main/java">
            <property name="enableSubPackages" value="false"/>
            <property name="constructorBased" value="true"/>
            <property name="trimStrings" value="true"/>
        </javaModelGenerator>
        <!-- 生成的mapper xml目录,默认模块.base,base表示基础的,迁移之前的都是ext扩展的,base尽量做到重用 -->
        <sqlMapGenerator targetPackage="mapper"
                         targetProject="./src/main/resources">
            <property name="enableSubPackages" value="false"/>
        </sqlMapGenerator>
        <!-- 生成的mapper目录,默认模块.base,base表示基础的,迁移之前的都是ext扩展的,base尽量做到重用 -->
        <javaClientGenerator targetPackage="me.ele.uni.mapper"
                             targetProject="./src/main/java" type="XMLMAPPER">
            <property name="enableSubPackages" value="false"/>
        </javaClientGenerator>

        <table tableName="TOM_QUESTIONNAIRE_PAPER" domainObjectName="TomQuestionnairePaper">
            <generatedKey column="QUESTIONNAIRE_PAPER_ID" sqlStatement="JDBC"  />
        </table>
    </context>
</generatorConfiguration>

通过上述配置中指定数据库表,以及生成对应文件的位置,通过mybatis-generator:generate 就可以一次性生成模板代码。不过要注意的是,生成的Mapper.xml文件中,是增量式生成,会有重复的内容,需要手动进行删除,否则项目无法正常部署。


SSE流式输出

是什么/为什么

SSE(Server-Sent Events,服务器发送事件) 是一种让服务器能够实时、连续地向浏览器推送数据的技术。

使用模型生成内容时,生成一个词就传回一个词,你立刻就能看到内容,感知到的延迟大大降低。

基于 HTTP 协议:它不需要像 WebSocket 那样建立复杂的协议转换,直接走普通的 HTTP 端口。

SSE 是一种基于文本的协议,每条消息由多个字段组成,字段之间用换行符 \n 分隔,常见字段有data、event、id、retry。

单条消息的结束,消息与消息之间用两个换行符 \n\n 分隔;

整个会话/流的结束(业务级约定)包含两种方案:

第一种方案:OpenAI 等大模型公司采用的标准,发送一个约定好的特殊字符串

data: {"content": "再见"}

data: [DONE]

第二种方案:业务 JSON 对象里增加一个字段,标记是否结束。

data: {"content": "你好", "is_finished": false}

data: {"content": "啊", "is_finished": true}

怎么用

实际开发中的需求是,实现接口,调用小e工作流提供的流式接口,然后转发获取到的内容给前端。

在团队中找对此比较熟悉的同学咨询,发现对方为了让流式输出能工作,除了java代码外,还修改了包含nginx配置等多处。

模型总结在Java中可用的方案有4种类,团队采用了方案1,网关层同学推荐采用方案2,自己最终借助大模型使用方案4完成的功能。

方案 使用场景 优点 缺点
方案一:HttpURLConnection(原生) 小型项目、快速测试、无框架依赖 JDK 自带,无需引入依赖;控制逻辑简单 代码较繁琐,需要手动解析 SSE 格式;阻塞 I/O,不适合高并发;无自动重连等高级功能
方案二:Spring WebFlux(响应式) 高并发、需要非阻塞流处理、微服务架构 基于 Reactor,支持响应式和背压;与 Spring 生态集成好;自动将 Flux 转成 SSE 输出 学习成本相对高;需要引入 Spring WebFlux 依赖;对传统 Servlet 思维需要适应
方案三:JAX-RS SSE API(Jersey/Java EE) Java EE/Jakarta EE 项目、已有 JAX-RS 架构 标准 API,跨实现通用;SseEventSink 和 SseBroadcaster 提供简单推送机制;与 REST 接口直接结合 依赖 Java EE 运行环境/相关实现;非响应式,吞吐量受限;部分服务器实现的 SSE 支持有限
方案四:OkHttp(第三方网络库) 需要简化 SSE 读取逻辑、已有 OkHttp 项目 API 简洁,易用;社区资源丰富;连接管理好,支持超时和自动重连配置 需额外引入 OkHttp 依赖;输出端仍需自己处理 SSE 格式;属于阻塞 I/O

方案1 原生方案

默认地我们让模型实现一个接口,转发来着小E的流流式输出,它采用的是HTTPRequest类(来自 cn.hutool.http.HttpRequest),最后发现读取到的内容在头传给浏览器端存在消息挤压的问题。也即虽然小e的流接收到后,在java端打印日志是 10:01 10:02 ... 10:40 陆续抵达,但是在浏览器端,他们会是在 10.40 一次性到达。

团队在运行的项目都是采用 HttpURLConnection 方案,本质上 HTTPRequest 应该是在 HttpURLConnection 上进行了更上层的封装,不排除有参数设置导致不可用的原因。

方案4 OkHttp方案

最终让模型重构,它使用了方案4,并且得到了功能良好的版本。核心的services层代码如下:

/**
     * 生成上下文
     *
     * @param generateContextReqDTO
     * @return
     */
@Override
public SseEmitter generateContext(GenerateContextReqDTO generateContextReqDTO) {
    UserSubject userSubject = Subject.get();

    // 设置超时时间为5分钟
    SseEmitter emitter = new SseEmitter(300000L);

    // 设置完成和超时回调
    emitter.onCompletion(() -> log.info("SSE连接正常完成"));
    emitter.onTimeout(() -> log.warn("SSE连接超时"));
    emitter.onError((ex) -> log.error("SSE连接发生错误", ex));

    // 使用线程池异步处理SSE流
    frontendMessagePoolExecutor.execute(() -> {
        java.io.InputStream inputStream = null;
        java.io.BufferedReader reader = null;

        try {
            // 构建请求参数
            Map<String, Object> requestBody = new HashMap<>();

            // 构建inputs对象
            Map<String, String> inputs = new HashMap<>();

            requestBody.put("inputs", inputs);
            requestBody.put("response_mode", "streaming");
            requestBody.put("user", userSubject.getWorkNo());

            String requestBodyJson = JSONObject.toJSONString(requestBody);
            log.info("生成上下文请求参数:[{}]", requestBodyJson);

            // 调用远程SSE接口,使用流式读取
            String url = "http://ialpha.rajax-inc.com/v1/workflows/run";

            // 使用 OkHttp 实现真正的流式读取
            okhttp3.OkHttpClient client = new okhttp3.OkHttpClient.Builder()
                .connectTimeout(300, java.util.concurrent.TimeUnit.SECONDS)
                .readTimeout(300, java.util.concurrent.TimeUnit.SECONDS)
                .writeTimeout(300, java.util.concurrent.TimeUnit.SECONDS)
                .build();

            okhttp3.RequestBody okHttpRequestBody = okhttp3.RequestBody.create(
                okhttp3.MediaType.parse("application/json; charset=utf-8"),
                requestBodyJson
            );

            okhttp3.Request request = new okhttp3.Request.Builder()
                .url(url)
                .post(okHttpRequestBody)
                .addHeader("Authorization", "Bearer app-ASTTcYqEQRF2785WvYQ9MzuA")
                .addHeader("Content-Type", "application/json")
                .build();

            okhttp3.Response response = client.newCall(request).execute();

            if (response.isSuccessful() && response.body() != null) {
                inputStream = response.body().byteStream();
                reader = new java.io.BufferedReader(new java.io.InputStreamReader(inputStream, "UTF-8"));

                String line;
                while ((line = reader.readLine()) != null) {
                    if (StringUtils.isNotBlank(line)) {
                        try {
                            // 解析 SSE 格式数据,去掉 "data: " 前缀
                            String actualData = line;
                            if (line.startsWith("data: ")) {
                                actualData = line.substring(6); // 去掉 "data: " 前缀(6个字符)
                            } else if (line.startsWith("data:")) {
                                actualData = line.substring(5); // 去掉 "data:" 前缀(5个字符)
                            }

                            // 发送SSE数据(SseEmitter 会自动添加 "data: " 前缀)
                            emitter.send(SseEmitter.event().data(actualData));
                            System.out.println(new java.text.SimpleDateFormat("yyyy-MM-dd HH:mm:ss.SSS").format(new java.util.Date()) + " - 原始: " + line + " | 发送: " + actualData);

                        } catch (java.io.IOException e) {
                            // 客户端断开连接,记录日志并退出循环
                            log.warn("客户端已断开SSE连接: {}", e.getMessage());
                            break;
                        }
                    }
                }

                emitter.complete();
                log.info("生成上下文SSE流处理完成");
            } else {
                String errorBody = response.body() != null ? response.body().string() : "无响应体";
                log.error("生成上下文失败,状态码:[{}],响应:[{}]", response.code(), errorBody);
                emitter.completeWithError(new BizException("-1", "生成上下文失败"));
            }

            // 关闭 response
            if (response != null) {
                response.close();
            }

        } catch (Exception ex) {
            log.error("生成上下文异常", ex);
            try {
                emitter.completeWithError(ex);
            } catch (Exception e) {
                log.error("发送SSE错误失败", e);
            }
        } finally {
            // 关闭资源
            if (reader != null) {
                try {
                    reader.close();
                } catch (Exception e) {
                    log.error("关闭reader失败", e);
                }
            }
            if (inputStream != null) {
                try {
                    inputStream.close();
                } catch (Exception e) {
                    log.error("关闭inputStream失败", e);
                }
            }
        }
    });

    return emitter;
}

技巧/注意事项

转发SSE内容需要避免data结构重复

上述的场景中,要在 Java 中实现"接收 -> 裁剪 -> 转发"过程,要检查是否以 data: 开头,如果是,提取后面的 JSON 字符串。足厚将裁剪后的数据重新拼装成 data: {新内容}\n\n 的格式。


面向AOP的的异常处理

是什么/为什么

异常可以避免程序因意外崩溃(容错),也可以通过堆栈信息快速定位和修复bug(溯源),并且将正常业务逻辑和错误处理逻辑分离(解耦)。关于异常,任何一门程序中都有,这里引用 Java异常处理和最佳实践 中的一个检查异常和非检查异常,检查异常是开发者需要关注的异常。

检查异常:也称为"编译时异常",编译器在编译期间检查的那些异常。由于编译器"检查"这些异常以确保它们得到处理,因此称为"检查异常"。如果抛出检查异常,那么编译器会报错,需要开发人员手动处理该异常,要么捕获,要么重新抛出。除了RuntimeException之外,所有直接继承 Exception 的异常都是检查异常。

非检查异常:也称为"运行时异常",编译器不会检查运行时异常,在抛出运行时异常时编译器不会报错,当运行程序的时候才可能抛出该异常。Error及其子类和RuntimeException 及其子类都是非检查异常。

基于切面的异常处理(AOP Exception Handling) 是指将异常处理代码从业务逻辑中抽离出来,集中在一个地方进行管理。它的核心思想是:不再在每个方法里写 try-catch,而是通过"切面"在程序运行出错时自动拦截异常。

怎么用

异常在任何开发语言中都可以单开一篇进行介绍,本节单独介绍下面向切面的异常处理,理解它有助于知道异常是如何在发声位置被传递到客户端的。

有两种方式实现面向切面的异常处理,这里介绍Web 层专用的 @ControllerAdvice 方式。这是 Spring MVC 提供的一种"声明式"切面,专门用于拦截 Controller 抛出的异常。它本质上是 AOP 的一种高度封装。包含:

  1. 定义统一响应结构(比如你之前的 code, errorMessage)。
  2. 创建全局异常处理器:
@RestControllerAdvice // 这是一个特殊的切面类
public class GlobalExceptionHandler {

    // 拦截所有的业务异常
    @ExceptionHandler(BusinessException.class)
    public Result handleBusinessException(BusinessException e) {
        // 返回统一的错误格式
        return Result.fail(e.getCode(), e.getMessage());
    }

    // 拦截未知的系统异常(如空指针、数据库宕机)
    @ExceptionHandler(Exception.class)
    public Result handleException(Exception e) {
        log.error("系统崩溃了:", e);
        return Result.fail(500, "服务器开小差了,请稍后再试");
    }
}

当项目中发生异常时,都会经过切面类的统一处理,并且按照"就近原则",或者叫 "最匹配原则(Most Specific Match)",Spring 会计算抛出的异常与处理方法中定义的异常之间的"距离",并选择距离最短的那一个。然后按照方法的处理完后返回给前端。

设计一个高质量的异常切面类(如 @ControllerAdvice)是构建稳健系统的关键。建议遵循以下 6 个核心原则:

1. 统一响应结构原则 (Unified Response Structure)

原则描述:无论发生什么异常,返回给前端或下游系统的 JSON 格式必须是固定的。

为什么:方便前端或调用方编写统一的拦截器来处理错误。如果有时返回 JSON,有时返回 HTML 错误页,调用方会崩溃。

做法:定义一个通用的 Result<T> 或 Response 类,包含 code、message、data 等字段。

2. 异常分层处理原则 (Layered Exception Handling)

原则描述:区分 业务异常 (Business Exception) 和 系统异常 (System Exception)。

  • 业务异常(如:余额不足、密码错误):属于预期内的错误。处理:无需记录堆栈详情(Error Log),直接返回错误码和友好提示。
  • 系统异常(如:数据库连接失败、空指针):属于非预期错误。处理:必须记录详细的堆栈信息(log.error),方便排查排,但给用户返回通用的"系统繁忙"。

3. 信息安全原则 (Security & Privacy)

原则描述:严禁将底层的堆栈信息(Stack Trace)或敏感 SQL 报错直接返回给客户端。

为什么:泄露代码路径、类名、数据库表名等信息会给黑客提供攻击线索(SQL 注入、漏洞利用)。

做法

  • 在开发环境(dev)可以显示详细信息。
  • 在生产环境(prod),对非业务异常一律通过"兜底方法"拦截,转换成模糊的提示语(如:"服务器开小差了")。

4. 记录有意义的日志原则 (Meaningful Logging)

原则描述:异常切面不仅是为了给用户反馈,更是为了给开发者排查问题。

做法

  • 上下文记录:在 Log 中记录发生异常时的请求 URL、当前用户 ID、关键入参等。
  • 分级记录:业务异常记 WARN 或 INFO,系统异常记 ERROR。
  • 静默处理:绝不要在切面里用 e.printStackTrace(),必须使用成熟的日志框架(Logback/Log4j2)。

5. 兜底与降级原则 (Fail-Safe & Fallback)

原则描述:必须有一个"万能补丁"方法拦截 Exception.class

为什么:程序员无法预知所有的异常。如果没有拦截最顶层的 Exception,当出现未预料到的错误时,框架会默认抛出 500 错误,甚至可能返回服务器容器(Tomcat/Nginx)的默认页面。

做法:始终在切面类的最后放一个处理 Exception.class 的方法。

6. 不要"吃掉"异常原则 (No Silent Catching)

原则描述:异常切面的职责是"转化"异常为"响应",而不是平白无故地消除它。

为什么:如果你在切面里捕获了异常却只返回了一个 success: true,上游系统会误以为操作成功,导致数据一致性问题。

做法:确保切面处理后,HTTP 状态码(或业务状态码)能真实反映错误的本质。

技巧/注意事项

按照切面类处理异常返回

在实际开发中,有这么一个/generate-context接口,它会调用流式接口,然后同样按照流式进行内容输出。在实际联调中,偶现了调用出现了 406 not acceptable 的异常,无其他可用的日志信息。本能的对 developTaskService.generateContext 进行排查无果;最后定位在问题发生在 @Valid @RequestBody GenerateContextReqDTO 上。

@Slf4j
@Api(tags = "开发任务api", value = "开发任务相关接口")
@RestController
@RequestMapping("/api/develop-tasks")
public class DevelopTaskController {

    @ApiOperation(value = "生成上下文", extensions = {@Extension(properties = {@ExtensionProperty(name = "author", value = "唐叶飞")})})
    @PostMapping("/generate-context")
    public SseEmitter generateContext(@Valid @RequestBody GenerateContextReqDTO generateContextReqDTO) {
        return developTaskService.generateContext(generateContextReqDTO);
    }

}

在GenerateContextReqDTO中定义了一个字段,进行了@notblank的检查,当用户没有输入合法的参数,触发了异常。

@ApiModelProperty("API Schema类型")
@NotBlank(message = "apiSchemaType不能为空")
private String apiSchemaType;

异常最终会进入到 DefaultGlobalExceptionHandler.javamethodArgumentNotValidException 处理,虽然获取到了异常信息,也构建了返回对象,但是 RestfulResult的结构跟 SseEmitter generateContext 方法是不兼容的,因此导致了 406 not acceptable 的异常。处理办法可以是:改为手动对参数进行校验,然后通过在 SseEmitter的输出中提供错误信息。

    /**
     * 参数异常拦截
     *
     * @param methodArgumentNotValidException
     * @return
     */
    @ExceptionHandler(MethodArgumentNotValidException.class)
    public RestfulResult methodArgumentNotValidException(MethodArgumentNotValidException methodArgumentNotValidException) {
        HttpServletRequest httpServletRequest = ServletRequestUtil.getHttpServletRequest();
        String remoteHost = httpServletRequest.getRemoteHost();
        String url = httpServletRequest.getRequestURL().toString();

        log.error("客户端host地址:[{}]访问url地址:[{}]和统一拦截参数异常msg:[{}]", remoteHost, url, methodArgumentNotValidException.getMessage());
        String msg = "参数不合法";
        BindingResult bindingResult = methodArgumentNotValidException.getBindingResult();
        if (bindingResult.hasErrors()) {
            msg = handleError(bindingResult);
        }

        return fail(msg);
    }

再说枚举

是什么/为什么

在上一篇枚举基本使用中,我们认定直接使用Integer是一种偷懒的做法,在实际项目中,发现这确实普遍存在的。不可避免地,follow现有项目的做法,却遇到了一些问题。

怎么用

现在定义了如下的枚举类,在Request VO中,使用 PartnerRoleEnum partnerRole 来接受参数,结果导致插入数据库的值,跟用户实际传递的不一样,例如:前端选择的是媒体类型,传递1,到数据库存入的却是3。为什么?

public enum PartnerRoleEnum {

  MEDIA(1, "媒体"),

  LAWYER(2, "律师"),

  POLICE(3, "公安局(警方)"),

}

这是因为 Jackson反序列化枚举的默认行为 导致的。

Jackson 是 Java 中最流行的 JSON 解析库:

  • 序列化 (Serialization):把 Java 对象 转换成 JSON 字符串。
  • 反序列化 (Deserialization):把 JSON 字符串 转换成 Java 对象。

当前端发送的是整数值(如 1),而你的VO中定义的是枚举类型时,Jackson默认会尝试匹配枚举的 name(枚举常量名),而不是你定义的 code 值。在项目中如果直接传递 {"partnerRole": "MEDIA"} 是能成功反序列化的。

public class PartnerRequestVO {
    private PartnerRoleEnum partnerRole;  // 前端发送: 1,但jackson无法反序列化
}
// 前端请求:{"partnerRole": 1}
// Jackson 尝试找名为"1"的枚举常量 → 失败,使用了一个错误的值插入书库

技巧/注意事项

分层使用,不全部用枚举

如下概括了再各个层集中如何使用枚举,并提供了 Entity 转 VO 和 VO 转 Entity 的示例代码。

层级 字段类型 优点 缺点
Entity Integer ✅ 与DB保持一致 ✅ ORM映射简洁 ❌ 业务逻辑需转换
RequestVO Integer ✅ 兼容性好 ✅ 客户端简单 ❌ 需验证
ResponseVO Integer + String ✅ 减少枚举依赖 ✅ 前端易用 ❌ 需要两个字段
业务代码 RoleEnum ✅ 类型安全 ✅ IDE提示 ❌ 需转换

最简单做法是 Entity 和 VO 中都用 Integer:

对于 Entity -> responseVO,核心是使用枚举获取具体的desc值:

// 5. Service层(业务逻辑)- 枚举和Integer互转
@Service
public class UserService {

    // Entity -> responseVO
    public UserResponseVO convertToVO(UserEntity entity) {
        UserResponseVO vo = new UserResponseVO();
        vo.setId(entity.getId());
        vo.setName(entity.getName());

        // Integer 转 枚举,获取信息
        RoleEnum roleEnum = RoleEnum.getByCode(entity.getRole());
        vo.setRoleCode(entity.getRole());
        vo.setRoleDesc(roleEnum != null ? roleEnum.getDesc() : "未知");

        return vo;
    }

    // requestVO -> Entity
    public UserEntity convertToEntity(UserRequestVO vo) {
        UserEntity entity = new UserEntity();
        entity.setName(vo.getName());

        // 验证角色有效性
        RoleEnum roleEnum = RoleEnum.getByCode(vo.getRole());
        if (roleEnum == null) {
            throw new IllegalArgumentException("无效的角色码:" + vo.getRole());
        }

        entity.setRole(vo.getRole());
        return entity;
    }
}

其他技巧

业务逻辑为什么要写在service中,而不是controller中?

  • 复用性:逻辑写在 Service 里,不仅 Controller(网页请求)能用,定时任务、消息队列、命令行工具也都能直接调。如果写在 Controller 里,逻辑就被 HTTP 接口绑死了。
  • 解耦:Controller 是"前台":只负责接客(接收请求、校验参数、返回响应)。Service 是"后厨":只负责做菜(算账、改数据库、走业务流程)。前台不该进厨房炒菜,后厨也不该去大堂收银。
  • 事务管理:数据库的"要么全成功,要么全失败"(事务)通常加在 Service 上。如果逻辑散落在 Controller 里,很难保证复杂操作的原子性。
  • 易测试性:测试 Service 只需要写简单的 Java 代码;测试 Controller 却要模拟网络环境和 HTTP 请求,非常麻烦。

一句话:Controller 应该"薄",Service 应该"厚"。

用上了Thread.sleep但是不是好的做法

实际开发中遇到了这么么一个问题:开发了 接口a 用于提交信息,并且会触发外部接口c;开发了 接口b 用于同步信息,外部接口c会调用接口b。在实际调用过程中,由于a接口,调用外部接口c后,数据库还没有来得及落库,c接口就开始调用接口b回写,导致抛出异常,大致流程如下:

这里后端开发常用的梗 sleep 一下 就派上用场,如下代码进行5 * 1s的尝试,直到数据库有数据,才开始执行后续逻辑。

  // 根据主键 id 查询角色,如果不存在则重试,最多重试5次,每次间隔1秒
  TrainingRole existingRole = null;
  int maxRetries = 5;
  for (int i = 0; i < maxRetries; i++) {
      existingRole = trainingRoleMapper.selectByPrimaryKey(roleId);
      if (existingRole != null) {
          break;
      }
      if (i < maxRetries - 1) {
          try {
              log.info("角色不存在,等待1秒后重试,当前重试次数: {}", i + 1);
              Thread.sleep(1000);
          } catch (InterruptedException e) {
              Thread.currentThread().interrupt();
              throw new RuntimeException("等待重试时被中断", e);
          }
      }
  }

在接口 B 中使用 sleep 是一种"凭运气编程",在生产环境下存在严重隐患:

  • 不可靠性:你无法确定 sleep 多久才是安全的。如果数据库负载高、网络延迟或事务执行变慢,sleep 时间可能不够,导致依然查不到数据。
  • 性能瓶颈:sleep 会阻塞线程。在高并发场景下,大量回调请求会迅速耗尽 Web 服务器的线程池,导致系统雪崩。
  • 资源浪费:线程被占用却不干活,极大地降低了系统的吞吐量。
  • 逻辑缺陷:如果接口 A 的事务最终因为后续逻辑报错而回滚了,接口 B 无论 sleep 多久都查不到数据,但它依然白白等待了。

如下更好的做法,其中最推荐的是事务同步器的方式:

方案名称 实现原理 优点 缺点 适用场景 推荐等级
事务同步器 (afterCommit) 在接口 A 中利用框架工具(如 Spring 的 TransactionSynchronizationManager),确保只有在事务成功提交后,才发起对外部服务 S 的调用。 最优雅、最彻底。从根源上保证了 B 回调时数据一定已落库。 需要依赖特定的框架特性(如 Spring)。 绝大多数业务场景。 ⭐⭐⭐⭐⭐ (首选)
拆分事务边界 将接口 A 拆分为两个方法:方法 1 只负责写数据库并提交事务;方法 2 调用方法 1 之后再调用外部服务 S。 逻辑清晰,不依赖复杂的框架技巧,代码可读性好。 需要重构代码结构,确保外层调用方没有开启大事务。 逻辑较简单的业务模块。 ⭐⭐⭐⭐
异步消息队列 (MQ) 接口 A 提交事务后发送消息到 MQ,由消费者去执行外部调用;或者接口 B 收到回调后发消息异步处理。 高吞吐、解耦。支持失败重试,能应对高并发和网络波动。 增加了系统复杂度(引入 MQ 组件),需考虑消息丢失或重复消费问题。 高并发、对响应时间敏感、链路较长的系统。 ⭐⭐⭐⭐
智能重试机制 (Retry) 接口 B 查不到数据时,不直接报错,而是采用退避算法(如每隔 100ms 重试,共 5 次)进行查询。 不需要改动接口 A 的逻辑,对外部系统透明。 依然占用线程资源(虽然比 sleep 效率高),且若事务回滚,重试仍会失败。 无法改动接口 A 或外部服务 S 的存量系统补救。 ⭐⭐⭐
冗余字段回传 接口 A 调用 S 时带上核心数据,要求 S 在回调 B 时原样传回,使 B 无需查询数据库。 完全消除时序问题,B 不再依赖数据库查询即可处理逻辑。 需要外部服务 S 配合修改接口定义;如果数据量大,传参不便。 外部服务可控且传递的数据量较小的场景。 ⭐⭐⭐
使用 Redis 暂存 接口 A 在写库的同时,将数据写入 Redis。接口 B 优先从 Redis 获取。 查询速度极快,避开了数据库事务隔离级别的限制。 增加了 Redis 维护成本,需处理 Redis 与数据库的数据一致性。 对实时性要求极高、并发压力大的场景。 ⭐⭐

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Projects

    No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions