☰
基于 Apache Flink 的实时日志分析系统实战:从 Kafka 日志消费、侧输出分流到定时器告警(flink-learning-project-log 源码解析)
2026/10/5 6:49:48 网站建设 项目流程
  • 示例工程
  • 大数据

【免费下载链接】flink-learning

flink learning blog. http://www.54tianzhisheng.cn/ 含 Flink 入门、概念、原理、实战、性能调优、源码解析等内容。涉及 Flink Connector、Metrics、Library、DataStream API、Table API & SQL 等内容的学习案例,还有 Flink 落地应用的大型项目案例(PVUV、日志存储、百亿数据实时去重、监控告警)分享。欢迎大家支持我的专栏《大数据实时计算引擎 Flink 实战与性能优化》

项目地址:https://gitcode.com/gh_mirrors/fl/flink-learning
点击查看免费下载

本指南以 flink-learning 仓库中的flink-learning-project-log模块为主体,完整讲解一个"日志消费 → 分流统计 → 错误告警"的端到端实时日志分析系统:LogAnalysisJob负责基于事件时间窗口的日志量统计,ErrorLogAlertJob负责按服务粒度的 ERROR 日志计数与阈值告警。读完本文,你将掌握 Side Output、AggregateFunction + ProcessWindowFunction 组合、KeyedProcessFunction + Timer + ValueState 定时器告警、Watermark 乱序容忍等 Flink 流处理核心技能,并能在本地 Kafka + Flink 1.20 环境下直接运行这两个作业。

模块定位与整体架构

flink-learning-project-log是 flink-learning-project 聚合工程下的一个独立子模块,其职责非常聚焦:从 Kafka 实时消费应用日志,做实时分析与告警。模块目录结构如下:

flink-learning-project/flink-learning-project-log/ ├── pom.xml └── src/main/java/com/zhisheng/project/log/ ├── ErrorLogAlertJob.java # 错误日志告警作业 ├── LogAnalysisJob.java # 日志分析作业 └── model/ ├── AppLogEvent.java # 应用日志事件模型 └── LogStatistics.java # 日志统计结果模型

整个系统的数据流向(摘自 README):

Kafka (log-topic) │ ├── LogAnalysisJob │ ├── 主流 → 按 serviceName+level 窗口统计 → 输出 LogStatistics │ └── 侧输出 → ERROR/FATAL 日志 → 单独处理 │ └── ErrorLogAlertJob └── ERROR 日志 → 按 serviceName 分组 → 定时器计数 → 超阈值告警 → AlertEvent

两个作业共享同一个 Kafka 日志主题,各自独立消费、互不干扰:一个解决"日志量分析"问题,一个解决"错误告警"问题,这正是生产环境将读多份 topic 的常见做法。

数据模型:日志与统计、告警三要素

AppLogEvent:应用日志事件

AppLogEvent.java 使用 Lombok@Data/@Builder定义了应用日志的完整字段,是流经整个系统的核心 POJO:

字段类型说明
logIdString日志 ID
levelString日志级别:DEBUG、INFO、WARN、ERROR、FATAL
messageString日志消息
serviceNameString产生日志的服务名称(分组与告警的关键维度)
classNameString产生日志的类名
threadNameString线程名
exceptionString异常堆栈信息(可选)
timestampLong日志时间戳(事件时间来源)
hostString主机名
traceIdStringTrace ID,用于链路追踪

其中timestamp被两个作业用作事件时间(Event Time),serviceName与level是统计和告警的分组键。

LogStatistics:窗口统计结果

LogStatistics.java 是 LogAnalysisJob 窗口聚合的输出模型,包含统计维度与窗口元信息:

  • serviceName/level:统计维度(服务名 + 日志级别)
  • count:窗口内日志数量
  • windowStart/windowEnd:窗口开始与结束时间(由 ProcessWindowFunction 从窗口上下文获取)

AlertEvent:告警事件

