在这里插入图片描述

前言

今天我们来整理项目中的笔记维度计数系统,也就是知文的点赞数、收藏数,以及“我是否点赞 / 收藏过”这种用户态判断。

普通做法可能是直接在 MySQL 里加两个字段:

like_count
fav_count

每次点赞就 update count = count + 1。但是在高并发内容社区里,热门笔记会把某一行打成热点,数据库行锁、redo log 和缓存一致性压力都会比较明显。

所以这个项目里采用的是:

Redis 分片位图做状态事实
Lua 脚本保证点赞/取消点赞原子幂等
Kafka 承接计数增量事件
Redis Hash 做临时聚合桶
定时任务把增量折叠到 Redis SDS 二进制计数快照
读取时优先读 SDS,异常时基于位图自愈重建

这套设计的核心思想是:

位图负责“谁点过”,SDS 负责“有多少”,事件负责把状态变化转换成计数增量。


一、计数系统整体流程

笔记维度计数的核心代码主要在:

src/main/java/com/tongji/counter
├── api
│   ├── ActionController.java
│   └── CounterController.java
├── service
│   ├── CounterService.java
│   └── impl/CounterServiceImpl.java
├── schema
│   ├── BitmapShard.java
│   ├── CounterKeys.java
│   └── CounterSchema.java
└── event
    ├── CounterEvent.java
    ├── CounterEventProducer.java
    ├── CounterAggregationConsumer.java
    └── CounterRebuildConsumer.java

一次点赞的完整流程可以理解为:

用户点赞
  ↓
定位 Redis 位图分片
  ↓
Lua 执行 GETBIT + SETBIT
  ↓
只有状态从 0 -> 1 时,才产生 delta = +1 事件
  ↓
Kafka 发送 counter-events
  ↓
消费者写入 Redis 聚合桶 agg:v1:{etype}:{eid}
  ↓
定时任务每 1 秒把聚合增量刷入 SDS
  ↓
读取计数时直接读取 cnt:v1:{etype}:{eid}

这里有两个重要分层:

层级 作用
位图事实层 保存用户是否点赞 / 收藏
SDS 汇总层 保存点赞数 / 收藏数快照
Kafka 事件层 把状态变化异步聚合成计数

二、Redis Key 与 SDS Schema 设计

1. 位图分片

// src/main/java/com/tongji/counter/schema/BitmapShard.java
public final class BitmapShard {
    public static final int CHUNK_SIZE = 32_768;

    public static long chunkOf(long userId) {
        return userId / CHUNK_SIZE;
    }

    public static long bitOf(long userId) {
        return userId % CHUNK_SIZE;
    }
}

这里每个分片是 32768 位,也就是 4KB。

这样做是为了避免 Redis 单个位图 Key 因为用户 ID 过大而变得很稀疏。

Key 规则如下:

// src/main/java/com/tongji/counter/schema/CounterKeys.java
public static String bitmapKey(String metric, String entityType,
                               String entityId, long chunk) {
    return String.format("bm:%s:%s:%s:%d",
            metric, entityType, entityId, chunk);
}

示例:

bm:like:knowpost:10001:0
bm:fav:knowpost:10001:0

2. SDS 固定结构计数

// src/main/java/com/tongji/counter/schema/CounterSchema.java
public final class CounterSchema {
    public static final String SCHEMA_ID = "v1";
    public static final int FIELD_SIZE = 4;
    public static final int SCHEMA_LEN = 5;

    public static final int IDX_LIKE = 1;
    public static final int IDX_FAV = 2;

    public static final Map<String, Integer> NAME_TO_IDX = Map.of(
            "like", IDX_LIKE,
            "fav", IDX_FAV
    );
}

SDS 结构是:

5 个字段 * 4 字节 = 20 字节

当前开放的字段是:

下标 含义
0 预留 read
1 like
2 fav
3 预留 comment
4 预留 repost

SDS Key:

public static String sdsKey(String entityType, String entityId) {
    return String.format("cnt:%s:%s:%s",
            CounterSchema.SCHEMA_ID, entityType, entityId);
}

示例:

cnt:v1:knowpost:10001

相比 Redis Hash:

HSET cnt:10001 like 100 fav 20

SDS 的好处是结构紧凑。5 个计数字段只占 20 字节,按固定偏移读写即可。


