本文记录了将用户行为分析平台的 ETL 系统从 Spark Streaming 迁移到 Flink 1.17 的完整历程,涵盖架构设计、技术选型、性能优化等多个维度,最终实现 10s 内端到端时效、单并行度 4000+ QPS 的高性能目标。

目录

  1. 项目背景与挑战
  2. 架构设计
  3. 解决方案专题
  4. 技术实现细节
  5. 性能调优实战案例
  6. 数据一致性保障
  7. 监控指标体系
  8. 优化效果总结
  9. 经验与教训

一、项目背景与挑战

1.1 原系统存在的问题

原有的 Spark Streaming ETL 系统在业务快速增长下暴露出多个问题:

问题

表现

影响

性能过低

部分逻辑逐条与外部系统交互,未使用批量

吞吐量上不去,延迟高

实时性不足

Spark 微批模式

端到端延迟 > 5分钟

MySQL 高并发

元数据频繁读取 MySQL

数据库负载飙升,影响稳定性

SSDB 缺陷

停止维护、不支持高可用

运维困难,单点故障风险

覆盖更新

设备/用户表直接覆盖

NULL 值冲掉原有数据

1.2 优化目标

指标

目标值

高峰 QPS

单并行度 4000+

端到端时效

< 10s

数据一致性

Exactly-Once

1.3 技术栈升级

组件

旧方案

新方案

流处理引擎

Spark Streaming

Flink 1.17.2

分布式缓存

SSDB

KVRocks (Redis 兼容)

本地缓存

Caffeine

JSON 库

FastJSON 1.x

FastJSON2

数据仓库

Doris

Doris (优化写入方式)


二、架构设计

2.1 三层流水线架构

系统采用三模块分离架构,各模块职责清晰,通过 Kafka 解耦:

Kafka(SDK) → [Gate Job] → Kafka → [ID Job] → Kafka → [DW Job] → Doris
               数据清洗            ID映射             数据富化

模块

职责

核心操作

Gate Job

数据清洗层

JSON 校验、格式标准化、解密解压、异常数据分流

ID Job

ID 处理层

设备ID/用户ID/平台ID 映射与绑定、属性处理

DW Job

数据仓库层

IP解析、UA解析、关键词富化、数据路由、写入 Doris

架构优势

  • 模块解耦:各模块独立部署,故障隔离,便于扩缩容
  • 可观测性:中间 Topic 便于问题排查和数据回溯
  • 灵活性:各模块可独立升级,不影响其他模块

三、解决方案专题

本章针对原系统的核心问题,按"问题 → 目标 → 方案"的结构梳理完整的解决思路。

3.1 MySQL 高负载解决方案

维度

内容

问题

元数据频繁读取 MySQL,数据库负载飙升,影响稳定性

原因

每条消息处理都直接查询 MySQL 获取配置信息

目标

降低 MySQL 负载 90%+,同时保证数据一致性

方案

元数据定期同步 KVRocks + Hash Tag 原子替换 + 两级缓存

方案架构

元数据访问链路(三级):
1. L1: Caffeine 本地缓存 - 毫秒级响应,refreshAfterWrite 异步刷新
2. L2: KVRocks 分布式缓存 - 同步服务定期从 MySQL 全量同步
3. MySQL 元数据源 - 只在缓存冷启动或新建元数据时访问

Cache-Sync 同步服务(独立进程):
● 定期全量同步 MySQL → KVRocks
● Hash Tag 原子替换:{key}_new → {key}
● 版本号递增通知 Flink 刷新 L1

关键技术点

1) Hash Tag 原子替换
// KVRocks 集群模式下,使用 {} 包裹的部分作为 hash tag
// 保证相同 hash tag 的 key 落在同一个 slot,支持原子操作

// 同步服务写入流程:
// Step 1: 写入新数据到 {EVENT_ID_MAP}_new
kvrocksClient.hset("{EVENT_ID_MAP}_new", field, value);

// Step 2: 原子替换 (RENAME 命令)
kvrocksClient.rename("{EVENT_ID_MAP}_new", "{EVENT_ID_MAP}");

// 好处: 切换瞬间完成,无中间状态,读取方无感知
2) 定期全量刷新 L1 缓存
// ConfigCacheService.java - 定时全量刷新
scheduler.scheduleAtFixedRate(
    this::refreshAllCachesParallel,
    refreshIntervalMinutes,  // 默认 10 分钟
    refreshIntervalMinutes,
    TimeUnit.MINUTES
);

// 并行刷新所有缓存 (HGETALL/SMEMBERS 批量拉取)
private void refreshAllCachesParallel() {
    List<RefreshTask> tasks = Arrays.asList(
        new RefreshTask("eventId", () -> refreshHashCache(EVENT_ID_MAP, eventIdCache)),
        new RefreshTask("eventAttrId", () -> refreshHashCache(EVENT_ATTR_ID_MAP, eventAttrIdCache)),
        // ... 更多缓存
    );
    
    // 多线程并行执行
    tasks.parallelStream().forEach(RefreshTask::execute);
}

效果

  • MySQL 查询量降低 99%+
  • 元数据访问延迟从 5-10ms 降至 < 1ms
  • 数据库 CPU 从 80% 降至 < 10%

3.2 并发写 MySQL 解决方案

维度

内容

问题

多个 Flink 并行实例同时创建新事件/属性,导致 MySQL 主键冲突或重复插入

原因

分布式环境下缺少协调机制,多实例同时检测到"不存在"并尝试创建

目标

避免重复创建,保证元数据唯一性

方案

KVRocks 分布式锁 + Double Check

方案流程

步骤

实例 A(先获取锁)

实例 B(后获取锁)

1

查询缓存:不存在

查询缓存:不存在

2

获取分布式锁 ✓

获取锁失败,等待...

3