AlertEvent.java 位于公共模块flink-learning-project-common,是 ErrorLogAlertJob 触发告警时的输出模型,字段包括:alertId(UUID)、level(INFO/WARNING/CRITICAL)、ruleName(规则名,如error-log-spike)、message、metricName、metricValue、threshold、timestamp、host。该模型与监控告警类模块(如 flink-learning-project-monitor-alert)同源,便于后续统一接入告警中心。

作业一:LogAnalysisJob —— 日志分流与窗口聚合统计

LogAnalysisJob.java 实现了完整的三段式处理链路,源码主流程如下:

// 1. 配置 Kafka Source KafkaSource<String> kafkaSource = ProjectKafkaUtil.buildKafkaStringSource( ProjectConstants.TOPIC_LOG, "log-analysis-group"); // 2. Watermark 策略:允许 3 秒乱序 WatermarkStrategy<AppLogEvent> watermarkStrategy = WatermarkStrategy .<AppLogEvent>forBoundedOutOfOrderness(Duration.ofSeconds(3)) .withTimestampAssigner( (SerializableTimestampAssigner<AppLogEvent>) (event, ts) -> event.getTimestamp()); // 3. 消费 + 反序列化 + 分配水位线 + Side Output 分流 SingleOutputStreamOperator<AppLogEvent> logStream = env .fromSource(kafkaSource, WatermarkStrategy.noWatermarks(), "Kafka Log Source") .map(json -> GsonUtil.fromJson(json, AppLogEvent.class)) .assignTimestampsAndWatermarks(watermarkStrategy) .process(new ProcessFunction<AppLogEvent, AppLogEvent>() { @Override public void processElement(AppLogEvent log, Context ctx, Collector<AppLogEvent> out) { out.collect(log); // 所有日志输出到主流 if ("ERROR".equals(log.getLevel()) || "FATAL".equals(log.getLevel())) { ctx.output(ERROR_LOG_TAG, log); // ERROR/FATAL 额外进侧输出流 } } }); // 4. 获取侧输出流单独处理 DataStream<AppLogEvent> errorLogStream = logStream.getSideOutput(ERROR_LOG_TAG); // 5. 按 serviceName + level 分组,每分钟滚动窗口聚合 DataStream<LogStatistics> logStats = logStream .keyBy(log -> log.getServiceName() + "|" + log.getLevel()) .window(TumblingEventTimeWindows.of(Time.minutes(1))) .aggregate(new LogCountAgg(), new LogStatsWindowFunction()); env.execute("实时日志分析系统");

Kafka Source:消费日志主题

作业通过公共模块的 ProjectKafkaUtil.java 构建 KafkaSource,其底层封装了 Flink 新版 Kafka Connector:

public static KafkaSource<String> buildKafkaStringSource(String topic, String groupId, String brokers) { return KafkaSource.<String>builder() .setBootstrapServers(brokers) .setTopics(topic) .setGroupId(groupId) .setStartingOffsets(OffsetsInitializer.latest()) .setValueOnlyDeserializer(new SimpleStringSchema()) .build(); }

关键点:

  • 消费主题来自 ProjectConstants.java 的TOPIC_LOG = "project-log-topic",两个作业分别使用独立消费者组log-analysis-group与error-log-alert-group,互不影响位点;
  • setStartingOffsets(OffsetsInitializer.latest()):作业启动时从最新位点消费,适合只关心新日志的分析/告警场景;
  • 默认 broker 地址为localhost:9092(DEFAULT_BROKER_LIST),本地单机 Kafka 可直接运行;
  • 数据以 JSON 字符串形式到达,随后通过GsonUtil.fromJson(json, AppLogEvent.class)反序列化为AppLogEvent(见 flink-learning-common 中的 Gson 工具)。

Watermark 策略:容忍 3 秒乱序

WatermarkStrategy.<AppLogEvent>forBoundedOutOfOrderness(Duration.ofSeconds(3)) .withTimestampAssigner((event, ts) -> event.getTimestamp());
  • 使用forBoundedOutOfOrderness生成水位线,允许日志事件乱序 3 秒,水位线 = 已见最大事件时间 − 3 秒;
  • 事件时间取自AppLogEvent.timestamp;
  • 这里注意一个细节:fromSource阶段使用WatermarkStrategy.noWatermarks(),待 map 反序列化完成后再用assignTimestampsAndWatermarks(watermarkStrategy)绑定带时间戳提取器的水位线策略,这是反序列化后提取事件时间的标准写法。