三、点赞与收藏接口

行为接口在 ActionController 中。

// src/main/java/com/tongji/counter/api/ActionController.java
@PostMapping("/like")
public ResponseEntity<Map<String, Object>> like(@Valid @RequestBody ActionRequest req,
                                                @AuthenticationPrincipal Jwt jwt) {
    long uid = jwtService.extractUserId(jwt);
    boolean changed = counterService.like(req.getEntityType(), req.getEntityId(), uid);

    return ResponseEntity.ok(Map.of(
            "changed", changed,
            "liked", counterService.isLiked(req.getEntityType(), req.getEntityId(), uid)
    ));
}

请求体:

// src/main/java/com/tongji/counter/api/dto/ActionRequest.java
@Data
public class ActionRequest {
    @NotBlank
    private String entityType;

    @NotBlank
    private String entityId;
}

点赞请求示例:

{
  "entityType": "knowpost",
  "entityId": "10001"
}

响应:

{
  "changed": true,
  "liked": true
}

这里的 changed 很关键。

如果用户已经点过赞,再重复点击点赞,changed=false,也就不会重复增加计数。


四、Lua 实现位图原子幂等

1. Service 层入口

// src/main/java/com/tongji/counter/service/impl/CounterServiceImpl.java
@Override
public boolean like(String entityType, String entityId, long userId) {
    return toggle(entityType, entityId, userId,
            "like", CounterSchema.IDX_LIKE, true);
}

@Override
public boolean unlike(String entityType, String entityId, long userId) {
    return toggle(entityType, entityId, userId,
            "like", CounterSchema.IDX_LIKE, false);
}

@Override
public boolean fav(String entityType, String entityId, long userId) {
    return toggle(entityType, entityId, userId,
            "fav", CounterSchema.IDX_FAV, true);
}

@Override
public boolean unfav(String entityType, String entityId, long userId) {
    return toggle(entityType, entityId, userId,
            "fav", CounterSchema.IDX_FAV, false);
}

这四个方法最终都会走 toggle


2. 位图定位与事件发送

private boolean toggle(String etype, String eid, long uid,
                       String metric, int idx, boolean add) {
    long chunk = BitmapShard.chunkOf(uid);
    long bit = BitmapShard.bitOf(uid);

    String bmKey = CounterKeys.bitmapKey(metric, etype, eid, chunk);

    Long changed = redis.execute(
            toggleScript,
            List.of(bmKey),
            String.valueOf(bit),
            add ? "add" : "remove"
    );

    boolean ok = changed == 1L;

    if (ok) {
        int delta = add ? 1 : -1;
        CounterEvent event = CounterEvent.of(etype, eid, metric, idx, uid, delta);

        eventProducer.publish(event);
        eventPublisher.publishEvent(event);
    }

    return ok;
}

这里的逻辑是:

  1. 根据 userId 定位分片 chunk
  2. 根据 userId 定位分片内偏移 bit
  3. Lua 脚本原子判断并修改位图
  4. 只有真的发生状态变化时,才发送计数事件

3. Lua 脚本

-- src/main/java/com/tongji/counter/service/impl/CounterServiceImpl.java
local bmKey = KEYS[1]
local offset = tonumber(ARGV[1])
local op = ARGV[2]

local prev = redis.call('GETBIT', bmKey, offset)

if op == 'add' then
  if prev == 1 then return 0 end
  redis.call('SETBIT', bmKey, offset, 1)
  return 1
elseif op == 'remove' then
  if prev == 0 then return 0 end
  redis.call('SETBIT', bmKey, offset, 0)
  return 1
end

return -1

这段 Lua 解决了两个问题:

  1. GETBITSETBIT 必须原子执行。
  2. 重复点赞、重复取消不能重复产生计数。

也就是说,位图不只是缓存,它是点赞/收藏状态的事实层。


五、计数事件与 Kafka 聚合

点赞状态变化后,会产生一个 CounterEvent

// src/main/java/com/tongji/counter/event/CounterEvent.java
@Data
public class CounterEvent {
    private String entityType;
    private String entityId;
    private String metric;
    private int idx;
    private long userId;
    private int delta;
}

发送到 Kafka:

