一.用户模块

user 模块是项目中的基础用户领域模块,基于 MyBatis 实现,包含 domain.User 领域对象、UserMapper 数据访问接口及其 XML SQL 映射、UserService/UserServiceImpl 基础服务封装。它本身不直接提供 HTTP 接口,而是为 auth、profile、relation 等上层模块提供用户数据查询、创建、密码更新等底层能力。

其中,UserService 当前封装了按手机号/邮箱/用户 ID 查询用户、判断手机号/邮箱是否存在、创建用户、更新密码等方法;UserMapper 则声明了更底层的数据访问能力,实际 SQL 主要定义在 UserMapper.xml 中。需要注意的是,目前项目并未把所有用户操作都统一收口到 UserService,部分模块仍直接依赖 UserMapper。

另外,@Transactional(readOnly = true) 的作用主要是声明只读事务和提供优化提示,并不直接解决并发问题;项目中真正用于保证手机号、邮箱、知光号唯一性的核心约束是数据库唯一索引。

二.认证登录注册模块

1. 认证配置

  • AuthProperties类 使用 @ConfigurationProperties(prefix = "auth"),负责把 application.yml 中的 auth.jwt、auth.verification、auth.password 配置绑定到 Java 对象中。

  • AuthConfiguration 使用 @EnableConfigurationProperties(AuthProperties.class) 让这个配置类生效。这样 Spring 启动时会自动读取 application.yml 中的 auth.xxx,并注入到 AuthProperties 中,实现配置与业务解耦。

  • PemUtils 负责读取 RSA 公钥和私钥。它先把 Resource 读成字符串,再去掉 PEM 头尾标记和空白字符,得到真正的 Base64 密钥内容,之后用 Base64.getDecoder().decode() 解码为原始字节数组。

  • 私钥使用 PKCS8EncodedKeySpec,公钥使用 X509EncodedKeySpec,最后交给 KeyFactory.getInstance("RSA") 生成 RSAPrivateKey 和 RSAPublicKey 对象。

  • AuthConfiguration 不是“解密文件”,而是认证相关的 Bean 配置类。它主要注册了 PasswordEncoder、JwtEncoder、JwtDecoder。

  • jwtEncoder 的创建过程是:先读取公私钥,再组装成 RSAKey,再放入 JWKSet,最后创建 NimbusJwtEncoder。这样 JwtEncoder 内部就持有了 RSA 私钥,可用于后续 JWT 签名。

  • jwtDecoder 的创建过程是:读取公钥并创建 NimbusJwtDecoder。这样后续收到 JWT 时就可以用公钥验签。

2. 验证码模块

  • saveCode 会用 scene + identifier 生成唯一 key,格式是 auth:code:{scene}:{identifier},然后用 Redis Hash 存储 code、maxAttempts、attempts,并设置验证码 TTL。

  • verify 校验时也是先根据 scene + identifier 找到对应 key,再取出数据进行判断。

  • 如果 key 不存在,返回 NOT_FOUND。当前实现里验证码过期后 Redis key 会直接消失,因此实际表现也是 NOT_FOUND,并没有单独返回 EXPIRED。

  • 如果尝试次数已经达到上限,返回 TOO_MANY_ATTEMPTS。

  • 如果验证码匹配成功,会删除 Redis key,防止重复使用,并返回 SUCCESS。

  • 如果验证码不匹配,则把 attempts 加一;如果加一后达到上限,会把该 key 的 TTL 额外设置为 30 分钟,相当于短时锁定。

  • VerificationService 负责验证码业务逻辑。

  • sendCode 会先校验参数,再执行发送间隔限制。它使用的 key 是 auth:code:last:{scene}:{identifier}。如果这个 key 还存在,说明还没到下一次发送时间,就直接拒绝发送。

  • sendCode 还会执行每日次数限制。它使用的 key 是 auth:code:count:{scene}:{identifier}:{yyyyMMdd}。这表示它是“按日期计数”,不是严格的“24 小时滚动封禁”。到了第二天会换一个新 key,计数自动重新开始。

  • 然后它会生成一个随机的数字验证码,长度由配置 codeLength 决定,默认是 6 位。

  • 生成后会调用 RedisVerificationCodeStore.saveCode(...) 存储验证码,再调用 CodeSender.sendCode(...) 执行发送。