Double Check:仍不存在

-

4

写入 MySQL + 同步 KVRocks

-

5

释放锁

获取锁成功

6

-

Double Check:已存在 ✓

7

-

直接返回,无需创建

分布式锁配置

  • lockKey = "dlock:event:{appId}_{eventName}"
  • lockValue = UUID(用于安全释放)
  • expireMs = 30000(防止死锁)

关键代码

// LockManager.java - 分布式锁管理器
public <T> CompletableFuture<T> executeWithLockAsync(
        String lockKey, long waitMs, long expireMs,
        Supplier<CompletableFuture<T>> supplier) {
    
    DistributedLock lock = new DistributedLock(kvrocksClient, lockKey);
    
    return lock.tryLockAsync(waitMs, expireMs)
            .thenCompose(acquired -> {
                if (!acquired) {
                    LOG.warn("获取锁超时: key={}", lockKey);
                    return CompletableFuture.completedFuture(null);
                }
                
                return supplier.get()
                        .whenComplete((result, ex) -> {
                            // 无论成功失败都释放锁
                            lock.unlockAsync();
                        });
            });
}

// DistributedLock.java - 锁释放使用 Lua 脚本保证原子性
String luaScript = 
    "if redis.call('get', KEYS[1]) == ARGV[1] then " +
    "    return redis.call('del', KEYS[1]) " +
    "else " +
    "    return 0 " +
    "end";

效果

  • 消除重复创建问题
  • 锁等待时间 P99 < 100ms
  • MySQL 写入冲突率从 5% 降至 0%

3.3 覆盖更新解决方案

维度

内容

问题

设备表/用户表直接覆盖写入,NULL 值冲掉原有非空数据

场景

用户首次访问设置了 city=北京,后续访问 city=NULL,原数据被覆盖丢失

目标

只更新本次有值的字段,NULL 字段保留原值

方案

Doris Unique Key + 部分列更新 + 动态列签名 + 按列签名分批

版本限制与解决思路

Doris 版本

部分列更新支持

限制

3.0+

灵活列更新

每行可以更新不同的列,无需分批

2.1.x (当前使用)

固定列更新

同一批次必须更新相同的列集合

由于当前使用 Doris 2.1.3,同一个 Stream Load 请求中所有行必须更新相同的列。这意味着:

  • 如果 Row1 更新 {id, field1},Row2 更新 {id, field1, field2},不能放在同一批
  • 需要按"列签名"分组,相同列集合的数据放同一批次

解决方案:动态提取非空列 + 按列签名分批 + 多批次并行提交

方案流程

示例场景:
● 原始数据:device_id=1, field1="A", field2="B", field3="C"
● 本次更新:device_id=1, field1="X", field2=NULL, field3=NULL
● 期望结果:field1="X",field2、field3 保留原值

处理步骤:
1. 提取非空字段:getNonNullColumns() → {device_id, field1}
2. 生成列签名:columnSignature = "device_id,field1"
3. Stream Load:Header: partial_columns=true, columns=device_id,field1
4. 结果:field1 更新为 "X",field2、field3 保留原值 ✓

关键代码

// PartialUpdateRow.java - 部分列更新接口
public interface PartialUpdateRow {
    Set<String> getNonNullColumns();  // 返回所有非空列名
    Map<String, Object> toNonNullMap();  // 转换为非空字段 Map
}

// BatchBuffer.java - 按列签名分组 (解决 Doris 2.1.x 限制的核心)
public Batch addRecord(PartialUpdateRow row) {
    // Step 1: 生成列签名 (排序保证一致性)
    String columnSignature = row.getNonNullColumns().stream()
            .sorted()
            .collect(Collectors.joining(","));
    
    // Step 2: 相同表 + 相同列签名 → 同一批次
    // 例如: b_device_1@device_id,field1 和 b_device_1@device_id,field1,field2 是不同批次
    String batchKey = table + "@" + columnSignature;
    Batch batch = batches.computeIfAbsent(batchKey, 
            k -> new Batch(database, table, columnSignature));
    
    batch.addRecord(serializeToBytes(row.toNonNullMap()));
    return batch;
}

// flush 时按批次分别提交
public void flush() {
    for (Batch batch : batches.values()) {
        // 每个批次独立 Stream Load,指定各自的 columns
        dorisClient.streamLoad(batch.getTable(), batch.getColumnSignature(), batch.getData());
    }
}

Doris 2.1.x vs 3.0 对比

版本

处理方式

示例

Doris 2.1.x

按列签名分批

Row1(id,f1) + Row3(id,f1) → Batch1;Row2(id,f1,f2) → Batch2

Doris 3.0+

灵活列更新

[Row1, Row2, Row3] → 单次 Stream Load

效果

  • NULL 覆盖问题完全解决
  • 数据完整性从 85% 提升至 100%
  • 网络传输量降低 30%+(只传非空字段)
  • 虽然分批增加了请求数,但通过并行提交保证性能

3.4 性能过低综合解决方案

维度

内容

问题

部分逻辑逐条与外部系统交互,吞吐量上不去,延迟高

原因

每条数据都触发网络 IO,序列化/反序列化频繁,对象创建过多

目标

吞吐量提升 4 倍+,延迟降低 97%+

方案

批量化 + 去重 + 算子合并 + 序列化优化

优化手段矩阵

优化类型

具体措施

效果

批量查询

N 次网络往返 → 1 次 HMGET

延迟降低 80%+

批量去重

按 cuid/IdPair 去重后批量处理

查询量降低 70%+

算子合并

富化三合一、ID三合一

序列化开销降低 40%

JSON 优化

FastJSON2 + byte[] 直接解析

解析性能提升 30%+

对象复用

避免频繁创建临时对象

GC 频率降低 80%

字符级操作

正则 → substring,多次 replace → 单次遍历

