toc 目录

管理员 - 1

Object2
Object2 摆烂中
tagHome
arrow_back返回
Object2

实时数据处理:从批处理到流处理的范式转移

数据处理正在从'攒一批再处理'向'来一条处理一条'转变。流处理技术让企业能够实时地从数据中获取洞察和采取行动。本文分析流处理的核心技术和应用场景。

实时数据处理:从批处理到流处理的范式转移

实时数据处理:从批处理到流处理的范式转移

传统的数据处理是"批处理"模式。每天凌晨跑一次 ETL 任务,把前一天的数据处理完,生成报表。业务人员第二天早上看报表,分析昨天发生了什么。

这种模式在很多场景下已经不够用了。当一个用户在你的电商平台上浏览商品时,你希望在他浏览的过程中就实时推荐他可能感兴趣的商品,而不是第二天再推送。当一个欺诈交易发生时,你希望在交易完成前就拦截它,而不是事后才发现。

这就是流处理要解决的问题:实时地处理数据,实时地做出反应。

批处理 vs 流处理

Database architecture

批处理和流处理的根本区别在于数据的处理时机。

批处理是"攒够了再处理"。数据先存储起来,积累到一定量后统一处理。就像邮局攒了一天的信件,晚上统一分拣。

流处理是"来一条处理一条"。数据产生后立即被处理,不经过存储环节。就像快递员收到一个包裹就立即派送,不等攒够一批。

两种模式各有优劣。批处理的优势是吞吐量高、实现简单、适合历史数据分析。流处理的优势是延迟低、实时性强、适合需要即时反应的场景。

在很多实际场景中,两种模式是互补的。流处理负责实时的告警和决策,批处理负责深度的分析和报表。Lambda 架构就是这种思路的体现:同时维护批处理和流处理两条链路,结果合并后对外提供服务。

Apache Kafka:数据管道的基石

Kafka 是流处理生态系统中最重要的基础设施。它本质上是一个分布式的、持久化的消息队列。

Kafka 的核心概念是"主题"(Topic)。生产者把消息发送到主题,消费者从主题订阅消息。消息在 Kafka 中持久化存储,消费者可以按自己的速度消费。

Kafka 的几个关键特性让它成为了流处理的基石。第一是高吞吐量。单个 Kafka 集群可以处理每秒数百万条消息。第二是持久化。消息被写入磁盘,可以保留任意长的时间。第三是可重放。消费者可以重新消费之前的消息,这对于故障恢复和数据重新处理非常重要。

Kafka 不只是消息队列,它正在成为"事件流平台"。企业的所有业务事件都可以发布到 Kafka,不同的消费者按需订阅。这种"中心化事件流"的架构让数据的流动变得清晰和可控。

Apache Flink:流处理引擎

有了 Kafka 作为数据管道,还需要一个流处理引擎来对数据进行计算。Apache Flink 是目前最流行的流处理引擎。

Flink 的核心能力是"有状态的流处理"。它不只是对每条消息做简单的转换,还能维护状态信息,进行跨消息的聚合和关联。

一个简单的例子:计算每分钟的交易总额。Flink 需要维护一个"当前分钟的累加器",每收到一条交易消息就更新累加器,当一分钟结束时输出结果。这就是"有状态"的含义。

更复杂的例子:检测欺诈交易。Flink 需要维护每个用户的"最近一小时的交易记录",当新交易到来时,检查它是否符合用户的正常交易模式。如果不符合,触发告警。

Flink 的另一个重要特性是"事件时间"处理。在分布式系统中,消息可能乱序到达。Flink 用"事件时间"(消息实际产生的时间)而不是"处理时间"(消息被处理的时间)来处理数据,保证了结果的正确性。

流处理的应用场景

Data storage visualization

流处理在很多场景中都有应用。

实时推荐是最典型的场景。用户在电商平台上浏览商品时,流处理引擎实时分析用户的行为,更新用户画像,生成个性化推荐。推荐结果在用户浏览的过程中动态更新,而不是静态的。

实时风控是另一个关键场景。金融交易、网络支付、账户登录,这些事件通过流处理引擎实时分析,检测异常行为并即时拦截。延迟从"天级"降低到"秒级"甚至"毫秒级"。

实时监控也是常见应用。服务器的指标数据通过流处理引擎实时分析,检测异常指标并触发告警。不需要等到指标存储后再查询分析。

物联网数据处理也是流处理的重要场景。数以万计的传感器每秒产生大量数据,流处理引擎实时分析这些数据,检测设备异常、优化运行参数、触发控制指令。

流处理的挑战

流处理虽然强大,但也面临独特的挑战。

数据质量是首要挑战。流处理的数据是实时到达的,没有机会做离线的数据清洗。脏数据、乱序数据、重复数据都需要在流处理的过程中处理。

状态管理是另一个挑战。有状态的流处理需要维护大量的状态信息。当状态量超过单机内存时,需要使用外部存储,这会增加延迟。状态的备份和恢复也是一个复杂的问题。

Exactly-once 语义是流处理的经典难题。在分布式环境中,消息可能被重复处理。保证每条消息恰好被处理一次(不多不少)需要复杂的事务机制。

调试和测试也比批处理困难。流处理程序是持续运行的,不像批处理程序那样可以方便地重跑。测试需要模拟流式的数据输入,比批处理的测试复杂得多。

Data processing pipeline

我的判断

流处理正在从"高端技术"变成"基础设施标配"。越来越多的业务场景需要实时的数据处理能力,流处理的需求在持续增长。

对于新项目,建议在数据架构设计时就考虑流处理的需求。即使当前只需要批处理,也可以用 Kafka 作为数据管道,为未来的流处理需求预留架构空间。

对于已有项目,从最需要实时处理的场景开始引入流处理。实时监控和实时告警通常是最容易切入的场景。

Flink + Kafka 是目前最成熟的流处理技术栈。对于大部分场景,这个组合已经足够。

流处理不是要替代批处理,而是补充它。两者各有适用场景,最有效的方法是根据业务需求选择合适的处理模式。

吉祥物