3. JWT 模块

  • JwtService 负责签发和解析 JWT。

  • issueTokenPair 会先生成一个 refreshTokenId,也就是 refresh token 的 jti,然后分别生成 access token 和 refresh token。

  • encodeToken 和 encodeRefreshToken 都会先用 JwtClaimsSet 构造 claims,也就是 JWT payload 中要存放的数据。

  • access token 中保存了 iss、iat、exp、sub、jti、token_type=access、uid、nickname 等信息。

  • refresh token 中保存了 iss、iat、exp、sub、jti、token_type=refresh、uid,信息更少。

  • 这里要特别注意:JwtClaimsSet 只是“准备载荷数据”,并不是只生成了 JWT 的第二部分。

  • 真正生成完整 JWT 的地方是 jwtEncoder.encode(JwtEncoderParameters.from(claims)).getTokenValue()。这一步会在内部完成三件事:构造 header、构造 payload、再用 RSA 私钥做 RS256 签名,最终拼成完整的 header.payload.signature。

  • 所以私钥并不是你在业务代码里手写调用的,而是已经提前装进了 NimbusJwtEncoder,当你调用 jwtEncoder.encode(...) 时,框架内部自动使用私钥完成签名。

  • JWT 这里不是“加密/解密”,而是“签名/验签”。JWT 的 header 和 payload 默认只是 Base64URL 编码,不是加密,拿到 token 的人都可以解码看到里面的内容,因此敏感数据不能直接放进去。

  • SecurityConfig 开启了 .oauth2ResourceServer(oauth -> oauth.jwt(...)),因此受保护接口上的 access token 会由 Spring Security 自动使用 JwtDecoder 进行验签。

  • JwtDecoder 是在 AuthConfiguration中基于 RSA 公钥创建的,所以 access token 的验签本质上是“框架自动用公钥验签”。

  • refresh token 的校验则是在业务中手动调用 jwtDecoder.decode(token) 完成的,例如 refresh 和 logout 流程。

4. Refresh Token 存储与认证业务

  • RedisRefreshTokenStore 并没有把整个 refresh token 保存到 Redis,而是只保存它的 jti 白名单,key 格式是 auth:rt:{userId}:{tokenId}。

  • 登录或注册成功后,系统会生成 token pair,并把 refresh token 的 jti 写入 Redis,TTL 与 refresh token 过期时间一致。

  • 刷新 token 时,系统先验签,再检查 token_type 必须是 refresh,再检查 Redis 白名单中是否存在对应 jti。

  • 如果合法,就签发一对新的 token,同时撤销旧的 refresh token,并把新的 jti 存入 Redis,这就是 refresh token 轮换。

  • AuthService 是整个认证模块的核心业务层。

  • sendCode:校验手机号/邮箱格式,按场景判断账号是否应该存在,然后调用验证码服务发送验证码。

  • register:校验是否同意协议,校验验证码,创建用户,若传入密码则进行密码策略校验和 BCrypt 哈希,然后签发 token pair,保存 refresh token 白名单,记录审计日志。

  • login:支持两种登录方式,密码登录或验证码登录。成功后签发 token pair,保存 refresh token 白名单,并记录登录日志。

  • refresh:校验 refresh token 的合法性和白名单状态,成功后签发新的 token pair,撤销旧 refresh token。

  • logout:如果 refresh token 合法,则撤销对应的 refresh token 白名单记录。

  • resetPassword:通过验证码重置密码,并撤销该用户的全部 refresh token,强制所有会话重新登录。

  • me:从当前 JWT 中提取 uid,查询并返回当前用户信息。

三.计数模块

1. 模块核心思路

  • 事实层用位图存“某个用户是否点赞/收藏过某个实体”。

  • 汇总层用 SDS 二进制字符串存“这个实体当前有多少点赞/收藏总数”。

  • Kafka 负责把“状态变化”变成异步计数增量,最终折叠到 SDS。

  • 本地 Spring 事件不是做总数统计的,而是做当前服务实例里的快速处理,比如缓存失效、缓存旁路更新。

2. 关键类和配置

  • BitmapShard:做位图分片。

  • 每个分片固定 32768 个 bit,也就是 4KB。

  • 分片规则不是“userid / % chunksize”,准确说法是:

  • chunk = userId / 32768

  • bit = userId % 32768

  • CounterKeys:定义 Redis key。

  • SDS 汇总 key:cnt:v1:{etype}:{eid}

  • 位图 key:bm:{metric}:{etype}:{eid}:{chunk}

  • 聚合桶 key:agg:v1:{etype}:{eid}

  • UserCounterKeys:定义用户维度计数 key,这里先不展开实现。

  • CounterSchema:定义 SDS 固定结构。

  • SCHEMA_ID = v1

  • FIELD_SIZE = 4,不是 8 字节

  • SCHEMA_LEN = 5

  • 当前只真正用了两个指标:

  • like -> idx=1

  • fav -> idx=2

  • 其他位置是预留的。

  • CounterEvent:Kafka 和本地事件共同使用的事件对象,里面有 entityType、entityId、metric、idx、userId、delta。

  • CounterTopics.EVENTS:Kafka 主题名,当前是 counter-events。