CPU 降低 40%

内存映射

IP 库使用 MappedByteBuffer + 单例共享

内存降低 90%

批量去重示例

// 一条消息包含 10 个事件,其中 cuid 分布: A×6, B×3, C×1

// 优化前: 10 次 KVRocks 查询
for (event : events) {
    userId = kvrocksClient.get(event.cuid);  // 10 次网络往返
}

// 优化后: 先去重,只需 3 次查询
Map<String, List<Event>> cuidToEvents = events.stream()
        .collect(Collectors.groupingBy(e -> e.cuid));

// 去重后只有 3 个不同的 cuid: A, B, C
for (Entry<cuid, eventList> : cuidToEvents.entrySet()) {
    userId = kvrocksClient.get(cuid);  // 只需 3 次
    for (event : eventList) {
        event.userId = userId;  // 回填到所有对应事件
    }
}

算子合并示例

对比

优化前

优化后

算子数

3 个(IP → UA → 关键词)

1 个(EnrichOperator)

JSON 解析

3 次

1 次

JSON 序列化

3 次

1 次

开销

-

减少 40%

效果汇总

指标

优化前

优化后

提升

QPS(单并行度)

1000

4000+

4x

端到端延迟

5分钟+

< 10s

-97%

CPU 使用率

90%

50%

-44%

GC 频率

10+/s

1-2/s

-85%

内存占用

8GB

2GB

-75%

3.5 数据倾斜解决方案

维度

内容

问题

Gate Job 按 appId 分区发送下游 Topic,大客户数据集中在少数分区

原因

头部 App 数据量占比高(Top 10 客户可能占 80% 数据),导致 ID Job 部分实例负载过高

目标

数据均匀分布,同时保证同一设备的事件顺序

方案

使用 appId + deviceId 组合作为 Kafka 分区 key

问题示意

分区方式

数据分布

问题

按 appId

Partition 0: App A (80%)

热点分区,ID Job 实例 0 过载

Partition 1: App B (10%)

其他实例空闲

Partition 2: App C (10%)

资源利用率低

解决方案

分区 key = appId + "_" + deviceId
  • 均匀分布:设备 ID 天然分散,数据均匀分布到各分区
  • 顺序保证:同一设备的所有事件落在同一分区,保证处理顺序
  • 兼容性:下游 ID Job 无需修改,透明升级

关键代码

// GateJob 发送下游 Kafka 时设置分区 key
String partitionKey = appId + "_" + deviceId;
ProducerRecord<String, String> record = 
    new ProducerRecord<>(topic, partitionKey, jsonData);

与分布式雪花算法配合

数据重新分区后,同一 App 的数据会分散到多个 Flink 实例处理。此时需要保证生成的设备 ID、用户 ID 全局唯一:

  • 每个 Flink TaskManager 启动时,通过 KVRocks 原子自增获取唯一 workerId
  • 雪花算法结合 workerId 生成 64 位唯一 ID
  • 即使同一 App 的数据被多个实例并行处理,ID 也不会冲突

效果

指标

优化前

优化后

分区数据量

极不均匀(1:8:1)

基本均匀

热点分区延迟

分钟级

< 10s

资源利用率

~30%

90%+

事件顺序

✓ 保证

✓ 保证(同设备同分区)


四、技术实现细节

4.1 SSDB → KVRocks + 两级缓存

背景:SSDB 已停止维护,且单机性能有限

方案:KVRocks (Redis 兼容) + Caffeine 两级缓存

层级

组件

特点

L1

Caffeine

毫秒级响应、异步加载、refreshAfterWrite、支持负缓存

L2

KVRocks

Redis 兼容、跨节点共享、数据持久化、支持集群

关键代码实现

// ConfigCacheService.java - 两级缓存构建
private <V> AsyncLoadingCache<String, CacheValue<V>> buildHashCacheWithNegative(
        int maxSize, int refreshMinutes, int expireMinutes,
        String kvrocksKey, Function<String, V> parser) {

    String actualKey = getActualKey(kvrocksKey);

    return Caffeine.newBuilder()
            .maximumSize(maxSize)                              // 容量限制
            .refreshAfterWrite(refreshMinutes, TimeUnit.MINUTES)  // 异步刷新(不阻塞)
            .expireAfterWrite(expireMinutes, TimeUnit.MINUTES)    // 过期时间
            .executor(cacheExecutor)                           // 异步执行器
            .recordStats()                                     // 统计信息(命中率等)
            .buildAsync((field, executor) ->
                    kvrocksClient.asyncHGet(actualKey, field)
                            .thenApply(value -> {
                                if (value != null) {
                                    V parsed = parser.apply(value);
                                    return CacheValue.of(parsed);
                                }
                                return CacheValue.notFound();  // 负缓存
                            })
            );
}

4.2 负缓存机制:防止缓存穿透

问题:如果某个 key 不存在,每次查询都会穿透到 KVRocks,造成大量无效请求。

解决方案:缓存"不存在"状态本身。

// 包装类,区分"找到值"和"确认不存在"
private static class CacheValue<T> {
    private final T value;

    static <T> CacheValue<T> of(T value) {
        return new CacheValue<>(value);  // 找到了值
    }

    static <T> CacheValue<T> notFound() {
        return new CacheValue<>(null);   // 确认不存在(负缓存)
    }

    boolean isPresent() {
        return value != null;
    }
}

使用方式

public CompletableFuture<Integer> getEventId(...) {
    return eventIdCache.get(field)
        .thenApply(cv -> {
            if (cv != null && cv.isPresent()) {
                return cv.getValue();  // 命中正缓存
            }
            return null;  // 命中负缓存,直接返回,不再穿透
        });
}

效果对比

场景

无负缓存

有负缓存

key 存在

缓存 value ✓

缓存 value ✓

key 不存在

