实时数据处理折腾手记

实时数据处理相关的坑,多半出在边界条件上。

这次按现象往回追,不先画大图。

批处理的时代终结不了,但真不够用了

先说说我们从哪儿来。两年前,我们这套数据平台就是典型的 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/) 转载或引用必须申明原指尖魔法屋来源及源地址!