3. 写入主流程

  • 用户点赞/取消点赞/收藏/取消收藏,都会走 CounterServiceImpl 的 like/unlike/fav/unfav。

  • 这些方法最终都进入 toggle()。

  • toggle() 先根据 userId 算出 chunk 和 bit,找到这个用户在对应位图里的位置。

  • 然后执行 Lua 脚本,对 bitmap 做原子切换:

  • 点赞时,如果原来是 0,改成 1,返回成功。

  • 如果原来已经是 1,就不改,返回失败。

  • 取消点赞时同理。

  • 这一步是幂等关键:只有状态真的发生变化,后面才会继续发事件。

4. 为什么点赞关系会立刻生效

  • 点赞关系是否成立,看的是 bitmap,不是 Kafka。

  • 也就是说,toggle() 里 Lua 成功把 bit 从 0 改成 1 的那一刻:

  • “这个用户已经点赞了这篇内容”这个事实就成立了。

  • isLiked() / isFaved() 立刻就能从 bitmap 读到新状态。

  • 所以点赞/收藏关系是实时生效的。

  • 延迟的不是“关系”,而是“总数”。

5. Kafka 计数聚合流程

  • toggle() 成功后会发两种事件:

  • 一种发 Kafka:用于异步总数统计。

  • 一种发本地 Spring 事件:用于快速旁路处理。

  • Kafka 这条链路里,CounterAggregationConsumer.onMessage() 监听 counter-events。

  • 收到消息后,把 JSON 转成 CounterEvent。

  • 然后构建聚合桶 key:agg:v1:{etype}:{eid}

  • 用 idx 作为 hash field,用 delta 作为增量,执行 HINCRBY。

  • 这里的 ack 是“写入聚合 hash 成功后 ack”,不是“写入 SDS 成功后 ack”。

6. flush 流程

  • flush() 每隔 1 秒执行一次。

  • 它会扫描所有聚合桶 agg:v1:*。

  • 对每个聚合桶,取出里面所有 field=idx, value=delta。

  • 再从聚合桶 key 里解析出 etype 和 eid,构造对应的 SDS key。

  • 然后执行 Lua,把某个 idx 上累积的 delta 折叠进 SDS 二进制结构里。

  • 成功后删除这个 hash field。

  • 如果整个 hash 空了,再删掉这个聚合桶 key。

7. SDS 是怎么存总数的

  • SDS 这里本质上是 Redis String 里存的一段固定长度二进制。

  • 当前总长度是 5 * 4 = 20 字节。

  • 每个指标占 4 个字节。

  • like 和 fav 分别写在固定偏移位置。

  • 所以总数查询不是扫位图,而是直接按偏移读取 SDS。

8. getCounts() 读取流程

  • getCounts() 先拿到 SDS key。

  • 然后用原始连接 getRaw() 直接把 Redis 里的二进制原始字节拿出来。

  • 这里必须读原始字节,因为这是固定结构二进制,不适合按普通字符串逻辑处理。

  • 如果 SDS 存在且长度正确,就按固定偏移读取 like/fav 对应位置的值,返回结果。

  • 如果 SDS 不存在,或者长度不对,就认为需要重建。

9. SDS 重建流程

  • 如果需要重建,先看当前是否在退避期。

  • 如果还在退避期,直接把请求的指标先返回 0,避免热点击穿。

  • 如果不在退避期,再看限流器是否允许这次重建。

  • 如果限流器不允许,也先返回 0,并提升退避等级。

  • 如果允许,再尝试加分布式锁。

  • 没抢到锁,同样走退避降级。

  • 抢到锁后,才真正根据 bitmap 分片做 BITCOUNT 汇总。

  • 汇总出每个指标的真实值后,重新写回 SDS。

  • 然后删除对应聚合桶中的相关 field,避免重复叠加。

  • 最后重置退避状态。 