每次穿透 ✗

缓存 notFound ✓

防穿透攻击

4.3 批量查询优化:N 次网络往返 → 1 次

优化前:N 个属性,N 次网络往返

for (String propId : propIds) {
    String column = cache.get(field).get();  // 每次都走网络
}

优化后:先查本地,再批量查远程

public CompletableFuture<Map<String, Integer>> batchGetEventAttrColumnIndex(
        String eventId, List<String> propIds) {

    Map<String, Integer> result = new ConcurrentHashMap<>();
    List<String> missedPropIds = new ArrayList<>();
    List<String> missedFields = new ArrayList<>();

    // Step 1: 批量检查 L1 本地缓存
    for (String propId : propIds) {
        String field = CacheKeyConstants.eventAttrColumnField(eventId, propId);
        CacheValue<String> cached = cache.synchronous().getIfPresent(field);
        
        if (cached != null && cached.isPresent()) {
            // L1 缓存命中
            result.put(propId, parseColumnIndex(cached.getValue()));
        } else if (cached != null && !cached.isPresent()) {
            // 负缓存命中(确认不存在),跳过
        } else {
            // 缓存未命中,需要查询 KVRocks
            missedPropIds.add(propId);
            missedFields.add(field);
        }
    }

    // Step 2: 如果全部命中,直接返回
    if (missedPropIds.isEmpty()) {
        return CompletableFuture.completedFuture(result);
    }

    // Step 3: 批量查询 KVRocks(只有 1 次网络往返)
    return kvrocksClient.asyncHMGet(actualKey, missedFields)
            .thenApply(values -> {
                // Step 4: 处理结果,更新缓存
                for (int i = 0; i < missedFields.size(); i++) {
                    String field = missedFields.get(i);
                    String propId = missedPropIds.get(i);
                    String column = values.get(i);

                    if (column != null) {
                        // 找到了,更新 L1 缓存
                        cache.synchronous().put(field, CacheValue.of(column));
                        result.put(propId, parseColumnIndex(column));
                    } else {
                        // 未找到,设置负缓存
                        cache.synchronous().put(field, CacheValue.notFound());
                    }
                }
                return result;
            });
}

效果:N 次网络往返 → 1 次,延迟降低 80%+

4.4 批量去重优化:减少 70%+ 重复 ID 查询

背景:一条消息的 data[] 数组中可能有多个事件,这些事件可能使用相同的 cuid 或 (deviceId, userId) 组合。

问题:如果逐条处理,同一个 cuid 会重复查询 KVRocks N 次

方案:先按 key 去重收集,再批量处理,最后回填结果

// OneIdAsyncOperator.java - 批量获取用户ID(按 cuid 去重)

private CompletableFuture<Void> processUserIds(IdProcessContext ctx) {
    
    // ========== Step 1: 收集所有 cuid 并去重 ==========
    // 记录 cuid → 对应的所有 pr 对象(一个 cuid 可能对应多个事件)
    Map<String, List<Map<String, Object>>> cuidToPrMap = new LinkedHashMap<>();

    for (Object itemObj : ctx.dataList) {
        JSONObject item = (JSONObject) itemObj;
        JSONObject pr = item.getJSONObject("pr");
        String cuid = pr.getString("$cuid");
        
        if (StringUtils.isNotBlank(cuid)) {
            // 相同 cuid 的多个 pr 对象收集到一起
            cuidToPrMap.computeIfAbsent(cuid, k -> new ArrayList<>()).add(pr);
        }
    }

    // ========== Step 2: 去重后批量异步查询 ==========
    // N 个事件可能只需 M 次查询,M << N
    List<CompletableFuture<Void>> futures = new ArrayList<>(cuidToPrMap.size());

    for (Map.Entry<String, List<Map<String, Object>>> entry : cuidToPrMap.entrySet()) {
        String cuid = entry.getKey();
        List<Map<String, Object>> prList = entry.getValue();

        CompletableFuture<Void> future = oneIdService.getOrCreateUserId(ctx.appId, cuid, uidMetrics)
                .thenAccept(result -> {
                    if (result != null) {
                        Long zgUserId = result.getId();
                        
                        // ========== Step 3: 回填到所有对应的 pr 对象 ==========
                        for (Map<String, Object> pr : prList) {
                            pr.put("$zg_uid", zgUserId);
                        }
                        
                        // 新用户归档
                        if (result.getIsNew()) {
                            archiveKafkaService.sendToKafka(ArchiveType.USER, ctx.appId, cuid, zgUserId);
                        }
                    }
                });
        futures.add(future);
    }

    // 等待所有查询完成
    return CompletableFuture.allOf(futures.toArray(new CompletableFuture[0]));
}

ZgId 同样按 (deviceId, userId) 组合去重

// 收集所有 (zgDid, zgUid) 组合并去重
Map<IdPair, List<Map<String, Object>>> idPairToPrMap = new LinkedHashMap<>();

for (Object item : ctx.dataList) {
    JSONObject pr = ((JSONObject) item).getJSONObject("pr");
    Long zgDeviceId = pr.getLong("$zg_did");
    Long zgUserId = pr.getLong("$zg_uid");  // 可为 null
    
    IdPair idPair = new IdPair(zgDeviceId, zgUserId);
    idPairToPrMap.computeIfAbsent(idPair, k -> new ArrayList<>()).add(pr);
}

// IdPair 类,用于去重
private static class IdPair {
    final Long zgDeviceId;
    final Long zgUserId;  // 可为 null
    
    // hashCode 和 equals 实现...
}

效果示例

假设一条消息有 10 个事件,cuid 分布:A×6、B×3、C×1

对比

查询次数

说明

优化前

10 次

每个事件都查询一次

优化后

3 次

去重后只查 A、B、C

减少

70%

-

