打开招聘网站搜大数据岗位,十个里有八个要求熟悉流处理框架,不管你是做平台开发、数据清洗还是实时大屏,绕不开Spark Streaming或者Flink。但很多人学流处理是割裂的——今天看个Kafka入门的博客,明天抄一段WordCount代码,真正要在生产环境处理业务数据的时候,又在窗口怎么开、状态怎么存、背压怎么调这些细节上卡壳。
这篇文章不打算讲完整的大数据学习路线,也不扯架构师视角的宏大设计,就聚焦一件事:把流处理代码写出工程水准。我会从流处理和批处理的本质差异讲起,帮你在脑子里建立正确的编程模型,然后对比Flink与Spark Streaming这两套主流框架的取舍,再把窗口、时间语义、状态管理这些绕不开的核心技巧逐一拆解,最后用一个Kafka到实时大屏的完整案例,把整个链路串起来。适合已经能用Hive或者Spark写离线任务、打算上手实时计算的读者,也适合项目里刚引入流处理、正在踩坑的工程师。
1. 流处理与批处理的本质差异——想清楚再动手
1.1 无界数据与有界数据的编程思维切换
先聊一个最基本的问题:流处理和批处理到底差在哪?
批处理面对的是有界数据,数据已经躺在HDFS或者OSS上,跑一个MapReduce或者Spark任务,启动、读数据、计算、写结果、结束,生命周期清清楚楚。你可以把同一份数据反复读取重算,出错了重跑一遍就行。流处理面对的是无界数据,数据像水龙头一样源源不断地流进来,没有起点也没有终点,任务一旦启动就得7乘24小时跑下去。水龙头不能关,数据没法重放(至少没法随意重放),所以每个算子都在跟时间赛跑。
这个差异直接决定了编程思维的转变。写批处理的时候,你想的是"怎么把数据算对";写流处理的时候,你得同时想"怎么把数据及时算出来,还要保证算对,而且任务不能挂"。举一个生活化的例子:批处理像月底盘点仓库,货物都到齐了,你慢慢数就行;流处理像快递分拣流水线,包裹一件接一件到,你必须在包裹经过你面前的那几秒内做出判断,扔到正确的格口里,同时还要保证没有包裹被漏掉或者扔错。
这个根本差异衍生出一堆具体问题:数据乱序怎么办?迟到的数据还要不要?状态存在哪里?任务重启之后怎么恢复?哪个环节处理不好,轻则计算结果不准,重则整个链路卡死。
1.2 算子与DAG:流处理的底层抽象模型
理解了有界无界之后,再看编程模型就顺了。不管是Flink还是Spark Streaming,核心抽象都是把数据流看成一条管道,管道上串联着各种算子。Map、Filter、FlatMap这些单条数据转换,KeyBy、Window、Aggregate这些聚合操作,Connect、Join这些多流操作,每个算子负责数据加工流水线上的一个环节。多个算子连接起来构成一张有向无环图,这就是DAG。
DAG这个概念不是流处理独有的,离线计算也有,但流处理里的DAG有一个关键特点:数据是逐条"流过"算子的,而不是整批"喂给"算子的。打个比方,离线DAG像一条传送带上放着装满货物的箱子,每到一个工位就整箱处理;流式DAG像每个包裹单独在传送带上跑,每个工位处理完单个包裹马上往下传。这个"逐条"的特性决定了流处理框架必须考虑单条数据的序列化效率、网络传输开销、内存占用,因为你没有"攒一批再算"的缓冲空间(微批架构除外,这一点后面讲Spark Streaming的时候细说)。
理解DAG模型还有个实际好处:排查问题的时候,你得按照算子链路一层层拆。数据延迟大了,是Source读Kafka慢了,还是Window算子计算耗时,还是Sink写外部系统阻塞了?没有算子链路的概念,排查就跟无头苍蝇一样。我见过太多人遇到流任务性能问题,第一反应是加并行度,结果瓶颈在下游Sink,并行度加了反而把下游数据库打挂。
2. 主流流处理框架选型——Flink与Spark Streaming的取舍
2.1 两种编程模型的核心差异
选框架这件事,网上吵得不可开交,Flink派和Spark派各说各话,但抛开社区热度,本质上就是两种编程模型的路线之争。
Spark Streaming走的是微批(micro-batch)路线,把源源不断的流数据按时间片切成小批量,比如每2秒切一个RDD,然后还是用Spark那套批处理引擎去算。好处是跟Spark批处理共用一套生态,RDD、DataFrame、SQL的API都能用,坏处是你拿到的延迟天然就是2秒起步(切得更小也不是不行,但切得越小调度开销越大,实际项目里1秒以下很难稳)。更关键的是,微批模型下"事件时间"语义支持得比较别扭,因为数据被切进了固定批次里。
Flink走的是真正的流处理路线,数据一条一条地经过算子,每个算子独立推进,配合分布式快照机制实现精确一次(exactly-once)语义。延迟可以压到毫秒级。更重要的是,Flink对流处理场景针对性极强:原生支持事件时间与水位线、自带窗口机制、状态管理内置在框架里、支持增量检查点。这些能力是Spark Streaming后来通过Structured Streaming追赶的。
我自己的体会是,如果你的实时场景主要是"准实时报表",延迟容忍度在分钟级,团队又已经重度使用Spark,那Structured Streaming完全够用,没必要为了追热点引入新框架。反过来,如果你要做实时风控、实时推荐这类对延迟和状态一致性要求都高的业务,Flink几乎是唯一靠谱的选项。
2.2 状态、容错与一致性:选型时真正需要关注的对比维度
很多人选型只看延迟,其实对生产系统来说,状态管理和容错才是更关键的差异点。
先看状态。窗口聚合要到窗口结束才出结果,结果要记着中间值;去重要记着已经见过的key;实时推荐要维护用户最近的行为序列。这些"中间记住的东西"就是状态。Flink把状态做成了框架的一等公民,你可以直接用ValueState、ListState、MapState这些API,框架负责状态的存储、备份和恢复,还支持把状态落到RocksDB,让单任务的状态规模超过内存上限。Spark Streaming的传统DStream API在状态管理上要弱不少,updateStateByKey这类操作在性能和易用性上都不够理想。Structured Streaming后来的状态管理好了一些,但跟Flink相比,在状态后端可插拔、增量checkpoint这些能力上还是有差距。
再看一致性。所谓的exactly-once,指的是故障恢复之后,每条数据对结果的影响恰好只有一次。批处理天然好做到,大不了重跑;流处理很麻烦,因为数据已经消费了一半,任务重启之后Kafka的offset从哪儿提交、状态从哪儿恢复,都得协调好。Flink用分布式快照加两阶段提交,配合Kafka的幂等Producer和事务性Sink,可以达到端到端exactly-once。Spark Streaming的微批模型做exactly-once相对容易一些,因为它本身就是"一批批处理",一批要么全成功要么全重算。这里要提醒一句:框架号称支持exactly-once,不代表你的整个链路就是exactly-once,下游Redis、MySQL这些外部系统如果不支持事务性写入,照样可能重复。
给你一个粗略的选型参考表:
| 维度 | Flink | Spark Structured Streaming |
|---|---|---|
| 处理模型 | 真正的逐条流处理 | 微批处理 |
| 延迟 | 毫秒级 | 百毫秒到秒级 |
| 事件时间/水位线 | 原生支持,成熟 | 支持但能力较弱 |
| 状态管理 | 强,支持RocksDB状态后端 | 中等,状态规模受限 |
| 精确一次语义 | 端到端支持较完善 | 支持,但依赖下游事务 |
| 与批处理生态的融合 | 需要单独维护 | 复用Spark生态 |
3. 核心编程技巧——时间语义、窗口与状态管理
3.1 时间语义与水位线:处理乱序数据的基石
聊完框架,进到真正写代码要面对的核心问题:时间。流处理里有三个时间概念,处理时间(Processing Time)、事件时间(Event Time)和摄入时间(Ingestion Time)。处理时间是算子所在机器的当前时间,事件时间是数据本身携带的业务时间,摄入时间是数据进入流处理框架的时间。
为什么事件时间这么重要?因为现实世界里的数据经常是乱序的。用户在手机上点了一下按钮,这个埋点数据可能要经过移动网络、网关、消息队列好几跳,才最终到达流处理引擎,中间随便哪里抖一下,先产生的数据可能后到。如果你用处理时间来做窗口,算出来的结果就是"数据到达时刻"的统计,而不是"业务发生时刻"的统计。做实时大屏还能忍,做统计报表和对账就完全不能接受。所以凡是对准确性有要求的场景,一律用事件时间。
用了事件时间,就必须面对水位线(Watermark)。水位线的本质是一个"我确信在它之前的数据都到了"的时间标记。比如你设置水位线等于当前观察到的事件时间减去10秒,意思是"我暂时认为比这个时间早10秒的数据都已经到了,窗口可以触发了"。这个10秒也就是乱序容忍度。设小了,更多迟到的数据会被当成乱序丢掉;设大了,窗口结果迟迟不发,延迟变高。我见过很多新手一上来就抄网上配置,把延迟容忍度设成5分钟甚至更长,结果窗口结果经常比预期晚很久才出来。实际项目里,这个值必须结合数据源端的延迟分布来定,一般先按30秒到1分钟起步,然后看数据到达延迟的P95和P99指标再调。
处理迟到数据还有一招叫allowedLateness,设置窗口在触发之后还可以等待一段时间接收迟到数据,来了就触发增量更新。这招适合对准确性要求高、能接受结果多次修正的场景,但别滥用——窗口已经关闭很久了还在等迟到数据,内存和下游写的压力都会变大。
3.2 窗口机制:滚动、滑动与会话窗口的实战选择
窗口是流处理最常用的计算模式,没有窗口你就没法做"每5分钟统计一次"这类需求。三种基本窗口类型要搞清楚。
滚动窗口(Tumbling Window)是固定的时间长度,互不重叠,比如每5分钟一个窗口,1点到1点05是一个,1点05到1点10是下一个。适合做周期性的独立统计,比如每5分钟算一次订单量。实现最简单,性能也最好。
滑动窗口(Sliding Window)是固定长度加固定步长,窗口之间会重叠。比如窗口长度10分钟,滑动步长5分钟,那1点到1点10是一个窗口,1点05到1点15又是一个窗口。这适合做平滑指标,比如最近10分钟的滚动平均值。但要注意重叠带来的计算放大——窗口重叠越多,一条数据进出的窗口越多,计算和内存开销越大。设窗口长度是步长的整数倍是常见配置,但如果你发现某个实时指标需要"每30秒看一次最近1小时的数据",算算这个窗口重叠度,再掂量一下集群扛不扛得住。
会话窗口(Session Window)按数据的活跃间隙来切分,间隙超过设定值就算一次会话结束。这个在用户行为分析里很常用,比如计算一次用户在App里的完整操作序列。Flink的Session Window实现是按key维护每个活跃会话,gap越大,内存里的会话对象越多,要控制好超时时间。
给个经验值参考:做实时大屏的指标统计,滚动窗口段选30秒到1分钟最常见;做在线监控告警,滑动窗口更实用;做用户行为路径分析,才会用到会话窗口。窗口类型选错,后面的调优全是白费劲,所以动手写代码之前先想清楚业务到底要哪种切分逻辑。
3.3 状态管理与检查点机制
窗口聚合只是状态的一个典型场景。更普遍的,只要流处理算子需要"记住"跨数据的信息,就在用状态。去重需要记住所有见过的key,累计求和需要记住中间值,实时特征计算需要记住用户最近的行为。
Flink里状态分两种:算子状态(Operator State)和键控状态(Keyed State)。键控状态是跟key绑定的,每个key一份,用ValueState、ListState、MapState这些API来读写。实际操作里注意几点。第一,状态别存大对象,如果单key状态数据量很大,比如用户行为序列很长,考虑先做聚合压缩再入库。第二,状态不会自动清理,你用了Keyed State但key越来越多,状态只会无限增长下去,必须配合TTL或者业务逻辑上清理。Flink从1.6开始支持状态的TTL,设置合理的过期时间能省很多事。第三,Checkpoint是异步的,但状态后端序列化大状态会有CPU开销,状态规模过大的时候得考虑从内存状态后端切到RocksDB。
我之前带过一个项目,实时去重任务跑了两个月,状态涨到几个T,RocksDB磁盘都快写爆了。查下来是业务上已经过期的key没有及时清掉。后来给状态设置了TTL,把一批弃用的历史key过滤掉,状态降了一个数量级,任务GC也明显好转。这些都是不跑到生产环境不会遇到的坑。
4. 实操案例——从Kafka到实时指标大屏的完整链路
4.1 场景定义与整体链路设计
理论讲一堆,不如跑一个完整的项目。我拿之前做过的一个网约车实时监控大屏来拆解。业务场景很简单:实时统计每个区域的订单量、活跃车辆数和平均接驾时长,前端大屏每30秒刷新。数据源头是App埋点产生的订单事件,先进Kafka,然后流处理消费、清洗、聚合,结果写Redis或者MySQL,后端接口提供JSON数据,前端用ECharts画图。
链路是:Kafka -> Flink(或Spark Streaming) -> Redis/MySQL -> Flask后端API -> ECharts大屏。这个链路覆盖了流处理接入、清洗、计算、对外输出和后端可视化,是大数据领域很典型的一套实时项目组合。搜索引擎里那些"网约车大数据综合项目——数据可视化Flask+ECharts""基于Spark的数据清洗"的帖子,做的差不多就是这条链路。
设计整体链路的时候有个关键决策:聚合结果写到哪。第一个版本我直接让Flink写MySQL,结果实时高峰时段的写入压力太大,MySQL的写入延迟直接拖慢了整个Flink作业的Sink。后来改成写Redis,用Hash结构存每个区域的统计值,后端接口直接读Redis,响应快得多,Flask只做数据透传和格式转换。写数据大屏的场景,Redis这种内存存储比MySQL合适得多;如果必须持久化,也建议先写Redis再加一个异步任务把Redis数据定期刷到MySQL。
4.2 Flink核心作业实现与参数配置
下面看核心代码,我用的Java和Flink 1.17的DataStream API。第一步是配置Kafka Source。这里的几个参数值值得注意:key.deserializer和value.deserializer决定了从Kafka读出来的字节流怎么变成对象;group.id用于管理消费进度;auto.offset.reset建议设成latest,因为实时大屏只需要看新数据,从头消费会把历史数据全读一遍,浪费资源。
Properties kafkaProps = new Properties(); kafkaProps.setProperty("bootstrap.servers", "10.0.0.11:9092,10.0.0.12:9092"); kafkaProps.setProperty("group.id", "ride-order-realtime-group"); kafkaProps.setProperty("auto.offset.reset", "latest"); kafkaProps.setProperty("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); kafkaProps.setProperty("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); DataStream<String> sourceStream = env.addSource( new FlinkKafkaConsumer<>("ride-orders", new SimpleStringSchema(), kafkaProps) );第二步是清洗和解析。线上埋点数据质量没有那么好,字段缺失、格式错误、经纬度越界都常见。我习惯先用一个Filter算子过滤明显脏数据,再用FlatMap解析JSON并转换成业务对象。解析这一步要捕获异常,解析失败的数据单独路由到侧输出流,方便后面审视数据质量,而不是直接把作业搞挂。
SingleOutputStreamOperator<OrderEvent> parsedStream = sourceStream .filter(line -> line != null && !line.trim().isEmpty()) .flatMap(new JsonToOrderEventFlatMap()) .name("parse-order-event"); // 侧输出保存脏数据,便于排查 OutputTag<String> dirtyTag = new OutputTag<String>("dirty-data") {}; parsedStream.getSideOutput(dirtyTag).map(new DirtySinkFunction());第三步是核心聚合。按键分区之后用滚动窗口算每个区域30秒内的订单数。这里有个性能优化细节:如果对实时性要求没那么高,可以先在窗口内做增量聚合,用aggregate的累加器,而不是等到窗口触发的时候再遍历窗口里的所有数据。Flink的增量聚合能显著降低内存和CPU开销。
DataStream<RegionMetric> metricStream = parsedStream .keyBy(OrderEvent::getRegionId) .window(TumblingEventTimeWindows.of(Time.seconds(30))) .aggregate(new OrderCountAggregate(), new OrderMetricWindowFunction()) .name("region-30s-agg");第四步是Sink写Redis。注意两个细节:一是Redis连接要复用,不能每条数据都新建连接,否则连接开销比数据本身还大;二是写入要设计好Redis的key结构。我用的方案是key为"region:{区域ID}:metric",field用时间戳,value是统计指标的JSON,并且设置了一个过期时间,保证数据不会无限堆积。
region:1001:metric -> { "windowEnd": 1720000000000, "orderCount": 128, "activeVehicleCount": 45, "avgPickupTime": 312.5 }4.3 如果技术栈是Spark Streaming/Structured Streaming
如果你团队的技术栈是Spark,这个案例用Structured Streaming也能做,核心逻辑差别不大,就是API风格从算子变成了DataFrame。读Kafka用readStream,聚合用groupBy加window函数,写外部系统用foreachBatch或者writeStream。
val sourceDF = spark.readStream .format("kafka") .option("kafka.bootstrap.servers", "10.0.0.11:9092,10.0.0.12:9092") .option("subscribe", "ride-orders") .option("startingOffsets", "latest") .load() val query = sourceDF .selectExpr("CAST(value AS STRING) as json") .selectExpr("json_tuple(json, 'regionId', 'timestamp', 'orderId') as (regionId, ts, orderId)") .withColumn("window", window(col("ts").cast("timestamp"), "30 seconds")) .groupBy("regionId", "window") .agg(count("*").alias("orderCount")) .writeStream .outputMode("update") .format("redis") .option("host", "10.0.0.21") .option("port", "6379") .option("key.prefix", "region:") .start()Structured Streaming的项目里,"Append"和"Update"两种输出模式要分清楚。聚合结果随着窗口推进不断更新,用Update模式;只导出新增的明细数据,用Append模式。用错模式会导致下游重复写入或者数据缺失。这是Structured Streaming最容易踩的坑之一。
5. 常见问题与排查技巧实录
5.1 数据倾斜:key分布不均导致的拖垮
流处理作业跑了一段时间之后,最常遇到的问题就是数据倾斜。某几个key的数据量巨大,导致这几个key所在的子任务积压大量数据,其它子任务闲着,整体延迟飙升。我遇到过最夸张的一次,某个热门区域的活动订单量是其它区域的几百倍,那个区域所在的窗口任务处理不过来,整个作业延迟从秒级涨到分钟级。
排查数据倾斜,先看Flink UI上各子任务的积压记录数,哪几个子任务明显偏高,十有八九就是倾斜了。再进一步,把key的分布跑一个统计任务,看看有没有热点key。解决方案有几种。如果热点是少数几个固定的key,可以在上游打散,加随机盐再聚合一层,最后再解盐聚合一次;如果热点key本身就对应一个真实的热点实体,比如某个中心城区,那就只能增加并行度并考虑按区域拆分资源配额。记住:打散是一种精巧但不万能的方案,用了随机盐会牺牲一些准确性,加了两次聚合也会增加延迟,权衡之后再上。
5.2 背压:下游处理不过来的信号
背压是流处理里特有的现象。上游算子产生的数据速度大于下游算子的处理速度,数据在下游算子前堆积。Flink UI上能看到背压状态,高背压的时候Source到Sink整条链路都很慢。最常见的触发原因有三个。第一,某个算子的逻辑太慢,比如窗口计算里调用了耗时的外部服务;第二,聚合算子状态太大,访问状态的开销过高;第三是Sink阻塞,写外部数据库的连接池满了。
排查思路我一般按"从下游往上游找"的顺序。先看Sink对应的外部存储,连接池有没有打满,写入耗时有没有飙升。再往上看窗口算子,窗口里的数据量是否过大,是不是key打散不够。再看Source,Kafka消费速率是否满足需求。这里有个误区:看到背压就盲目加并行度。如果瓶颈在下游Sink,加并行度只会让更多数据同时压向Sink,问题更严重。先定位瓶颈在哪里,再针对性调整。
5.3 数据重复与结果漂移
实时任务还有一个经典问题:明明框架保证了exactly-once,最终结果还是和离线对不上。这多半出现在两个环节。一是Kafka的offset提交和结果写入不是原子的,比如先写Redis再提交offset,中间作业崩溃重启,Redis里就有了重复写入。二是下游目标系统没有幂等能力,Flink的exactly-once再怎么保证,如果Redis里的写入不是幂等的,照样有重复计数。
解决思路有两个层面。第一层,让下游支持幂等,Redis可以用时间戳做幂等判断,MySQL可以用唯一键加定时任务对账。第二层,框架端配好两阶段提交,Flink的KafkaSink配合Kafka的transaction可以实现端到端exactly-once。但这要求Kafka集群本身开启了transaction,很多自建的Kafka集群没有开这个特性。所以千万别以为代码里写了exactly-once就万事大吉,最终一致性如何保证,必须从上游到下游全链路审视一遍。
5.4 任务重启恢复的正确姿势
流作业跑着跑着挂了,重启的时候最大的问题是恢复的一致性。如果是无状态作业,重启无脑从最新offset消费就行。有状态作业就得看检查点。Flink提供了从最近一次检查点恢复的能力,但要注意恢复之后的结果可能有一个窗口期和之前不一致,因为未完整处理的中间数据发生了回退。这个窗口期通常很短,实时大屏这种场景能容忍,但对账类场景就要准备好补偿机制。
实际操作上,我建议给每个流作业配置好启动参数,job name、checkpoint目录、状态后端这些写清楚,并且监控检查点的完成情况。检查点频繁失败往往预示着状态后端有问题或者外部系统不稳定,这比业务数据出错更隐蔽,要重点盯。
6. 生产环境部署与调优经验
6.1 资源规划与并行度设置
流作业的资源规划跟离线作业差别很大。离线任务跑完就释放资源,流作业常驻运行,资源占用是持续的。之前在热词里看到"大数据集群部署策略",这里结合流处理说几句。
并行度的设置要分情况。简单的无状态算子并行度可以大一些,反正只做转换。状态算子的并行度受限于状态的分区方式,键控状态的并行度其实就是key的分布,改并行度会触发状态的重分布,开销很大,所以对一个长期运行的作业,建议状态算子的并行度一上来就定好,尽量减少中途调整。一个通用的做法是:状态算子的并行度根据单key的数据量乘以key的数量估算内存需要,然后再加30%的余量。宁可集群多留一些资源,也别因为并行度设得太低导致单个子任务内存不够。
还有一个经验是区分CPU密集和IO密集。纯计算的窗口聚合,CPU是瓶颈,并行度可以跟物理核数匹配。大量读写外部系统的Sink算子,IO是瓶颈,并行度可以低一些,但每次写入的批量要大,减少IO次数。这个思路跟Web服务调优是一样的——先搞清楚瓶颈在哪。
6.2 状态后端的选型与参数调整
状态后端这个细节太容易被忽视,但它在生产环境里决定作业稳不稳定。Flink两种主要状态后端,HashMapStateBackend进内存,快,但状态一多GC就痛苦;RocksDBStateBackend落盘,状态可以很大,但序列化和反序列化有开销,读的速度比内存慢一个量级。
怎么选?状态规模在单个子任务几百MB以内的,用内存后端,性能最好。状态规模上到GB甚至TB级别,老老实实用RocksDB。RoCksDB的调优是个深坑,但至少要注意几个参数:state.backend.rocksdb.memory.managed设为true,让RocksDB使用托管内存,避免堆外内存失控;block cache大小和write buffer大小如果不确定,先用默认,再通过监控看实际命中率来调。我踩过一次坑,RockDB默认的block size偏小,热点key多的统计任务读放大严重,后来调大block size和write buffer,写放大和读延迟都有明显改善。
6.3 监控告警配置
流作业不像离线任务,挂了你第二天起床看调度日志就能发现。流作业挂一分钟,可能就丢失了大量实时数据。监控告警是生产的必修课。
Flink自带的指标已经覆盖了核心信息:uptime、restartsCount、numRecordsInPerSecond、numRecordsOutPerSecond、checkpoint的完成时间和失败次数、背压状态。把这些指标接入Prometheus加Grafana,配几个核心告警:作业down了立刻告警,checkpoint连续失败3次告警,背压持续高超过5分钟告警,Kafka消费lag超过阈值告警。告警渠道最好是能打到人手机上的,比如企业微信、钉钉或者短信。我见过太多团队只在Grafana上配了看板,没有告警,结果作业挂了半小时,指标已经全红了,才被人发现。对实时数据业务来说,这半小时的数据缺口可能永远补不回来。
7. 几点排坑心得与技巧总结
讲了这么多,最后分享几个我在实际项目中总结的体会,算是对这篇文章的收尾。
第一,流处理的学习曲线比批处理陡很多,但核心不在于API怎么用,而在于脑子里有没有建立"数据是不间断流动的"这个模型。我见过不少有三年Spark经验的人,写Flink代码照样写批处理风格,窗口聚合前先collect一个list再去算,完全背离了流处理的初衷。这个思维转换没有捷径,就是多写、多跑、多观察数据在算子间的流动。
第二,生产环境里,稳定性比炫技重要。什么高级API、炫酷函数,都不如一个能稳定跑三个月的作业实在。代码里多写防御性逻辑,对脏数据有容错性;配置里多用保守参数,Kafka拉取频率不高不低,窗口配置跟业务需求对齐,checkpoint间隔设置合理。想开大窗口之前,先假设任务明天就会挂,想清楚怎么恢复。
第三,流处理很难完全离线和在线割裂。做实时计算,几乎一定要跟离线数据对账。每个实时指标,都建议定期跟离线任务算出来的口径做对比,实时和离线偏差太大,说明实时口径有问题。这个"双跑核对"机制,是很多线上故障的最后一道防线。
第四,工具链要舍得投入。Flink UI讲得再天花乱坠,也不如一套完整的监控告警体系帮助大。集群层面关注Kafka的磁盘吞吐和lag,作业层面关注checkpoint和背压,业务层面关注指标值和离线口径的偏差。三层监控都配齐了,才算真正把你的流处理作业放心交给生产环境。
我自己带过不少实时项目,踩过的坑比写出来的多得多。希望你在这个标题之下看到的,不只是几个API的用法,而是流处理工程化的完整轮廓。从模型认知到框架选型,从窗口状态到排错调优,每一步都有自己的"为什么"。把这个"为什么"想明白了,写代码只是最后那一下落地而已。