10. kafka灾难重建

  • 但如果发生以下灾难场景中,仅依赖在线修复机制

  • Redis 中关键计数 Key(SDS)丢失或被误删;

  • SDS 内容损坏或结构不可反序列化;

  • Redis 聚合桶整体不可用;

  • 需要跨实体、跨时间窗口进行全量计数恢复

  • 这个时候需要kafka从最开始数据进行读取

  • Redis SDS 作为最终计数承载结构;

  • Kafka 事件作为唯一可追溯事实;

  • 至少一次 + 幂等折叠确保恢复安全;

  • 与线上链路隔离,避免干扰业务流量。

四.用户关系模块

1.完整主流程

  • 用户发起关注操作。

  • 系统先做限流,防止同一个用户短时间内高频关注。

  • 限流通过后,先写入“关注表”,表示 A 关注了 B。

  • 同一个业务事务里,再写入一条 outbox 事件记录。

  • outbox 里的 payload 保存的是关系事件 JSON,比如:
    FollowCreated、FollowCanceled。

  • 数据库事务提交后,这条 outbox 记录会进入 MySQL binlog。

  • Canal 作为 MySQL 的伪从库订阅 binlog,拿到 outbox 表的增量变化。

  • Canal 把 binlog 解析成结构化数据,区分出这是哪张表、哪种事件类型、哪些列发生了变化。

  • 桥接层只关心 outbox 表的 INSERT/UPDATE 行变更。

  • 对每一条变更行,桥接层从列数据中提取 payload 字段。

  • 然后把这些 payload 重新包装成一个 JSON 消息。

  • 这个 JSON 最后会被序列化成字符串,发送到 Kafka。

  • Kafka consumer 收到消息后,先把外层 JSON 拆开。

  • 然后逐条取出其中的 payload。

  • 再把 payload 反序列化成关系事件对象。

  • 如果事件是“关注创建”,就执行后续同步逻辑:
    写粉丝表、更新 Redis 关注/粉丝列表、更新关注数和粉丝数。

  • 如果事件是“取消关注”,就执行反向逻辑:
    取消粉丝关系、删除 Redis 列表里的对应成员、扣减计数。

  • 最终形成:
    following 负责主写路径,
    follower + Redis + counter 负责异步补齐,
    达到最终一致。

2.Canal 到 Kafka 的详细流程

  • Canal 连接到本地 Canal 服务。

  • 建立连接时会配置主机、端口、实例名、用户名、密码等信息。

  • 然后设置订阅过滤规则,只关心 outbox 表。

  • 连接建立后,会先回滚到上一次未确认的位置,从而继续消费尚未完成确认的增量数据。

  • 接着进入循环拉取消息的过程。

  • 每次按批次拉取一批数据,比如一批最多 1000 条。

  • 如果这批没有数据,就线程休眠一小段时间,再继续拉取。

  • 如果拉到了数据,就遍历这一批里的每个 entry。

  • 只处理行级数据事件,忽略非行级事件。

  • 然后把 entry 中的二进制内容解析成 RowChange。

  • 这一步之所以会包一层异常处理,是为了避免某条脏数据或解析失败直接把整个消费循环打断。

  • 解析出 RowChange 后,再判断事件类型。

  • 这里只继续处理 INSERT 和 UPDATE。

  • 然后遍历每一行变更后的列数据。

  • 从列里找到 payload 字段,把它取出来。

  • 每一行提取出的 payload 会先放到一个 JSON 对象里。

  • 多行数据会组成一个 JSON 数组。

  • 最后再构造一个总消息:
    包含表名、事件类型、数据数组。

  • 因为 Kafka 接收的是字符串消息,所以最后要把 JSON 对象序列化成字符串。

  • 序列化完成后,把字符串发送到 Kafka topic。

  • 当前这批消息全部处理完后,再对 Canal 位点做 ack,表示这一批已经消费到这里。

3.Kafka consumer 处理关系消息的详细流程

  • consumer 监听 Kafka 中的关系 outbox 主题。

  • 收到消息后,先解析最外层 JSON。

  • 判断这条消息是不是 outbox 表的 INSERT/UPDATE 数据。

  • 如果不是,就直接忽略。

  • 如果是,就取出其中的 data 数组。

  • data 数组里的每一项代表一条 outbox 行。

  • 每一项里再取出 payload。

  • payload 本身是一个 JSON 字符串,里面保存真正的业务事件。

  • 然后把 payload 反序列化成关系事件对象。

  • 接下来先做去重,避免消息重复投递时把同一事件重复处理。

  • 如果去重发现已经处理过,就直接结束。

  • 如果是第一次处理:

    • 关注创建事件:补写粉丝表,更新 Redis 关注列表和粉丝列表,关注数 +1,粉丝数 +1。

    • 取消关注事件:取消粉丝表关系,删除 Redis 关注列表和粉丝列表中的成员,关注数 -1,粉丝数 -1。

  • 处理成功后,再确认 Kafka 消费位点。

  • 这样可以保证消息不会轻易丢掉,同时又避免用户请求线程被这些后置更新拖慢。