4.5 分布式 ID 生成:雪花算法 + KVRocks WorkerId

问题:Flink 多 TaskManager 场景下,如何保证雪花 ID 的 workerId 唯一?

方案:通过 KVRocks 原子自增生成全局唯一 workerId

// KvrocksClient.java - 原子自增
public Long incr(String key) {
    if (isCluster) {
        return clusterConnection.sync().incr(key);
    } else {
        return standaloneConnection.sync().incr(key);
    }
}

// OneIdService.java - 生成唯一 workerId
private int generateWorkerId() {
    // 通过 KVRocks 原子自增获取全局唯一 ID
    Long workerId = kvrocksClient.incr("snowflake:worker_id");
    return (int) (workerId % MAX_WORKER_ID);  // 取模保证范围 [0, 1023]
}

雪花算法实现(支持时钟回拨)

// SnowflakeIdGenerator.java
public class SnowflakeIdGenerator {
    
    // 64位结构: 1位符号位 + 41位时间戳 + 10位workerId + 12位序列号
    private static final long EPOCH = 1704067200000L;  // 2024-01-01
    private static final long WORKER_ID_BITS = 10L;    // 支持 1024 个节点
    private static final long SEQUENCE_BITS = 12L;     // 每毫秒 4096 个 ID
    
    // 时钟回拨配置
    private static final long SMALL_CLOCK_BACKWARDS_MS = 5L;   // 小回拨阈值
    private static final long MAX_CLOCK_BACKWARDS_MS = 100L;   // 最大容忍回拨
    
    public synchronized long nextId() {
        long timestamp = System.currentTimeMillis();

        // ========== 时钟回拨检测和处理 ==========
        if (timestamp < lastTimestamp) {
            long offset = lastTimestamp - timestamp;
            
            if (offset <= SMALL_CLOCK_BACKWARDS_MS) {
                // 小回拨(≤5ms):静默等待
                timestamp = waitForClockRecovery(lastTimestamp, offset);
            } else if (offset <= MAX_CLOCK_BACKWARDS_MS) {
                // 中等回拨(5-100ms):记录警告并等待
                LOG.warn("时钟回拨 {}ms,等待恢复", offset);
                timestamp = waitForClockRecovery(lastTimestamp, offset);
            } else {
                // 大回拨(>100ms):抛出异常,拒绝生成
                throw new RuntimeException("时钟回拨过大: " + offset + "ms");
            }
        }

        // 同一毫秒内序列号自增
        if (timestamp == lastTimestamp) {
            sequence = (sequence + 1) & MAX_SEQUENCE;
            if (sequence == 0) {
                // 序列号用完,等待下一毫秒
                timestamp = waitNextMillis(lastTimestamp);
            }
        } else {
            sequence = 0L;
        }

        lastTimestamp = timestamp;

        // 组装 64 位 ID
        return ((timestamp - EPOCH) << TIMESTAMP_SHIFT)
                | (workerId << WORKER_ID_SHIFT)
                | sequence;
    }
    
    private long waitForClockRecovery(long lastTimestamp, long offset) {
        try {
            Thread.sleep(offset);  // 等待回拨时间
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
        
        // 自旋确保恢复
        long timestamp = System.currentTimeMillis();
        while (timestamp <= lastTimestamp) {
            Thread.yield();
            timestamp = System.currentTimeMillis();
        }
        return timestamp;
    }
}

4.6 算子合并:减少 40% 序列化开销

4.6.1 富化算子三合一

优化前:三个独立算子串行处理

IpEnrichOperator → UAEnrichOperator → KeywordEnrichOperator
       ↓                 ↓                    ↓
   JSON解析          JSON解析             JSON解析
   JSON序列化        JSON序列化           JSON序列化

问题

  • 每个算子都要解析和序列化 JSON(3 次解析 + 3 次序列化)
  • 算子间数据传递有序列化开销
  • Flink 调度开销增加

优化后:合并为 EnrichOperator

/**
 * 统一富化算子 - 合并 IP、UA、搜索关键词三个功能
 * 
 * 优化点:
 * 1. JSON 只解析一次,只序列化一次
 * 2. 减少算子间数据传递的序列化开销
 * 3. IP 和 UA 是同步操作,关键词是异步操作
 */
public class EnrichOperator extends RichAsyncFunction<ZGMessage, ZGMessage> {

    private transient IpDatabaseLoader ipLoader;
    private transient UserAgentParser uaParser;
    private transient SearchKeywordParser keywordParser;
    private transient KvrocksClient kvrocksClient;
    
    // 本地缓存
    private transient Cache<String, UserAgentInfo> uaCache;
    private transient Cache<String, BaiduKeyword> keywordCache;

    @Override
    public void asyncInvoke(ZGMessage input, ResultFuture<ZGMessage> resultFuture) {
        metrics.in();
        
        // ========== 1. JSON 只解析一次 ==========
        JSONObject json = JSON.parseObject(input.getRawData());
        
        // ========== 2. IP 富化(同步,本地 IP 库)==========
        String ip = json.getString("ip");
        String[] ipResult = ipLoader.lookup(ip);
        json.put("country", ipResult[0]);
        json.put("province", ipResult[1]);
        json.put("city", ipResult[2]);
        
        // ========== 3. UA 富化(同步,本地解析 + 缓存)==========
        String ua = json.getString("ua");
        UserAgentInfo uaInfo = uaCache.get(ua, k -> uaParser.parse(k));
        json.put("os", uaInfo.getOs());
        json.put("browser", uaInfo.getBrowser());
        
        // ========== 4. 关键词富化(异步,需要查 KVRocks)==========
        String ref = json.getString("$ref");
        enrichKeywordAsync(json, ref)
            .whenComplete((v, ex) -> {
                // ========== 5. JSON 只序列化一次 ==========
                input.setRawData(json.toJSONString());
                metrics.out();
                resultFuture.complete(Collections.singleton(input));
            });
    }
}
4.6.2 ID 处理算子合并

优化前:三个独立的异步算子

DeviceIdAsyncOperator → UserIdAsyncOperator → ZgidAsyncOperator

优化后:合并为 OneIdAsyncOperator,使用链式 CompletableFuture

/**
 * 统一ID异步算子 - 合并 DeviceId、UserId、Zgid 三个算子
 * 
 * 优化点:
 * 1. 合并三个算子为一个,减少数据流转和序列化开销
 * 2. 使用链式 CompletableFuture,保证执行顺序
 * 3. UserId 按 cuid 去重后批量处理
 * 4. ZgId 按 (did, uid) 组合去重后批量处理
 */
public class OneIdAsyncOperator extends RichAsyncFunction<ZGMessage, ZGMessage> {

