☰
从业务需求到技术选型:《Flink 实战与性能优化》第 1.1 节深度解读——你的公司是否需要引入实时计算引擎
2026/10/3 8:40:53 网站建设 项目流程
  • 示例工程
  • 大数据

【免费下载链接】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 实战与性能优化》第一章 1.1 节展开,从公司日常的实时计算需求出发,完整梳理了实时数据「采集 → 计算 → 下发」的完整链路,对比了离线计算与实时计算、批处理与流处理的本质差异,并结合 flink-learning 开源仓库中flink-learning-monitor系列实战模块(监控采集、日志告警、PV/UV 统计、宕机检测等)的源码实现,帮助读者判断自身业务是否真的需要引入实时计算引擎,并为后续章节深入学习 Flink 奠定选型认知基础。读完本节,你将能独立梳理出典型的实时计算业务场景、评估引入实时计算需要面对的四大技术挑战,并知道如何在仓库中找到对应的落地案例代码。

实时计算需求:从业务一线的真实诉求说起

在公司里,作为数据开发工程师,你大概率收到过产品经理、运营甚至领导提出的这样一类需求:

小田,你看能不能做个监控大屏实时查看促销活动商品总销售额(GMV)? 小朱,搞促销活动的时候能不能实时统计下网站的 PV/UV 啊? 小鹏,我们现在搞促销活动能不能实时统计销量 Top5 商品啊? 小李,怎么回事啊?现在搞促销活动结果服务器宕机了都没告警,能不能加一个? 小刘,服务器这会好卡,是不是出了什么问题啊,你看能不能做个监控大屏实时查看机器的运行情况? 小赵,我们线上的应用频繁出现 Error 日志,但是只有靠人肉上机器查看才知道情况,能不能在出现错误的时候及时告警通知? 小夏,我们 1 元秒杀促销活动中有件商品被某个用户薅了 100 件,怎么都没有风控啊? 小宋,你看我们搞促销活动能不能根据每个顾客的浏览记录实时推荐不同的商品啊?

这些需求表面上五花八门,但最根本的业务本质只有一个——实时查看数据信息。而要满足这一本质诉求,整个处理链路必须满足三个环节的实时性:

  1. 实时采集数据:把业务系统产生的数据实时收集起来;
  2. 实时计算数据:对采集到的数据进行实时的加工、聚合、过滤;
  3. 实时下发结果:将计算结果实时推送或写入下游,供告警、存储与可视化展示使用。

只有采集、计算、下发三个环节全部保持实时,用户看到的数据才是最接近实时的,上述需求才算真正落地。

值得一提的是,这类需求并非停留在理论层面。在 flink-learning 仓库中,flink-learning-monitor模块正是文档所述需求场景的工程化落地,其 README 把整套监控体系划分为监控数据采集、告警、日志处理、监控数据存储、PV/UV 统计、监控数据可视化展示六个子模块,与本节「采集—计算—下发」的链路一一对应。

数据实时采集:到底要采集什么

针对上面的各类需求,我们需要实时采集的数据可以归纳为六类:

  • 用户搜索信息
  • 用户浏览商品信息
  • 用户下单订单信息
  • 网站的所有浏览记录
  • 机器 CPU/Mem/IO 信息
  • 应用日志信息

前四类属于业务数据,是促销活动分析、PV/UV 统计、销量 TopN 统计的数据来源;后两类属于运行数据,是服务器监控告警、应用 Error 日志告警的数据来源。

在仓库中,这两类数据都有对应的统一模型。业务运行指标被抽象为 MetricEvent,其结构包含四要素:name(指标名)、timestamp(指标时间戳)、fields(指标字段值,如 CPU 使用率、内存使用率)、tags(指标标签,如集群名、主机 IP),这样一个模型就能统一承载文档所说的机器 CPU/Mem/IO 信息;而日志数据则被抽象为LogEvent,并配套LogSchema用于 Kafka 反序列化,详见 flink-learning-common 的 model 与 schemas 目录。

另外,FlinkJobMetricCollect 演示了如何通过 HTTP 调用 Flink JobManager 的/jobs/overview接口采集 Flink 作业自身的运行指标,这可以视为「服务器运行状态监控」中针对 Flink 集群本身的采集示例。

数据实时计算:拿到数据之后算什么

采集到的数据实时上报后,需要实时的计算逻辑。文档列举了典型计算任务:

  • 计算所有商品的总销售额
  • 统计单个商品的销量,最后求 Top5
  • 关联用户信息和浏览信息、下单信息
  • 统计网站所有的请求 IP 并统计每个 IP 的请求数量
  • 计算一分钟内机器 CPU/Mem/IO 的平均值、75 分位数值
  • 过滤出 Error 级别的日志信息