Side Output:ERROR/FATAL 日志分流

侧输出流是本作业的第一个核心知识点。通过匿名ProcessFunction,对每一条日志同时执行两个动作:

out.collect(log); // 主流:全部日志 if ("ERROR".equals(...) || "FATAL".equals(...)) { ctx.output(ERROR_LOG_TAG, log); // 侧流:仅 ERROR/FATAL }

侧输出标签在类顶部定义:

private static final OutputTag<AppLogEvent> ERROR_LOG_TAG = new OutputTag<AppLogEvent>("error-log") {};

随后用logStream.getSideOutput(ERROR_LOG_TAG)取出侧流并print("error-log")独立打印。相比filter之后再发往 Kafka 新 topic,Side Output 让"主流全量日志 + 侧流错误日志"在一个作业内共存,避免重复消费、降低端到端延迟,是日志分流的推荐做法。

窗口聚合:AggregateFunction + ProcessWindowFunction 组合

这是本作业的第二个核心知识点,也是"增量聚合 + 全量窗口信息"的经典组合:

  1. keyBy:serviceName + "|" + level作为复合键,保证每个服务每个级别的日志进入各自的窗口桶;
  2. window:TumblingEventTimeWindows.of(Time.minutes(1))事件时间滚动窗口,每 1 分钟切分一次;
  3. aggregate:aggregate(new LogCountAgg(), new LogStatsWindowFunction())同时传入增量聚合函数与全量窗口函数。

LogCountAgg —— 增量聚合:实现AggregateFunction<AppLogEvent, Long, Long>,累加器就是一个 Long:

public Long createAccumulator() { return 0L; } public Long add(AppLogEvent value, Long accumulator) { return accumulator + 1; } public Long getResult(Long accumulator) { return accumulator; } public Long merge(Long a, Long b) { return a + b; }

每条日志到达只做一次+1,内存占用恒定,不缓存整窗口数据,性能远优于全量收集。

LogStatsWindowFunction —— 全量窗口函数:实现ProcessWindowFunction<Long, LogStatistics, String, TimeWindow>,窗口触发时拿到增量聚合的最终计数(elements中只有一个元素),并解析复合键、补上窗口元信息:

Long count = elements.iterator().next(); String[] parts = key.split("\\|", 2); String serviceName = parts.length > 0 ? parts[0] : "unknown"; String level = parts.length > 1 ? parts[1] : "unknown"; out.collect(LogStatistics.builder() .serviceName(serviceName).level(level).count(count) .windowStart(context.window().getStart()) .windowEnd(context.window().getEnd()) .build());

组合的意义:AggregateFunction 保证窗口内内存高效;ProcessWindowFunction 补充windowStart/windowEnd窗口边界信息并产出最终结果对象,兼顾性能与语义完整性。这是 Flink 窗口聚合生产级的标准姿势。

作业二:ErrorLogAlertJob —— 基于定时器的错误日志告警

ErrorLogAlertJob.java 实现"1 分钟内某服务 ERROR 数量超过阈值即告警"的逻辑。核心常量:

/** 告警阈值:1 分钟内超过此数量的 ERROR 日志触发告警 */ private static final int ERROR_THRESHOLD = 10; /** 统计窗口大小:60 秒 */ private static final long WINDOW_SIZE_MS = 60_000L;

主流程:

// Watermark 策略:允许 5 秒乱序(比分析作业更宽松) WatermarkStrategy.<AppLogEvent>forBoundedOutOfOrderness(Duration.ofSeconds(5)) .withTimestampAssigner((event, ts) -> event.getTimestamp()); DataStream<AlertEvent> alertStream = logStream .filter(log -> "ERROR".equals(log.getLevel()) || "FATAL".equals(log.getLevel())) .keyBy(AppLogEvent::getServiceName) .process(new ErrorCountAlertFunction()); env.execute("错误日志告警系统");