    @Override
    public void asyncInvoke(ZGMessage input, ResultFuture<ZGMessage> resultFuture) {
        metrics.in();
        
        JSONObject rawDataJson = JSON.parseObject(input.getRawData());
        IdProcessContext ctx = new IdProcessContext(input, appId, did, usr, dataList);

        // 链式异步处理: DeviceId → UserId → ZgId
        processDeviceId(ctx)
                .thenCompose(v -> processUserIds(ctx))    // UserId 依赖 DeviceId
                .thenCompose(v -> processZgIds(ctx))      // ZgId 依赖 DeviceId + UserId
                .whenComplete((v, throwable) -> {
                    if (throwable != null) {
                        LOG.error("处理失败: appId={}, did={}", appId, did, throwable);
                        metrics.error();
                    } else {
                        // 更新 rawData
                        input.setRawData(rawDataJson.toJSONString());
                        metrics.out();
                    }
                    resultFuture.complete(Collections.singleton(input));
                });
    }
}

效果

  • 减少 40% 序列化开销
  • 降低 Flink 调度开销
  • 代码更易维护

4.7 序列化优化

4.7.1 FastJSON → FastJSON2
<!-- 优化后 -->
<dependency>
    <groupId>com.alibaba.fastjson2</groupId>
    <artifactId>fastjson2</artifactId>
    <version>2.0.43</version>
</dependency>

效果:性能提升 30%+,安全性更好

4.7.2 直接从 byte[] 解析

优化前

String rawData = new String(record.value(), StandardCharsets.UTF_8);  // byte[] → String
JSONObject json = JSON.parseObject(rawData);  // String → JSON

优化后

// 直接从 byte[] 解析,跳过 String 转换
JSONReader.Feature[] READER_FEATURES = { JSONReader.Feature.IgnoreAutoTypeNotMatch };
JSONObject json = JSON.parseObject(valueBytes, READER_FEATURES);

// 只在需要时才转 String(下游确实需要 rawData 字符串)
if (json != null) {
    message.setRawData(new String(valueBytes, StandardCharsets.UTF_8));
}
4.7.3 BatchBuffer 零拷贝

优化前

String json = JSON.toJSONString(map);  // Map → String
byte[] bytes = json.getBytes(UTF_8);   // String → byte[]

优化后

// 直接 Map → byte[]
private byte[] serializeToBytes(Map<String, Object> map) {
    try (JSONWriter writer = JSONWriter.ofUTF8()) {
        writer.writeAny(map);
        return writer.getBytes();  // 零拷贝
    }
}
4.7.4 EventAttrRow 优化

List → Array

// 优化前
private final List<String> columns = new ArrayList<>(TOTAL_COLUMNS);

// 优化后
private final String[] columns = new String[TOTAL_COLUMNS];

JSON 转义字符级操作

// 优化前:每次 replace 创建新字符串
private String escapeJson(String s) {
    return s.replace("\\", "\\\\")
            .replace("\"", "\\\"")
            .replace("\n", "\\n");
}

// 优化后:直接操作 char,零中间对象
private void appendEscaped(StringBuilder sb, String value) {
    for (int i = 0, len = value.length(); i < len; i++) {
        char c = value.charAt(i);
        switch (c) {
            case '\\': sb.append("\\\\"); break;
            case '"':  sb.append("\\\""); break;
            case '\n': sb.append("\\n"); break;
            case '\r': sb.append("\\r"); break;
            case '\t': sb.append("\\t"); break;
            default:   sb.append(c);
        }
    }
}

快速 NULL 检查

// 优化前
if (NULL_VALUE.equals(value)) { ... }

// 优化后:先检查长度,快速失败
private static boolean isNullValue(String value) {
    if (value == null || value.isEmpty()) return true;
    // NULL_VALUE = "\\N",长度为 2
    if (value.length() == 2 && value.charAt(0) == '\\' && value.charAt(1) == 'N') 
        return true;
    return false;
}

4.8 内存优化:IP 库内存映射

问题:JDK 工具分析显示 IP 库加载占用大量堆内存,每个 slot 都加载一份,导致 OOM

方案

  1. 使用 MappedByteBuffer 内存映射,数据不占用 JVM 堆内存
  2. 单例 + 引用计数,多 slot 共享
/**
 * IP数据库加载器 - 内存映射版本 (单例 + 引用计数)
 *
 * 核心优化:
 * - 单例模式:同一 JVM 内所有 Flink slot/task 共享一个实例
 * - 引用计数:只有最后一个使用者关闭时才清理资源
 * - MappedByteBuffer:数据不占用 JVM 堆内存
 * - 文件锁:解决多进程并发下载问题
 */
public class IpDatabaseLoader implements Closeable {

    private static volatile IpDatabaseLoader INSTANCE;
    private static final Object LOCK = new Object();
    private static final AtomicInteger REF_COUNT = new AtomicInteger(0);