4.列表查询流程要分开看

4.1 offset / limit 分页查询流程
  • 这种分页适合“第 1 页、第 2 页、第 3 页”这种传统翻页方式。

  • 前端传的是:
    userId + limit + offset。

  • 后端先拼出 Redis key。

  • 如果查关注列表,就查“关注 ZSet”。

  • 如果查粉丝列表,就查“粉丝 ZSet”。

  • 第一步先从 Redis ZSet 按倒序区间读取:
    也就是直接取 offset ~ offset + limit - 1 这一段数据。

  • 如果 Redis 命中,说明这一页数据已经在缓存里了:

    • 直接返回这一页的用户 ID 列表。

    • 再批量查这些用户的资料。

    • 最后组装成前端需要的资料列表返回。

  • 如果 Redis 没命中,就进入第二层:
    看本地缓存里有没有这个用户的热点 Top 列表。

  • 这里的本地缓存不是所有用户都用,而是给超大列表用户做热点优化。

  • 如果本地缓存命中:

    • 直接从本地缓存的 Top 列表中按 offset/limit 截取。

    • 取到用户 ID 后,再去批量补用户资料。

    • 最后返回给前端。

  • 如果本地缓存也没命中,就要回源数据库。

  • 回源时不会只查 limit 条,而是会查 offset + limit 这么多,目的是把前面页的数据也一起补进缓存。

  • 同时会再做一个上限保护,防止一次查太深把数据库压垮。

  • 数据库查出来后:

    • 按记录时间顺序把这些数据回填到 Redis ZSet。

    • 给这个 ZSet 设置过期时间。

  • 然后判断这个用户是不是“大列表用户”。

  • 如果是,就把 Redis 里最前面的热点部分再同步到本地缓存里。

  • 本地缓存通常只保留前一小段热点数据,不会无限存。

  • Redis 回填完成后,再次从 Redis 按 offset/limit 读取目标页。

  • 这样做的目的,是统一返回路径,保证最终结果都从缓存结构取出。

  • 取到 ID 列表后,再批量查询用户资料。

  • 最终返回资料列表给前端。

流程:

  • Redis miss 后,再考虑本地热点缓存。

  • 本地缓存 miss 后,才回源数据库。

  • 数据库结果先回填 Redis,再重新按页读取。

  • 大列表用户才会额外维护本地热点缓存。

4.2 cursor 游标分页查询流程
  • 这种分页适合前端“下拉加载更多”。

  • 前端传的是:
    userId + limit + cursor。

  • 这个 cursor 本质上是“上一页最后一条记录的时间戳”。

  • 后端还是先定位 Redis ZSet。

  • 但是这次不是按 offset 截取,而是按 score 范围查。

  • Redis ZSet 里的 score 存的就是关系创建时间戳。

  • 所以后端会去取:
    score <= cursor 的最新一段数据。

  • 如果是第一页,没有 cursor,就相当于从最新数据开始取。

  • 第一步先查 Redis:
    按时间分数倒序取一页数据。

  • 如果 Redis 命中:

    • 直接拿到这页用户 ID。

    • 再批量查资料。

    • 返回给前端。

  • 如果 Redis 没命中:

    • 就去数据库查一批关系数据。

    • 把查到的数据按时间戳作为 score 回填到 Redis ZSet。

    • 然后设置过期时间。

    • 再按同样的 cursor 条件重新从 Redis 取一次。

  • 重新取到数据后:

    • 拿到这一页 ID。

    • 批量查资料。

    • 返回给前端。

  • 这种分页的核心不是“跳到第几页”,而是“从上一次停下的位置继续往后翻”。

  • 所以它更适合无限滚动,也比深度 offset 分页更稳定。

5.为什么要设计两种分页

  • offset 分页适合传统页码场景,前端容易理解,也方便直接跳页。

  • 但 offset 越深,查询和缓存压力越大,所以需要 Redis 和本地热点缓存兜底。

  • cursor 分页适合滚动加载,天然更适合按时间顺序读取。

  • cursor 不强调跳页,而强调连续加载,因此更轻量,也更适合关系流这种“按时间倒序浏览”的场景。

五.ES搜索模块

