实时数据处理折腾手记
实时数据处理相关的坑,多半出在边界条件上。
这次按现象往回追,不先画大图。
批处理的时代终结不了,但真不够用了
先说说我们从哪儿来。两年前,我们这套数据平台就是典型的 Lambda 架构:离线层用 Spark 每天凌晨两点跑全量任务,实时层用 Storm 跑几个简单指标,最后在服务层做合并。
架构图看着齐整,运维起来很折腾。最大的问题是那套"手动合并"的逻辑:
// 典型的 Lambda 架构合并逻辑(伪代码)
def getMetrics(start: Long, end: Long): Metrics = {
val realtime = realtimeStore.query(start, end) // 比如最近 6 小时
val batch = batchStore.query(start - 6hours, end) // 比如昨天今天
val corrected = batch.filter(_.timestamp > start) // 用批处理覆盖实时
realtime.filter(_.timestamp <= start - 6hours) ++ corrected
}
这段代码的问题不在语法,在于它假定"批处理的数据一定是准的"。但实际情况是:批处理任务可能会失败、重跑、补数据,到时候你去覆盖实时数据,就会碰到时间窗口对不齐、数据重复或丢失这种烂摊子。
更坑的是,那套 Spark 任务跑起来要两个小时左右。遇到数据倾斜或上游延迟,监控告警就会连续响,然后全组人爬起来看日志、排查问题。
第一次流计算尝试:Kafka Streams 的甜头
当时想的是:能不能先从简单场景开始,把几个关键指标改成实时计算?我们选了 Kafka Streams,原因很简单——项目里已经在用 Kafka,不想再引入一套独立集群。
代码写起来确实轻量:
// Kafka Streams 处理用户行为事件
StreamsBuilder builder = new StreamsBuilder();
KStream<String, String> events = builder.stream("user-events");
// 统计每分钟的 PV/UV
KTable<Windowed<String>, Long> pvCount = events
.groupByKey()
.windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(1)))
.count();
// UV 需要去重,这里用了一个简单的近似方案
KTable<Windowed<String>, Long> uvCount = events
.selectKey((k, v) -> extractUserId(v))
.groupByKey()
.windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(1)))
.aggregate(
() -> new MutableLong(0),
(k, v, agg) -> agg.incrementAndGet(),
Materialized.with(Serdes.String(), new MutableLongSerde())
)
.mapValues(MutableLong::get);
// 写入结果 topic
pvCount.toStream().to("pv-metrics", Produced.keySerde(Serdes.String()));
uvCount.toStream().to("uv-metrics", Produced.keySerde(Serdes.String()));
这套东西跑起来确实爽:任务延迟从小时级降到了分钟级,业务能实时看到活动效果。但很快就碰到了几个现实问题:
第一个是"迟到数据处理"。用户行为事件可能会因为网络延迟、客户端重试等原因迟到几分钟甚至更久,但我们的时间窗口早就关闭了。简单粗暴的办法是把窗口留大一点,但这对实时性没有帮助。
第二个是"状态膨胀"。Kafka Streams 默认把状态存在本地 RocksDB,随着数据量增长,磁盘占用越来越吓人。有一次重启任务,光恢复状态就花了二十分钟。
第三个是" Exactly-Once “。Kafka Streams 理论上支持 exactly-once,但前提是下游也支持事务。我们那套存储是 MySQL,结果就是:要么降级到 at-least-once,要么自己搞一套去重逻辑。
这些坑让我意识到:轻量化的代价是你要自己兜更多底。
上 Flink:从能跑跑不坏到生产级
大概在半年后,我们开始考虑把核心指标迁移到 Flink。不是 Kafka Streams 不好,而是需求已经压不住了:业务要求更多复杂聚合、窗口操作、状态管理,还要支持自定义函数和多种 sink。
第一次搭 Flink 集群的时候,我犯了个低级错误——直接用了默认配置。结果跑了一个小时的作业就 OutOfMemory,日志显示是状态后端撑爆了。
后来才明白,Flink 的状态管理不像 Kafka Streams 那样"开箱即用”,你得根据场景选择:
// Flink 状态后端选择(关键配置)
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 场景 1:小状态,低延迟 → HashMapStateBackend + JobManager
env.setStateBackend(new HashMapStateBackend());
env.getCheckpointConfig().setCheckpointStorage("file:///tmp/checkpoints");
// 场景 2:大状态,高可用 → EmbeddedRocksDBStateBackend + 分布式存储
env.setStateBackend(new EmbeddedRocksDBStateBackend());
env.getCheckpointConfig().setCheckpointStorage("hdfs:///flink/checkpoints");
// 场景 3:超大状态,需要增量 → 增量 checkpoint
env.getCheckpointConfig().enableIncrementalCheckpoints(true);
env.getCheckpointConfig().setMinPauseBetweenCheckpoints(500);
env.getCheckpointConfig().setCheckpointTimeout(60000);
我们最后选了第二种方案,状态存在 HDFS 上,配合 RocksDB 本地缓存。这样即便任务重启,状态恢复时间也控制在了五分钟以内。
另一个大头是"水位线"和"迟到数据处理"。刚开始我完全照抄官方文档,设了个固定延迟:
// 固定延迟的水位线(简单但不够灵活)
env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);
env.getConfig().setAutoWatermarkInterval(200);
SingleOutputStreamOperator<Event> withWatermarks = events
.assignTimestampsAndWatermarks(
WatermarkStrategy.<Event>forBoundedOutOfOrderness(Duration.ofSeconds(5))
.withTimestampAssigner((event, timestamp) -> event.getTimestamp())
);
跑了一段时间发现,有些迟到数据直接被丢弃了,指标准确率打了折扣。后来改成动态延迟,再加个侧输出流专门收集迟到事件:
// 动态延迟 + 侧输出流处理迟到数据
WatermarkStrategy<Event> watermarkStrategy = WatermarkStrategy
.<Event>forGenerator(ctx -> new PunctuatedWatermarkGenerator())
.withTimestampAssigner((event, timestamp) -> event.getTimestamp())
.withIdleness(Duration.ofMinutes(2)); // 2 分钟没有数据就推进水位线
SingleOutputStreamOperator<Event> mainStream = events
.assignTimestampsAndWatermarks(watermarkStrategy);
// 侧输出流收集迟到数据
OutputTag<Event> lateOutputTag = new OutputTag<Event>("late-events") {};
SingleOutputStreamOperator<Result> result = mainStream
.keyBy(Event::getUserId)
.window(TumblingEventTimeWindows.of(Time.minutes(5)))
.allowedLateness(Time.minutes(3)) // 允许 3 分钟延迟
.sideOutputLateData(lateOutputTag)
.aggregate(new MyAggregateFunction());
// 单独处理迟到事件
DataStream<Event> lateEvents = result.getSideOutput(lateOutputTag);
lateEvents.addSink(new LateEventSink());
这套东西调通后,迟到数据的处理率从 30% 提升到了 95% 以上,业务那边也没再因为"数据不对"找我们吵架。
状态管理:从"能用就行"到"扛得住压"
状态管理这块,我踩过最多的坑就是"没想清楚 T"。一开始为了赶进度,直接用 Keyed State 存所有用户的历史行为,跑了一个月后发现,单个 TaskManager 的状态已经到了 50GB,重启一次要恢复半天。
后来改成"分片 + TTL"的方案:
// 分片存储 + TTL 配置
public class UserBehaviorAggregator extends KeyedProcessFunction<String, Event, Result> {
private ValueState<UserProfile> userProfileState;
private ListState<Event> recentEventsState;
private MapState<String, Long> featureCountsState;
@Override
public void open(Configuration parameters) {
ValueStateDescriptor<UserProfile> profileDesc = new ValueStateDescriptor<>(
"userProfile",
UserProfile.class
);
profileDesc.enableTimeToLive(StateTtlConfig
.newBuilder(Time.days(7)) // 用户画像保留 7 天
.setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
.setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
.cleanupInRocksdbCompactFilter(1000) // RocksDB 压缩时清理过期数据
.build());
userProfileState = getRuntimeContext().getState(profileDesc);
ListStateDescriptor<Event> eventsDesc = new ListStateDescriptor<>(
"recentEvents",
Event.class
);
eventsDesc.enableTimeToLive(StateTtlConfig
.newBuilder(Time.hours(2)) // 最近事件保留 2 小时
.setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
.build());
recentEventsState = getRuntimeContext().getListState(eventsDesc);
}
@Override
public void processElement(Event event, Context ctx, Collector<Result> out) throws Exception {
// 实际处理逻辑
UserProfile profile = userProfileState.value();
if (profile == null) {
profile = new UserProfile();
}
profile.updateFrom(event);
userProfileState.update(profile);
// 其他逻辑...
}
}
这套方案把状态总量压到了原来的三分之一,最主要的是把冷数据和热数据分开了——用户画像这种相对稳定的东西保留久一点,最近事件这种高频更新的东西保留短一点, RocksDB 的压缩策略也能更有效工作。
还有一个容易被忽略的问题是"背压"。一开始我们用的 sink 是同步写入 MySQL,高峰期的时候,Kafka 消费速度跟不上生产速度,积压的数据越来越多,最后导致整个作业卡住。
后来改成异步 sink,再配合限流:
// 异步 sink + 限流
public class AsyncMySQLSink extends RichSinkFunction<Result> {
private transient ExecutorService executor;
private transient Semaphore rateLimiter;
@Override
public void open(Configuration parameters) {
executor = Executors.newFixedThreadPool(10);
rateLimiter = new Semaphore(1000); // 限制每秒最多 1000 条写入
}
@Override
public void invoke(Result value, Context context) {
if (!rateLimiter.tryAcquire()) {
return; // 超过限流阈值就丢弃(或记录到侧输出流)
}
CompletableFuture.runAsync(() -> {
try {
// 异步写入 MySQL
writeToMySQL(value);
} catch (Exception e) {
// 异常处理
} finally {
rateLimiter.release();
}
}, executor);
}
}
这么一改,作业的吞吐量提升了三倍,背压问题基本解决了。
流批一体:同一套代码两种模式
前阵子我们把 Flink 升级到 1.18,顺便试了一下 Flink 的流批一体 API。以前总觉得流和批是两套东西,现在看来,Flink 这套"一套代码,两种模式"的想法确实解决了不少痛点。
比如那段统计用户留存率的逻辑,以前要写两套:一套用 Spark 的 batch 模式跑历史数据,一套用 Flink 的流模式跑实时数据。现在可以合并:
// 流批一体 API:同一段代码既能跑批处理也能跑流处理
public class RetentionAnalysis {
public static void main(String[] args) throws Exception {
// 根据参数选择执行模式
boolean isStreaming = args.length > 0 && args[0].equals("streaming");
StreamExecutionEnvironment env = isStreaming
? StreamExecutionEnvironment.getExecutionEnvironment()
: StreamExecutionEnvironment.getExecutionEnvironment()
.setRuntimeMode(RuntimeExecutionMode.BATCH);
// 数据源适配
DataStream<UserEvent> events = isStreaming
? env.addSource(new KafkaSource<>(...))
: env.fromElements(readFromFile(...));
// 统一处理逻辑
DataStream<RetentionResult> results = events
.keyBy(UserEvent::getUserId)
.window(isStreaming
? TumblingEventTimeWindows.of(Time.days(1))
: GlobalWindows.create())
.aggregate(new RetentionAggregateFunction());
if (isStreaming) {
results.addSink(new KafkaSink<>(...));
} else {
results.writeAsText("output/retention_results.txt");
}
env.execute("retention-analysis");
}
}
这样写的好处是明显的:逻辑只维护一份,测试和调试也更简单。但也不是没有坑——有些操作在批处理和流处理下的行为不一致,比如窗口函数的时间语义,你得在代码里显式区分。
几点实际判断
“要不要上实时"得业务先拍板。报表类需求准比快重要;风控、推荐、实时监控那边,延迟直接绑业务指标。
流计算和批处理解决的是不同问题,别硬把离线全量任务塞进 Flink。状态管理是设计阶段就要想的事——留哪些、留多久、怎么 TTL,别等 OOM 了再补。Exactly-once 理想很满,我们线上基本按 at-least-once 跑,靠监控和补数兜底。工具从 Kafka Streams 起步,需求压不住了再上 Flink,比一开始堆全套靠谱。
还在路上的东西
平台日常流量能扛,但在线特征更新、作业血缘追踪、状态变更可观测性都还在待办里。批处理时代一个人能跑全链路,流计算得有人专门盯窗口、背压和 checkpoint。
翻两年前的笔记,很多"将来优化项"还在列表上。每解决一个具体问题,通常会冒出两三个新的。
后记:写这篇文章时翻了两年前的代码,不少"将来优化项"原封不动还在待办里。
版权声明: 本文首发于 指尖魔法屋-实时数据处理折腾手记(https://blog.thinkmoon.cn/post/136-batch-to-stream-realtime-data-processing/) 转载或引用必须申明原指尖魔法屋来源及源地址!
评论
使用 GitHub 账号登录后即可留言,支持 Markdown。