// src/main/java/com/tongji/counter/event/CounterEventProducer.java
public void publish(CounterEvent event) {
    try {
        String payload = objectMapper.writeValueAsString(event);
        kafka.send(CounterTopics.EVENTS, payload);
    } catch (JsonProcessingException e) {
        // 生产异常不影响主流程,可接入告警
    }
}

主题名:

public final class CounterTopics {
    public static final String EVENTS = "counter-events";
}

这里设计成异步事件,是为了让用户点赞动作尽快返回,不阻塞在计数汇总上。


六、Redis 聚合桶与 SDS 刷写

1. 事件写入聚合桶

// src/main/java/com/tongji/counter/event/CounterAggregationConsumer.java
@KafkaListener(topics = CounterTopics.EVENTS, groupId = "counter-agg")
public void onMessage(String message, Acknowledgment ack) throws Exception {
    CounterEvent evt = objectMapper.readValue(message, CounterEvent.class);

    String aggKey = CounterKeys.aggKey(evt.getEntityType(), evt.getEntityId());
    String field = String.valueOf(evt.getIdx());

    try {
        redis.opsForHash().increment(aggKey, field, evt.getDelta());
        ack.acknowledge();
    } catch (Exception ex) {
        // 不提交位点,下次重试
    }
}

聚合桶 Key:

agg:v1:knowpost:10001

Hash 结构:

field = 指标 idx
value = 当前待刷写 delta

例如:

HINCRBY agg:v1:knowpost:10001 1 1
HINCRBY agg:v1:knowpost:10001 2 -1

2. 定时刷写到 SDS

@Scheduled(fixedDelay = 1000L)
public void flush() {
    Set<String> keys = redis.keys("agg:" + CounterSchema.SCHEMA_ID + ":*");

    for (String aggKey : keys) {
        Map<Object, Object> entries = redis.opsForHash().entries(aggKey);
        if (entries.isEmpty()) {
            continue;
        }

        String[] parts = aggKey.split(":", 4);
        String cntKey = CounterKeys.sdsKey(parts[2], parts[3]);

        for (Map.Entry<Object, Object> e : entries.entrySet()) {
            String field = String.valueOf(e.getKey());
            long delta = Long.parseLong(String.valueOf(e.getValue()));
            int idx = Integer.parseInt(field);

            redis.execute(incrScript, List.of(cntKey),
                    String.valueOf(CounterSchema.SCHEMA_LEN),
                    String.valueOf(CounterSchema.FIELD_SIZE),
                    String.valueOf(idx),
                    String.valueOf(delta));

            redis.execute(decrScript, List.of(aggKey), field, String.valueOf(delta));
        }

        if (redis.opsForHash().size(aggKey) == 0L) {
            redis.delete(aggKey);
        }
    }
}

这一步每 1 秒执行一次,把 Redis Hash 中的增量折叠到 SDS。

也就是说,计数读数是秒级最终一致。


3. SDS 原子更新 Lua

local cntKey = KEYS[1]
local schemaLen = tonumber(ARGV[1])
local fieldSize = tonumber(ARGV[2])
local idx = tonumber(ARGV[3])
local delta = tonumber(ARGV[4])

local cnt = redis.call('GET', cntKey)
if not cnt then
  cnt = string.rep(string.char(0), schemaLen * fieldSize)
end

local off = idx * fieldSize
local v = read32be(cnt, off) + delta

if v < 0 then v = 0 end

local seg = write32be(v)
cnt = string.sub(cnt, 1, off) .. seg .. string.sub(cnt, off + fieldSize + 1)

redis.call('SET', cntKey, cnt)
return 1

这里有两个细节:

  • SDS 不存在时,会初始化为全 0 的固定长度字节串。
  • 如果减计数后小于 0,会归 0,避免出现负数。

七、读取计数与自愈重建

1. 读取接口

// src/main/java/com/tongji/counter/api/CounterController.java
@GetMapping("/{etype}/{eid}")
public ResponseEntity<CountsResponse> getCounts(@PathVariable("etype") String entityType,
                                                @PathVariable("eid") String entityId,
                                                @RequestParam(value = "metrics", required = false) String metricsStr) {
    List<String> metrics;

    if (metricsStr == null || metricsStr.isBlank()) {
        metrics = new ArrayList<>(CounterSchema.SUPPORTED_METRICS);
    } else {
        metrics = Arrays.stream(metricsStr.split(","))
                .map(String::trim)
                .filter(CounterSchema.SUPPORTED_METRICS::contains)
                .toList();
    }

    Map<String, Long> counts = counterService.getCounts(entityType, entityId, metrics);
    return ResponseEntity.ok(new CountsResponse(entityType, entityId, counts));
}