1.搜索模块完整主流程

  • 应用启动时,先检查 ES 里的搜索索引是否存在。

  • 如果索引不存在,就先创建索引和 mapping。

  • 索引建好后,再检查当前索引里有没有文档。

  • 如果索引为空,就从数据库中把历史公开已发布内容分页回灌到 ES。

  • 每条内容在写入 ES 前,会先从数据库查详情,再补齐作者信息、标签、图片、正文、计数等字段。

  • 用户搜索时,请求不会先打数据库,而是直接打 ES。

  • ES 返回命中文档后,后端再把命中文档映射成前端需要的搜索结果对象。

  • 如果用户已登录,还会额外补上当前用户对每条内容的 liked/faved 状态。

  • 用户输入前缀时,联想建议接口直接调用 ES 的 Completion Suggester,从 title_suggest 字段拿候选词。

  • 理想设计上,内容新增、修改、删除后,应该通过 outbox -> Canal -> Kafka -> 搜索消费者 -> ES 做增量更新。

  • 但按当前代码现状,更可靠的说法应该是:
    “搜索索引初始化和查询链路已经完整,增量更新消费者已写好,但内容侧事件生产链路没有完全看到。”

2.索引初始化与回灌流程

  • 应用启动。

  • 搜索模块先检查 zhiguang_content_index 这个索引是否存在。

  • 如果不存在,就创建索引。

  • 索引 mapping 里会定义:
    title、body、description、tags、作者信息、发布时间、计数字段、联想字段等。

  • title 和 body 使用 IK 分词器,所以 ES 集群需要提前装好 analysis-ik 插件。

  • 索引创建完成后,再统计当前索引中的文档数。

  • 如果索引已经有数据,就跳过历史回灌。

  • 如果索引为空,就从数据库分页查询公开且已发布的内容。

  • 每次分页拉一批内容。

  • 对每条内容调用一次 upsert 索引逻辑。

  • upsert 时先查数据库详情。

  • 然后组装 ES 文档字段:
    内容 ID、内容类型、标题、描述、作者 ID、作者头像、作者昵称、作者标签、发布时间、状态、标签数组、图片数组、置顶状态等。

  • 接着尝试根据 contentUrl 抓取正文内容。

  • 如果正文抓取失败,就退化为只使用 description。

  • 正文内容太长时会截断,避免索引文档过大。

  • 然后再从计数服务读取 like 和 fav 的聚合计数,写入 like_count 和 favorite_count。

  • view_count 当前直接写成 0。

  • 如果标题不为空,还会把标题写到 title_suggest,用于联想建议。

  • 最后把这份文档写入 ES,并使用 refresh=wait_for,保证写完之后很快就能被搜到。

3.搜索索引增量更新流程

  • 理想链路是:
    内容发生变化。

  • 业务层往 outbox 写一条搜索事件。

  • MySQL binlog 记录这条 outbox 变化。

  • Canal 订阅到 outbox 表的增量。

  • 桥接层把 payload 发到 Kafka 的 canal-outbox 主题。

  • 搜索模块的 Kafka consumer 监听这个主题。

  • consumer 解析消息里的 payload。

  • 如果 entity=knowpost,并且操作是普通更新,就执行 upsertKnowPost(id)。

  • 如果操作是删除,就执行 softDeleteKnowPost(id)。

  • upsert 会重新从数据库拉最新详情,再覆盖写入 ES。

  • softDelete 不是真删 ES 文档,而是把文档状态改成 deleted。

  • 这样可以保证索引更新是幂等的,同一个 ID 多次覆盖写不会乱。

