🚀实战 | 日志数据实时清洗与聚合:一文搞懂实时 ETL 全流程!(附架构+代码)

在大数据体系中,实时 ETL(Extract-Transform-Load) 已成为企业数据流转的命脉。
特别是在日志类业务中,实时清洗与聚合不仅是流式计算的入门场景,更是构建实时数仓、实时大屏、实时监控的基石。

本文将以一个日志实时处理案例为线索,完整讲解从 Kafka → Flink → Hive / StarRocks 的全流程实践,
带你理解实时 ETL 的 核心思路、架构设计、业务逻辑与优化技巧


一、为什么要做实时 ETL?

传统离线 ETL(如 Hive SQL + 调度任务)往往存在:

  • 数据延迟高(T+1 / T+N);

  • 批处理效率低;

  • 无法支撑实时看板、报警系统。

而实时 ETL 的价值在于:

  • 🔥 延迟低(秒级数据可见);

  • 🔁 数据持续更新

  • 💡 更好地服务实时业务场景(如 PV/UV 实时统计、实时销售分析)。


二、日志实时 ETL 典型架构

来看一张经典架构图👇

[Web / App 日志]
       ↓
     Kafka
       ↓
  Flink 实时清洗
       ↓
  1)清洗后落地 ODS Hive
  2)实时聚合入 DWS 层(StarRocks / Druid)
       ↓
     可视化大屏 / 实时监控

📊 架构要点:

  • Kafka 负责 日志数据缓冲与解耦

  • Flink 承担 实时抽取、转换、聚合

  • Hive / StarRocks 提供 查询与分析

  • 最终支撑大屏、报表、监控等实时应用。


三、日志数据样例

假设我们采集到的原始日志为:

{
  "user_id": "u12345",
  "event": "page_view",
  "page_id": "index",
  "duration": 5.2,
  "timestamp": "2025-10-11T12:23:45Z",
  "source": "app",
  "device": "android"
}

业务目标:

  1. 清洗掉空字段或异常数据;

  2. 解析时间、设备来源等维度;

  3. 实时统计 PV(访问次数)、UV(独立用户数)、平均停留时长。


四、实时清洗逻辑(Flink 实现)

DataStream<String> rawStream = env
    .addSource(new FlinkKafkaConsumer<>("log_topic", new SimpleStringSchema(), props));

DataStream<LogEvent> cleanStream = rawStream
    .map(json -> JSON.parseObject(json, LogEvent.class))
    .filter(e -> e != null && e.getUserId() != null && e.getPageId() != null)
    .assignTimestampsAndWatermarks(
        WatermarkStrategy.<LogEvent>forBoundedOutOfOrderness(Duration.ofSeconds(3))
            .withTimestampAssigner((event, ts) -> event.getTimestamp())
    );

✨ 优化要点:

  • 使用 Watermark 处理乱序数据;

  • filter() 提前过滤脏数据;

  • 将清洗逻辑尽量下推,减少后续计算负担。


五、实时聚合逻辑

我们以「每分钟」为窗口统计 PV、UV、平均时长:

cleanStream
    .keyBy(LogEvent::getPageId)
    .window(TumblingEventTimeWindows.of(Time.minutes(1)))
    .aggregate(new LogAggFunc(), new WindowResultFunc())
    .addSink(new FlinkKafkaProducer<>("agg_topic", new SimpleStringSchema(), props));

LogAggFunc 实现:

public class LogAggFunc implements AggregateFunction<LogEvent, Tuple3<Long, Set<String>, Double>, Tuple3<Long, Long, Double>> {
    public Tuple3<Long, Set<String>, Double> createAccumulator() {
        return Tuple3.of(0L, new HashSet<>(), 0.0);
    }
    public Tuple3<Long, Set<String>, Double> add(LogEvent value, Tuple3<Long, Set<String>, Double> acc) {
        acc.f0 += 1; // PV
        acc.f1.add(value.getUserId()); // UV
        acc.f2 += value.getDuration(); // 总时长
        return acc;
    }
    public Tuple3<Long, Long, Double> getResult(Tuple3<Long, Set<String>, Double> acc) {
        return Tuple3.of(acc.f0, (long) acc.f1.size(), acc.f2 / acc.f0);
    }
    public Tuple3<Long, Set<String>, Double> merge(...) { return null; }
}

这样每分钟就可以输出:

page_id=index, PV=3000, UV=2580, avg_duration=6.1s

六、结果写入与落地

  • ODS 层(Hive):存清洗后的原始日志(用于追溯与批分析);

  • DWS 层(StarRocks):存聚合结果,支持秒级查询;

  • ADS 层(大屏):直接从 DWS 获取最新数据。

📦 示例 Sink:

cleanStream.addSink(new HiveSinkFunction<>("ods_log_event"));
aggStream.addSink(new StarRocksSink<>("dws_page_agg"));

七、常见问题与优化建议

问题场景解决方案
数据乱序设置合理 Watermark 延迟(2~5 秒)
小文件过多通过分桶或合并写入方式减少文件数
Checkpoint 超时调整间隔 + RocksDB 状态后端
Kafka 堆积增加分区 / 并行度,优化消费者吞吐
精度问题使用窗口合并、幂等 Sink 防止重复写入

八、案例落地效果

在某旅游数据平台项目中,系统每秒处理 50,000+ 条日志事件

  • Kafka 延迟:< 1 秒;

  • Flink ETL 延迟:约 3 秒;

  • StarRocks 实时查询延迟:< 500 ms;

  • 大屏刷新频率:每 5 秒。

实现真正意义上的「分钟级实时分析」✅


九、总结:实时 ETL,不止是“清洗”

实时 ETL 的核心价值在于:

把数据从“混乱”变成“有序”,
从“原始”变成“可分析”,
从“日志流”变成“业务洞察”。

它不是单纯的数据通道,而是企业数据实时化的基础设施

📌 如果你觉得这篇文章对你有所帮助,欢迎点赞 👍、收藏 ⭐、关注我获取更多实战经验分享!
如需交流具体项目实践,也欢迎留言评论

更多推荐