    // IPv4 数据库(内存映射版本)
    private final AtomicReference<IpDatabaseMapped> ipv4DatabaseRef = new AtomicReference<>();

    public static IpDatabaseLoader getOrCreate(Builder builder) throws Exception {
        if (INSTANCE == null) {
            synchronized (LOCK) {
                if (INSTANCE == null) {
                    INSTANCE = builder.build();
                    INSTANCE.init();
                    LOG.info("IP数据库单例实例已创建");
                }
            }
        }
        
        int count = REF_COUNT.incrementAndGet();
        LOG.info("IP数据库引用计数增加: {}", count);
        return INSTANCE;
    }

    @Override
    public void close() {
        int count = REF_COUNT.decrementAndGet();
        LOG.info("IP数据库引用计数减少: {}", count);
        
        if (count <= 0) {
            cleanup();
            INSTANCE = null;
        }
    }
}

效果:IP 库内存占用降低 90%+

4.9 定时全量刷新:避免冷启动

private void startScheduledRefresh() {
    scheduler = Executors.newSingleThreadScheduledExecutor(
            r -> new Thread(r, "cache-refresh-scheduler"));
    
    scheduler.scheduleAtFixedRate(
            this::refreshAllCachesParallel,
            refreshConfig.getRefreshIntervalMinutes(),
            refreshConfig.getRefreshIntervalMinutes(),
            TimeUnit.MINUTES
    );
}

// 多线程并行刷新
private void refreshAllCachesParallel() {
    List<RefreshTask> tasks = Arrays.asList(
        new RefreshTask("eventId", () -> refreshHashCache(EVENT_ID_MAP, eventIdCache)),
        new RefreshTask("eventAttrId", () -> refreshHashCache(EVENT_ATTR_ID_MAP, eventAttrIdCache)),
        new RefreshTask("eventAttrColumn", () -> refreshHashCache(EVENT_ATTR_COLUMN_MAP, eventAttrColumnCache)),
        // ... 更多缓存
    );

    // 并行执行
    List<CompletableFuture<RefreshResult>> futures = tasks.stream()
        .map(task -> CompletableFuture.supplyAsync(task::execute, refreshPool))
        .collect(Collectors.toList());

    // 等待所有完成
    CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();
}

五、性能调优实战案例

5.1 案例一:对象复用 - QPS 4倍提升

问题现象:系统 QPS 始终卡在 1000 左右,加大并行度、增加资源都没用。

定位过程

# Step 1: 查看 GC 日志
-XX:+PrintGCDetails -XX:+PrintGCDateStamps -Xloggc:/tmp/gc.log

# 发现 Young GC 每秒执行 10+ 次,明显异常

根因分析:Young GC 频繁 = 短命对象太多 = JSON 解析时每次都创建大量临时对象

优化方案:升级 FastJSON 到 FastJSON2,启用对象复用

JSONReader.Feature[] features = {
    JSONReader.Feature.UseNativeObject
};
JSONObject json = JSON.parseObject(valueBytes, features);

效果

指标

优化前

优化后

QPS

1000

4000+

Young GC

10+/s

1-2/s

提升幅度

-

4倍

5.2 案例二:缓存命中率指标定位 Bug

问题现象:上线新版本后,KVRocks 的 QPS 暴涨,系统延迟飙升。

定位过程

# 通过 Flink Metrics 查看自定义指标
cache_hit_rate = 0%   ← 缓存命中率竟然是 0!
kvrocks_qps = 暴涨    ← 全部穿透到 KVRocks

根因:搜索代码 invalidateAll,发现版本检查代码有 bug

// 错误代码:版本号不变时也会清空!
public void checkVersion() {
    long newVersion = getVersion();
    if (newVersion >= currentVersion) {  // 应该是 >,不是 >=
        invalidateAllL1();  // 每次都清空
        currentVersion = newVersion;
    }
}

// 修复后
if (newVersion > currentVersion) {
    invalidateAllL1();
    currentVersion = newVersion;
}

效果:缓存命中率恢复到 85%+,KVRocks QPS 降低 80%

教训:监控指标是发现问题的眼睛,没有指标就是盲人摸象

5.3 案例三:JDK 工具定位内存溢出

问题现象:Flink TaskManager 频繁 OOM

定位过程

# 启动 Flink 作业时添加 JVM 参数,OOM 时自动生成 heap dump
flink run \
    -yD env.java.opts.taskmanager="-XX:+HeapDumpOnOutOfMemoryError -XX:HeapDumpPath=/tmp" \
    ...

# TaskManager OOM 后,在 /tmp 目录下会生成 heap dump 文件
# 文件名格式:java_pid<进程号>.hprof

# 使用 JDK 自带工具 jvisualvm 分析
jvisualvm --openfile /tmp/java_pid12345.hprof

jvisualvm 分析结果

查看 "类" 视图,按实例大小排序:
- IpDatabase: 800MB  ← 罪魁祸首!
- IpDatabase: 800MB  ← 又一份?
- IpDatabase: 800MB  ← 每个 slot 都有一份!

根因:IP 库文件每个 Flink slot 都加载一份到堆内存

解决方案:MappedByteBuffer + 单例共享(见 4.8 节)

效果:IP 库内存占用降低 90%+,OOM 问题解决

5.4 案例四:火焰图定位 CPU 热点

问题现象:CPU 使用率 90%+,但吞吐量上不去

定位过程

# 启动 Flink 作业时启用火焰图功能
flink run -yD rest.flamegraph.enabled=true ...

# 在 Flink Web UI 上直接查看火焰图
# 路径:Flink Web UI → 选择 Job → 选择 Task → 点击 "Flame Graph" 标签

火焰图解读

