Flink 实时 ETL 全链路优化实战:从 Spark Streaming 到 Flink 的升级之路
本文记录了将用户行为分析平台的 ETL 系统从 Spark Streaming 迁移到 Flink 1.17 的完整历程,涵盖架构设计、技术选型、性能优化等多个维度,最终实现 10s 内端到端时效、单并行度 4000+ QPS 的高性能目标。
目录
- 项目背景与挑战
- 架构设计
- 解决方案专题
- 技术实现细节
- 性能调优实战案例
- 数据一致性保障
- 监控指标体系
- 优化效果总结
- 经验与教训
一、项目背景与挑战
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
方案:
- 使用 MappedByteBuffer 内存映射,数据不占用 JVM 堆内存
- 单例 + 引用计数,多 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 性能优化方法论
- 先测量,后优化:没有指标就是盲人摸象(火焰图、GC日志、Metrics)
- 找热点:定位 Top N 问题,20% 代码造成 80% 性能问题
- 逐个击破:优化 → 验证 → 上线,一次只改一个变量
9.2 核心优化思路
|
思路 |
具体措施 |
|
缓存是王道 |
热点数据本地缓存、负缓存防穿透、批量查询减少网络往返 |
|
减少序列化 |
合并算子、复用对象、byte[] 直接操作 |
|
异步化 |
CompletableFuture 链式处理、异步 I/O 不阻塞 |
|
批量化 |
聚合后批量写入、批量去重、减少网络往返 |
|
去重 |
先去重再处理,减少重复计算和查询 |
9.3 踩坑总结
|
坑 |
教训 |
|
JSON 库性能差异大 |
FastJSON2 比 1 快 30%+ |
|
正则表达式很慢 |
高频场景用字符操作替代 |
|
缓存没监控 |
命中率是核心指标 |
|
大文件多份加载 |
单例 + 内存映射 |
|
NULL 覆盖数据 |
部分列更新 |
|
重复 ID 查询 |
先去重再批量处理 |
|
时钟回拨 |
雪花算法需要处理回拨 |
更多推荐
所有评论(0)