请求示例:

GET /api/v1/counter/knowpost/10001?metrics=like,fav

响应:

{
  "entityType": "knowpost",
  "entityId": "10001",
  "counts": {
    "like": 128,
    "fav": 35
  }
}

2. 正常读取 SDS

byte[] raw = getRaw(sdsKey);
boolean needRebuild = (raw == null || raw.length != expectedLen);

if (!needRebuild) {
    for (String m : metrics) {
        Integer idx = CounterSchema.NAME_TO_IDX.get(m);
        int off = idx * CounterSchema.FIELD_SIZE;
        long val = readInt32BE(raw, off);
        result.put(m, val);
    }
}

正常情况下,只需要一次 Redis GET,然后按固定偏移解析即可。


3. SDS 异常时基于位图重建

如果 SDS 不存在,或者长度不对,就进入重建逻辑:

long sum = bitCountShardsPipelined(m, entityType, entityId);
writeInt32BE(newSds, idx * CounterSchema.FIELD_SIZE, sum);
result.put(m, sum);

重建方式是扫描所有位图分片:

private long bitCountShardsPipelined(String metric, String etype, String eid) {
    String pattern = String.format("bm:%s:%s:%s:*", metric, etype, eid);

    Set<String> keys = redis.keys(pattern);
    if (keys.isEmpty()) return 0L;

    List<Object> res = redis.executePipelined((RedisCallback<Object>) connection -> {
        for (String k : keys) {
            connection.stringCommands().bitCount(k.getBytes(StandardCharsets.UTF_8));
        }
        return null;
    });

    long sum = 0L;
    for (Object o : res) {
        if (o instanceof Number n) {
            sum += n.longValue();
        }
    }
    return sum;
}

这里的事实来源是位图,所以即使 SDS 丢了,也能根据 BITCOUNT 重建。

同时项目还做了:

  • Redisson 分布式锁,避免并发重建
  • RateLimiter 限流,防止热点内容重建风暴
  • 指数退避,避免反复失败时打爆 Redis
  • 重建后删除聚合桶字段,避免重复加算

八、Feed 中如何使用计数

Feed 列表中不能把用户态信息直接缓存进公共缓存。

所以项目里是这样做的:

Map<String, Long> counts = counterService.getCounts(
        "knowpost",
        String.valueOf(base.id()),
        List.of("like", "fav")
);

boolean liked = uid != null && counterService.isLiked("knowpost", base.id(), uid);
boolean faved = uid != null && counterService.isFaved("knowpost", base.id(), uid);

也就是说:

  • 点赞数 / 收藏数:读 SDS
  • 我是否点赞 / 收藏:实时读位图

这样可以避免公共 Feed 缓存被某个用户的 liked/faved 状态污染。


九、知识点总结

1. 为什么不用 MySQL 计数列?

因为点赞、收藏是高频行为,热门笔记会造成数据库热点行更新。

计数允许秒级最终一致,没有必要每次都同步写 MySQL。

2. 为什么不用 Redis Hash 直接计数?

Redis Hash 更简单,但每个 field/value 都有额外结构开销。

项目用定长 SDS 字节串,5 个指标只需要 20 字节,更紧凑,也方便按 Schema 扩展。

3. 为什么位图是事实层?

因为位图天然适合表达“某个用户是否做过某个动作”。

同一个用户重复点赞,只会发现位图已经是 1,不会重复产生计数事件。

4. 自愈重建解决了什么问题?

当 SDS 丢失或结构异常时,可以通过位图分片 BITCOUNT 重新计算真实计数。

这让计数系统不仅能跑得快,还能在异常时恢复。


总结

这一篇主要整理了项目中笔记维度的点赞收藏计数系统。

这套方案不是简单的 Redis INCR,而是把状态和计数拆开:位图保存用户行为事实,Kafka 承接状态变化事件,Redis Hash 做短暂聚合,SDS 保存紧凑计数快照,读取异常时再通过位图自愈重建。

它牺牲了一点秒级强一致,换来了高并发写入、低成本读取、幂等控制和异常恢复能力。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