  • 横轴:函数 CPU 占比(越宽占用越高)
  • 纵轴:函数调用堆栈(从下到上是调用链)

发现两个大平顶(宽条):

热点函数

CPU 占比

问题

Pattern.matcher()

35%

正则表达式重复编译

String.replace()

20%

字符串多次替换

优化方案

// 问题1:字符串截断用了正则
// 优化前
public String ensureLength(String s, int max) {
    return s.replaceAll("(.{" + max + "}).*", "$1");
}

// 优化后:直接使用 substring
public String ensureLength(String s, int max) {
    return s.length() <= max ? s : s.substring(0, max);
}

// 问题2:JSON 转义用了多次 replace
// 优化前
s.replace("\\", "\\\\").replace("\"", "\\\"").replace("\n", "\\n");

// 优化后:字符级操作(见 4.7.4)

效果:CPU 使用率从 90% 降到 50%,吞吐量提升 40%


六、数据一致性保障

6.1 Flink Checkpoint + Doris 2PC

// Checkpoint 配置
env.enableCheckpointing(60000);  // 1 分钟
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
env.getCheckpointConfig().setMinPauseBetweenCheckpoints(30000);
env.getCheckpointConfig().setCheckpointTimeout(120000);

// Doris Sink 使用 2PC
DorisExecutionOptions.builder()
    .enable2PC()  // 两阶段提交
    .build();

6.2 故障恢复:状态保存 + 幂等重试

// PartialUpdateSinkWriter.java
private void recoverFromState(Collection<WriterState> states) {
    for (WriterState state : states) {
        LOG.info("恢复状态: {}", state);
        
        for (WriterState.BatchSnapshot snapshot : state.getPendingBatches()) {
            // 生成新 label 重试
            String label = generateLabel(snapshot.database, snapshot.table);
            
            DorisHttpClient.LoadResult result = client.streamLoad(
                    snapshot.database,
                    snapshot.table,
                    snapshot.columnSignature,
                    snapshot.jsonData,
                    label);
            
            // Label 已存在说明之前提交成功,幂等
            if (result.success || result.isLabelExists()) {
                LOG.info("恢复批次提交成功: {}", snapshot);
                totalRowsWritten += snapshot.recordCount;
            } else {
                LOG.error("恢复批次提交失败: {}, error={}", snapshot, result.message);
            }
        }
    }
}

七、监控指标体系

7.1 自定义业务指标

public class OperatorMetrics {
    
    private Counter inputCounter;
    private Counter outputCounter;
    private Counter errorCounter;
    private Counter skipCounter;
    private Counter cacheHitCounter;
    private Counter cacheMissCounter;
    
    public void register(MetricGroup metricGroup) {
        inputCounter = metricGroup.counter("input");
        outputCounter = metricGroup.counter("output");
        errorCounter = metricGroup.counter("error");
        skipCounter = metricGroup.counter("skip");
        cacheHitCounter = metricGroup.counter("cache_hit");
        cacheMissCounter = metricGroup.counter("cache_miss");
        
        // 缓存命中率
        metricGroup.gauge("cache_hit_rate", () -> {
            long total = cacheHitCounter.getCount() + cacheMissCounter.getCount();
            return total > 0 ? (double) cacheHitCounter.getCount() / total : 0;
        });
    }
    
    public void in() { inputCounter.inc(); }
    public void out() { outputCounter.inc(); }
    public void error() { errorCounter.inc(); }
    public void skip() { skipCounter.inc(); }
    public void cacheHit() { cacheHitCounter.inc(); }
    public void cacheMiss() { cacheMissCounter.inc(); }
}

7.2 关键指标告警

指标

含义

告警阈值

cache_hit_rate

缓存命中率

< 80%

error_rate

错误率

> 1%

latency_p99

P99 延迟

> 1s

kvrocks_qps

KVRocks 请求量

突增 50%+

7.3 监控方案对比

特性

Flink Web UI

Prometheus + Grafana

实时查看

✅ 支持

✅ 支持

历史数据

❌ 不支持

✅ 支持

趋势图表

❌ 不支持

✅ 丰富图表

告警

❌ 不支持

✅ 支持告警规则

聚合查询

❌ 只能看单个Task

✅ 可跨Job聚合

配置难度

✅ 开箱即用

⚠️ 需要配置


八、优化效果总结

优化项

效果

两级缓存 + 负缓存

减少 80%+ KVRocks 访问

批量查询

网络往返 N → 1

批量去重

减少 70%+ 重复 ID 查询

算子合并

减少 40% 序列化开销

FastJSON2 + byte[] 直接解析

性能提升 30%+

对象复用

QPS 1000 → 4000+

IP 库内存映射

内存占用降低 90%

正则优化 + 字符级操作

CPU 降低 40%

部分列更新

解决 NULL 覆盖问题

雪花 ID + KVRocks WorkerId

分布式 ID 唯一性保证


九、经验与教训

9.1 性能优化方法论

  1. 先测量,后优化:没有指标就是盲人摸象(火焰图、GC日志、Metrics)
  2. 找热点:定位 Top N 问题,20% 代码造成 80% 性能问题
  3. 逐个击破:优化 → 验证 → 上线,一次只改一个变量

9.2 核心优化思路

思路

具体措施

缓存是王道

热点数据本地缓存、负缓存防穿透、批量查询减少网络往返

减少序列化

合并算子、复用对象、byte[] 直接操作

异步化

CompletableFuture 链式处理、异步 I/O 不阻塞

批量化

聚合后批量写入、批量去重、减少网络往返

去重

先去重再处理,减少重复计算和查询

9.3 踩坑总结

教训

JSON 库性能差异大

FastJSON2 比 1 快 30%+

正则表达式很慢

高频场景用字符操作替代

缓存没监控

命中率是核心指标

大文件多份加载

单例 + 内存映射

NULL 覆盖数据

部分列更新

重复 ID 查询

先去重再批量处理

时钟回拨

雪花算法需要处理回拨


更多推荐