这些计算任务在 flink-learning 仓库中都有对应的可运行示例。比如「统计网站 PV/UV」的需求,flink-learning-monitor-pvuv子模块提供了三种实现思路,其中 HyperLogLogUvExample 展示了从 Kafka 读取用户访问事件(UserVisitWebEvent),按「日期_页面ID」生成 Redis Key,并通过RedisCommand.PFADD将用户 ID 写入 Redis 的 HyperLogLog 结构,从而以极低的内存开销完成 UV 去重统计的完整链路——这正是文档中「实时统计网站的 PV/UV」需求的标准答案。

而「过滤出 Error 级别的日志信息」的需求,在 LogEventAlert 中有最直接的体现:从 Kafka 读取日志事件后,一行filter(logEvent -> "error".equals(logEvent.getLevel()))即可实现 Error 日志的实时过滤,供后续告警或落库使用。

数据实时下发:计算结果去向何方

实时计算后的数据需要及时下发到下游,文档将下游明确划分为两类:

告警方式(邮件、短信、钉钉、微信)

在计算层将计算结果与阈值进行比较,超过阈值即触发告警,让运维提前收到通知并及时应对,从而减少故障带来的损失。仓库的flink-learning-monitor-alert子模块实现了完整的告警体系:

  • OutageAlert 是「服务器宕机告警」的工程化实现:从 Kafka 读取MetricEvent指标流,通过OutageProcessFunction(1000 * 10, 60)检测机器是否在指定时间窗口内(10 秒粒度、60 秒窗口)失联,从而判定宕机并构造AlertEvent;
  • 告警渠道方面,flink-learning-monitor-alert的 utils 目录 提供了DingDingGroupMsgUtil(钉钉群消息)、DingDingWorkspaceNoticeUtil(钉钉工作通知)、EmailNoticeUtil(邮件)、SMSNoticeUtil(短信)、PhoneNoticeUtil(电话)等多种通知渠道,覆盖了文档所说的邮件、短信、钉钉、微信等告警方式中的绝大部分。

存储(消息队列、DB、文件系统等)

计算结果写入存储后,监控大盘(Dashboard)从存储(如 ElasticSearch、HBase)中查询对应指标即可实时查看监控信息。这样运营可以知道哪些是爆款商品、哪些店铺成交额最高、哪些商品浏览量最多;运维可以时刻了解机器运行状况,出现宕机或不稳定可及时处理;开发可以依据 Error 日志定位项目 Bug;领导可以看到促销活动的成交情况。

仓库中与之对应的是flink-learning-monitor-storage与flink-learning-monitor-dashboard两个子模块:前者提供了将日志、指标数据写入 ElasticSearch 的 SQL 模板(flink_log_2es.sql、flink_metrics_2es.sql),后者负责可视化展示;而日志数据的完整流式处理管道可以参考 LogMain:从 Kafka 读取原始日志(OriginalLogEventSchema)→OriLog2LogEventFlatMapFunction结构化解析 → 告警分支(LogAlert.alert)→ 写入 ES(LogSink2ES.sink2es),一整套「采集—计算—下发」链路一目了然。

整个流程从数据采集到数据计算再到数据下发,任何一环出现问题都会影响最终效果,因此对实时性的要求非常高。

实时计算场景:四类典型业务归类

文档总结了实时计算的常见应用场景,包括交通信号灯数据、道路车流量统计(拥堵状况)、公安视频监控、服务器运行状态监控、金融证券实时跟踪股市波动计算风险价值、数据实时 ETL、银行或支付公司的金融盗窃预警等。作者调研到的行业实际使用场景还包括:业务数据处理(聚合、统计)、流量日志、ETL、安防(公安视频结构化数据、Flink 图片搜索)、风控(主要处理结构化数据)、业务告警、动态数据监控。

归纳起来,实时计算场景大致可以分为四类:

场景类别核心诉求典型业务
实时数据存储微聚合、字段过滤、数据脱敏、组建数仓实时 ETL、实时数仓构建
实时数据分析接入机器学习框架或算法建模分析商品推荐、广告推荐
实时监控告警实时检测异常并通知金融交易风控、车流量预警、服务器监控告警、应用日志告警
实时数据报表实时展示经营数据活动营销销售额/销售量大屏、TopN 商品

其中「实时数据存储 / 实时 ETL / 实时数仓」在仓库中有专门的项目级实践:flink-learning-project-real-time-data-warehouse 子模块正是围绕实时数据仓库场景搭建的案例。

离线计算 vs 实时计算:理解流处理与批处理