4.关键词搜索流程

  • 前端传入:
    q、size、可选 tags、可选 after。

  • q 是搜索关键词。

  • tags 是逗号分隔的标签过滤条件。

  • after 是上一页最后一条命中的游标。

  • 后端收到请求后,先解析当前登录用户。

  • 如果用户已登录,后面可以补 liked/faved。

  • 然后把 tagsCsv 解析成标签列表。

  • 再把 after 从 Base64URL 解码回 ES 的排序值数组。

  • 接下来构造 ES 查询。

  • 查询主体是 multi_match,主要搜 title 和 body。

  • 其中 title 权重更高,权重大约是 body 的 3 倍。

  • 然后加过滤条件:
    只搜 status=published 的内容。

  • 如果用户传了标签,再额外按 tags 做 terms 过滤。

  • 在相关性得分之外,还会叠加业务加权。

  • 当前加权字段主要是:
    like_count 和 view_count。

  • 加权方式是 field_value_factor + log1p,避免大数值把排序拉得过猛。

  • 最终排序顺序大致是:
    先按相关性分数,再按发布时间,再按点赞数,再按浏览数,最后按内容 ID 稳定排序。

  • 如果前端传了 after,搜索就不是从第一页开始,而是从上次最后一个命中的后面继续往下查。

  • ES 返回结果后,后端开始逐条组装返回对象。

  • 每条命中会取出标题、描述、标签、图片、作者信息、点赞数、收藏数等字段。

  • 如果 ES 返回了高亮片段,就把标题高亮和正文高亮合并成一个 snippet。

  • 如果没有高亮,就退回使用原始 description。

  • 图片列表里第一张图会被当作封面图。

  • 如果用户已登录,还会再去计数模块判断当前用户是否点过赞、是否收藏过。

  • 所以最终返回给前端的不只是“搜到了什么”,还带有用户态信息。

  • 最后,后端会把最后一条命中的排序值重新编码成 nextAfter。

  • 前端下一页继续传这个 nextAfter,就能做游标续翻。

  • 响应里还会带 hasMore,表示前端是否继续展示“加载更多”。

5.联想建议流程

  • 前端传一个标题前缀 prefix。

  • 后端不查数据库,直接查 ES。

  • ES 使用的是 Completion Suggester。

  • 联想字段是 title_suggest。

  • 也就是说,联想建议本质上是基于标题做前缀补全。

  • ES 返回候选项后,后端把候选标题提取出来,组成字符串列表返回前端。

  • 这条链路比主搜索更轻,没有高亮、没有复杂排序、没有用户态补充。

六.AI模块

AI / LLM 模块整体定位

  • 这个模块主要服务 knowpost。

  • 一条线负责根据正文自动生成一段简短描述。

  • 另一条线负责把一篇知文切片、向量化,然后基于这篇知文做 RAG 问答。

  • 聊天能力由 DeepSeek 提供。

  • 向量检索依赖 Spring AI + Elasticsearch VectorStore。

  • 向量 embedding 模型配置的是 text-embedding-v4,维度 1536。

  • 向量库索引名配置的是 zhiguang-ai-index。

1、描述生成功能的完整流程

  • 前端把知文正文发给后端。

  • 后端调用“描述生成”接口。

  • 接口只接收正文内容,不需要先落库。

  • 服务层先校验正文是否为空。

  • 如果正文为空,直接报错,不会调模型。

  • 校验通过后,系统会构造两段提示词:
    一段是 system prompt,
    一段是 user prompt。

  • system prompt 会明确告诉模型:
    你是中文文案编辑,
    需要基于正文生成一个简洁、有吸引力、且不超过 50 个汉字的中文描述,
    不要输出解释,不要分段。

  • user prompt 会把用户正文拼进去,再要求模型直接给出短描述。

  • 然后调用 DeepSeek 聊天模型。

  • 当前参数大致是:
    model = deepseek-chat
    temperature = 0.8
    maxTokens = 120

  • 模型返回结果后,后端不会直接原样返回。

  • 它还会做一次后处理。

  • 后处理第一步是做字符标准化。

  • 比如统一全角半角、清理换行、多空格。

  • 然后会去掉前后多余的引号和尾部标点。

  • 最后再按字符数做一次硬截断。

  • 截断规则不是简单按字节,而是按 code point 计数。

  • 超过 50 个字符就只保留前 50 个。

  • 处理完成后,后端把这段描述返回给前端。

  • 所以这个功能的本质就是:
    正文 -> 提示词 -> DeepSeek 生成 -> 后处理清洗 -> 返回短描述