处理链路为:Kafka 消费 → 反序列化 → 分配水位线(容忍 5 秒乱序)→filter只保留 ERROR/FATAL →keyBy(serviceName)按服务分组 →process进入带状态的有键控处理函数。

KeyedProcessFunction + Timer + ValueState 的计数告警

ErrorCountAlertFunction extends KeyedProcessFunction<String, AppLogEvent, AlertEvent>是告警作业的核心,其状态与处理逻辑如下。

状态初始化(open 方法):使用两个键控 ValueState 保存每个服务的运行状态:

errorCountState = getRuntimeContext().getState( new ValueStateDescriptor<>("error-count", Long.class)); timerTimestampState = getRuntimeContext().getState( new ValueStateDescriptor<>("timer-timestamp", Long.class));
  • error-count:当前 60 秒窗口内的 ERROR 计数;
  • timer-timestamp:当前已注册的定时器时间戳(用于去重,避免同一窗口重复注册定时器)。

processElement —— 计数 + 注册定时器 + 即时告警:

Long currentCount = errorCountState.value(); long newCount = (currentCount == null ? 0 : currentCount) + 1; errorCountState.update(newCount); if (timerTimestampState.value() == null) { long timerTs = log.getTimestamp() + WINDOW_SIZE_MS; ctx.timerService().registerEventTimeTimer(timerTs); timerTimestampState.update(timerTs); } if (newCount == ERROR_THRESHOLD) { out.collect(buildAlertEvent(ctx.getCurrentKey(), newCount)); LOG.warn("服务 {} 在 1 分钟内 ERROR 数量达到 {} 条,触发告警", ...); }

该函数同时实现了两种告警语义:

  1. 定时器告警:每条 ERROR 到来时计数 +1;如果是该窗口第一条 ERROR,则基于事件时间注册"日志时间 + 60 秒"的定时器。定时器触发(onTimer)时窗口结束,检查并记录窗口内累计的 ERROR 数量后清空状态,进入下一个窗口周期;
  2. 即时告警:当计数恰好达到阈值ERROR_THRESHOLD = 10时,不等定时器触发,立即输出AlertEvent并打 warn 日志,保证告警的实时性——这是"定时器兜底 + 阈值即时触发"的经典组合,避免错误爆发时告警延迟。

onTimer —— 窗口结束清理:

Long count = errorCountState.value(); if (count != null && count > 0) { LOG.info("服务 {} 定时器触发,窗口内 ERROR 数量: {}", ctx.getCurrentKey(), count); } errorCountState.clear(); timerTimestampState.clear();

定时器触发即代表一个 60 秒窗口结束:记录该窗口累计 ERROR 数,随后清空两个状态,开启新一轮计数。注意这里无需等待定时器就能在阈值命中时即时告警,而定时器则保证"即使未达阈值,窗口也会周期性闭合、状态不泄漏"。

buildAlertEvent —— 组装告警事件:

AlertEvent.builder() .alertId(UUID.randomUUID().toString()) .level(ProjectConstants.ALERT_LEVEL_CRITICAL) // CRITICAL 级别 .ruleName("error-log-spike") .message(String.format("服务 [%s] 在 1 分钟内产生 %d 条 ERROR 日志,超过阈值 %d", ...)) .metricName("error_log_count") .metricValue((double) errorCount) .threshold((double) ERROR_THRESHOLD) .timestamp(System.currentTimeMillis()) .host(serviceName) .build();

告警级别使用ProjectConstants.ALERT_LEVEL_CRITICAL(CRITICAL),规则名为error-log-spike,并携带指标名、指标值、阈值等元信息,可直接写入告警主题或对接监控告警模块进一步分发。

涉及 Flink 知识点速查

关联文档的 README 给出了知识点清单,结合源码补充"实际用法"一列如下:

知识点说明所在类源码中的实际用法
Side Output侧输出流,将数据分流到不同通道LogAnalysisJobOutputTag<AppLogEvent>("error-log"),ctx.output(ERROR_LOG_TAG, log)将 ERROR/FATAL 单独导出
AggregateFunction增量聚合函数,内存高效LogAnalysisJobLogCountAgg,累加器为Long,每条日志只做+1
ProcessWindowFunction获取窗口元信息的全量窗口函数LogAnalysisJobLogStatsWindowFunction,从context.window()取getStart()/getEnd()
KeyedProcessFunction有状态的键控处理函数ErrorLogAlertJobErrorCountAlertFunction,keyBy 后每服务独立维护状态与定时器
Timer事件时间/处理时间定时器ErrorLogAlertJobregisterEventTimeTimer(log.getTimestamp() + 60_000),onTimer窗口结束清状态
ValueState键控状态,存储单个值ErrorLogAlertJoberror-count与timer-timestamp两个 ValueState
WatermarkStrategy水位线策略处理乱序数据全部forBoundedOutOfOrderness:分析作业 3 秒、告警作业 5 秒

运行环境与启动方式

本模块依赖 Flink 1.20 生态。根 pom.xml 中的关键版本:

  • flink.version:1.20.3
  • flink-connector-kafka.version:3.4.0-1.20(新版 KafkaSource/KafkaSink Connector)

模块自身 pom.xml 仅依赖flink-learning-project-common,后者再依赖flink-learning-common(Gson 工具、常量等),因此编译运行时需保证本地已安装 Kafka(默认localhost:9092)并创建主题project-log-topic。

启动前提与步骤:

  1. 启动本地 Kafka,确认localhost:9092可访问;
  2. 创建日志主题:kafka-topics.sh --create --topic project-log-topic --partitions 4 --replication-factor 1 --bootstrap-server localhost:9092(分区数与作业并行度env.setParallelism(4)配合可提升吞吐);
  3. 按 JSON 格式向主题写入日志事件,字段需与AppLogEvent对应,例如:
    {"logId":"1","level":"ERROR","message":"timeout","serviceName":"order-service","className":"OrderService","threadName":"main","timestamp":1700000000000,"host":"host-1","traceId":"trace-1"}
  4. 以 Maven 方式运行(从仓库根目录执行):mvn -pl flink-learning-project/flink-learning-project-log -am clean package,随后分别启动LogAnalysisJob与ErrorLogAlertJob的main方法;本地调试时也可直接在 IDE 中运行两个 Job 类。

运行后可在控制台观察三类输出:error-log(侧输出流 ERROR/FATAL 日志)、log-stats(每分钟按 serviceName+level 的统计结果)、error-alert(超过 10 条阈值触发的告警事件)。

小结与扩展方向

flink-learning-project-log以不到 300 行代码覆盖了实时日志处理的两大核心场景,且每个场景都对应一组可复用的 Flink 知识点:日志分析使用 Side Output 分流 + 事件时间滚动窗口 + AggregateFunction/ProcessWindowFunction 组合;错误告警使用 KeyedProcessFunction + 事件时间定时器 + ValueState 实现"定时器兜底、阈值即时触发"的告警语义。

在实际生产落地时,可以在此基础上做以下扩展:把log-stats与error-alert输出从print改为写入 Kafka Sink(ProjectKafkaUtil 已内置buildKafkaStringSink)或 ClickHouse/ES 存储;将告警消息对接 flink-learning-project-monitor-alert 的告警分发链路;把阈值、窗口大小等常量参数化,交由配置中心动态调整。这些细节的代码实现均可作为深入学习与二次开发的基础。

  • 示例工程
  • 大数据

【免费下载链接】flink-learning

flink learning blog. http://www.54tianzhisheng.cn/ 含 Flink 入门、概念、原理、实战、性能调优、源码解析等内容。涉及 Flink Connector、Metrics、Library、DataStream API、Table API & SQL 等内容的学习案例,还有 Flink 落地应用的大型项目案例(PVUV、日志存储、百亿数据实时去重、监控告警)分享。欢迎大家支持我的专栏《大数据实时计算引擎 Flink 实战与性能优化》

项目地址:https://gitcode.com/gh_mirrors/fl/flink-learning
点击查看免费下载

相关推荐

上一篇:ScottPlot 安装配置指南:3 分钟跑通第一个交互式图表
下一篇:DashPlayer 视频下载教程:把长视频存到本地,离线也能学英语

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询