流处理与批处理

在对比离线与实时计算之前,需要先厘清流处理和批处理这一对基础概念:

  • 流处理:处理的数据是源源不断且实时到来的,是一种重要的大数据处理手段;
  • 批处理:历史比较悠久、使用场景较多,主要操作大容量的静态数据集,并在计算过程完成后返回结果。

实时计算的流程是:不断从 MQ 中读取采集的数据 → 处理计算(过滤、聚合等简单操作)→ 往 DB 里存储。在计算层你无法感知会有多少数据量过来,只能尽快处理并及时下发。而离线计算则是从 DB(不限 MySQL,还有各种存储介质)读取已固定的数据(前一天、前一星期、前一个月),再做复杂的计算或统计分析,最后生成可供直观查看的报表(Dashboard)。

离线计算的特点

  • 数据量大且时间周期长(一天、一星期、一个月、半年、一年)
  • 在大量数据上进行复杂的批量计算操作
  • 数据在计算之前已经固定,不再会发生变化
  • 能够方便的查询批量计算的结果

实时计算与流式数据的特点

离线计算的数据是固定的,任务通常是定时的(如每晚 0 点计算前一天数据生成报表);而实时计算的数据源是流式的。什么是流式数据?可以这样理解:你在淘宝下单或浏览某件商品后,页面会立刻给你推荐同类商品广告和相似店铺,这背后就是实时数据处理并作出推荐——系统需要不断从你在网页上的点击动作中获取数据,实时分析后给出推荐。

流式数据具有如下特点:

  • 数据实时到达
  • 数据到达次序独立,不受应用系统所控制
  • 数据规模大且无法预知容量
  • 原始数据一经处理,除非特意保存,否则不能被再次取出处理,或者再次提取数据代价昂贵

实时计算的优势

「实时计算一时爽,一直实时计算一直爽」。对于持续生成最新数据的场景,采用流数据处理非常有利:

  • 监控服务器运行指标时,能根据采集上来的实时数据判断,超出阈值立即发出警报;
  • 通过处理流数据生成简单报告,如五分钟窗口聚合数据平均值;
  • 在流数据中进行多维度关联、聚合、筛选,从复杂事件中找到根因;
  • 应用机器学习算法做复杂的数据分析,根据处理结果给出差异化推荐内容(千人千面)。

实时计算面临的四大挑战

实时计算虽好,落地时却要直面四类技术挑战:

  1. 数据处理唯一性:如何保证数据只处理一次?至少一次?还是最多一次?这对应 Flink 的精确一次(Exactly-once)语义能力;
  2. 数据处理的及时性:采集的实时数据量太大可能导致短时间处理不过来,如何保证数据及时处理、不出现数据堆积?这考验引擎的吞吐与背压(Backpressure)处理能力;
  3. 数据处理层和存储层的可扩展性:如何根据采集的实时数据量大小动态扩缩容?
  4. 数据处理层和存储层的容错性:如何保证处理层和存储层高可用,出现故障时服务依旧可用?

这些挑战正是后续 1.2 节重磅介绍 Flink、1.3 节对比 Spark Streaming、Structured Streaming 与 Storm 的重要铺垫——也正是因为有这些需求,才催生了不断涌现的实时计算框架。

小结与反思

本节从实时计算的需求作为切入点,分析了完成这类需求所需的完整过程:实时数据采集 → 实时数据计算 → 实时数据下发(告警 / 存储),随后总结了四类典型的实时计算场景(实时数据存储、实时数据分析、实时监控告警、实时数据报表),并系统对比了离线计算与实时计算的区别(数据固定与否、任务定时与否、结果产出方式),最后提出了实时计算落地的四大挑战。

对照 flink-learning 仓库,本节描述的每个需求几乎都能找到源码级的落地案例:PV/UV 实时统计见 HyperLogLogUvExample,日志告警见 LogEventAlert,宕机监控见 OutageAlert,完整日志处理管道见 LogMain。读者在动手选型之前,不妨先对照本节内容问自己两个问题:你们公司有文中讲到的类似需求吗?目前的方案是离线批处理还是已经引入了实时计算?想清楚这两个问题,再进入下一节正式认识 Flink,选型之路会清晰很多。

  • 示例工程
  • 大数据

【免费下载链接】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
点击查看免费下载
上一篇:Gatsby 多主题组合实战:用 gatsby-theme-blog、gatsby-theme-notes 与组件 Shadowing 构建组合式站点
下一篇:Agentic Awesome Skills 中的 Angular 状态管理:Signal、NgRx 与 RxJS 全模式实战指南

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

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

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

立即咨询