2、RAG 建索引流程

  • RAG 的目标不是给整站建统一问答索引。

  • 它的目标是给“某一篇知文”建立向量切片。

  • 当系统需要确保一篇知文可问答时,会触发一次 ensureIndexed。

  • 这个方法本质上会去执行“单篇重建索引”逻辑。

  • 重建开始前,先查数据库里的知文详情。

  • 如果这篇知文不存在,直接结束。

  • 如果这篇知文不是 published,或者不是 public,也直接跳过,不给它建向量索引。

  • 如果内容地址 contentUrl 不存在,也无法继续,因为没法抓正文。

  • 接着会做一次“指纹判断”。

  • 指纹优先用 contentSha256。

  • 如果没有 Sha256,就退而用 ETag。

  • 系统会先去向量索引里找这篇文章现有任意一条切片的 metadata。

  • 然后比较旧指纹和当前指纹是否一致。

  • 如果一致,说明正文没变,就跳过重建。

  • 如果不一致,才继续真正重建。

  • 真正重建时,先通过 contentUrl 把正文抓下来。

  • 当前默认当作文本直接拉,不像 search 模块那样做复杂字符集兜底。

  • 抓到正文后,先做切片。

  • 切片第一层是按 Markdown 标题分段。

  • 也就是遇到 # 开头的标题,会把前一段收起来,开始新段。

  • 分段之后,再做固定长度切片。

  • 每段如果不超过 800 字符,就直接作为一个 chunk。

  • 如果超过 800 字符,就继续切成多个 chunk。

  • chunk 之间保留 100 字符重叠。

  • 这样做是为了让语义上下文不要在边界处断得太狠。

  • 切好之后,系统会先删除这篇文章旧的所有向量切片。

  • 删除条件是 metadata.postId = 当前文章ID。

  • 然后再组装新的向量文档。

  • 每个 chunk 都会带一份 metadata。

  • metadata 里大致有:
    postId
    chunkId
    position
    contentEtag
    contentSha256
    contentUrl
    title

  • 然后批量写入向量库。

  • 写入完成后,这篇文章就可以被后续 RAG 检索到了。

  • 所以 RAG 建索引的主线就是:
    查知文 -> 校验状态 -> 比较指纹 -> 抓正文 -> Markdown 分段 -> 800字切片 + 100字重叠 -> 删旧切片 -> 写新向量文档

3、RAG 问答流程

  • 前端调用某一篇知文的问答接口。

  • 请求里会带:
    postId
    question
    topK
    maxTokens

  • 后端收到问题后,第一步不是立刻问模型。

  • 而是先做一次 ensureIndexed。

  • 这一步相当于问答前的兜底:
    如果索引不存在就建,
    如果指纹没变就跳过,
    如果已经是最新就不重复写。

  • 然后开始做向量检索。

  • 检索时会先做“宽召回”。

  • 也就是不是只取 topK 条,而是先取更大的集合。

  • 当前规则大致是:
    fetchK = max(topK * 3, 20)

  • 先多召回一些,是因为后面还要做一次服务端过滤。

  • 当前向量检索没有直接把 postId 过滤条件下推给向量库。

  • 所以后端会把召回结果拿回来后,再逐条筛掉“不是当前这篇知文”的切片。

  • 只保留 metadata.postId = 当前文章ID 的 chunk。

  • 然后最多留下前 topK 条作为上下文。

  • 接下来把这些上下文拼起来。

  • chunk 和 chunk 之间会用分隔符隔开,方便模型理解上下文块。

  • 然后组装新的提示词。

  • system prompt 会强调:
    你是中文知识助手,
    只能依据给定上下文回答,
    不能确定时就明确说不确定。

  • user prompt 会把“用户问题 + 检索到的上下文”一起发给模型。

  • 然后再调用 DeepSeek 聊天模型。

  • 这次参数和描述生成不同。

  • 当前更偏向稳健回答:
    temperature = 0.2
    maxTokens = 前端传入值

  • 返回方式也不是一次性整段返回。

  • 这里走的是流式输出。

  • 后端把模型输出包装成 Flux<String>。

  • 接口以 SSE 的方式持续往前端推送内容。

  • 所以前端体验上会像“边生成边显示”。

  • 整个问答流程本质就是:
    收到问题 -> 问答前兜底建索引 -> 向量召回 -> 过滤当前文章切片 -> 拼上下文 -> DeepSeek 流式生成答案 -> SSE 返回

4、和知文业务的触发关系

  • 当前 RAG 索引不是完全自动事件驱动同步。

  • 它主要有 3 个触发入口。

  • 第一个入口是“确认正文上传完成”之后。

  • 这时系统会尝试先做一次预索引。

  • 第二个入口是“知文发布成功”之后。

  • 这时系统也会尝试做一次预索引。

  • 第三个入口是“真正发起问答”之前。

  • 这时系统会再次 ensureIndexed,防止前两次漏掉或内容更新后未同步。

  • 此外还有一个手动接口。

  • 管理或调试时,可以显式调用“单篇重建索引”。

  • 它会返回本次重建产生了多少个 chunk。

Logo

AtomGit 是由开放原子开源基金会联合 CSDN 等生态伙伴共同推出的新一代开源与人工智能协作平台。平台坚持“开放、中立、公益”的理念,把代码托管、模型共享、数据集托管、智能体开发体验和算力服务整合在一起,为开发者提供从开发、训练到部署的一站式体验。

更多推荐