做日志实时分析,大部分人第一反应是ELK那一套,但等数据量真的冲到几十亿条一天、业务方又开始要分钟级聚合报表的时候,Elasticsearch那套全文检索加简单聚合的逻辑就不太够用了。我前两年负责过一套日志分析平台,每天接入的业务日志量在五十亿条左右,高峰期每秒峰值能到十几万条,要求是分钟级出PV、UV、接口成功率、耗时分布这些指标,还要支持按业务线、机房、接口维度实时切换。折腾一圈之后,我最后把核心计算引擎从Spark Streaming换成了Flink,整套链路在存储层保留ClickHouse,接入层继续用Kafka,中间用Flink SQL做实时清洗和聚合。这篇文章想把这些经验完整地梳理一遍,包括Flink在日志实时分析中的定位、整体架构设计、SpringBoot整合Flink的坑、MySQL同步到ClickHouse的常见问题,以及JDBC连接器异常这类高频故障的排查方法。
如果你正准备用Flink搭一套日志分析平台,或者作业已经上线但经常被反压、checkpoint失败、数据延迟折腾得焦头烂额,这篇内容应该能让你少走一些弯路。
1. 从日志平台的痛点看Flink的定位
1.1 日志实时分析难在哪里
日志数据看起来结构简单,但真要做实时分析,难点全在“规模”和“实时性”上。首先是流量不平滑,白天业务高峰和凌晨低谷之间的吞吐差距可能达到十倍,作业必须能扛住突发流量,不能因为高峰就把反压打满。其次是日志格式不统一,同一个系统里可能有Nginx访问日志、Java服务logback输出、客户端埋点日志,字段名和类型经常对不上,解析规则各有各的坑。第三是乱序问题非常明显,客户端上报的数据经常晚到几十秒甚至几分钟,如果不处理乱序,统计出来的分钟级数据就会忽高忽低,报表根本没法看。
这三个问题叠加在一起,意味着实时计算引擎必须具备几个能力:可扩展的吞吐能力、稳定的状态管理、能处理事件时间乱序的窗口机制。而这些恰好是Flink的强项。我最早用的Spark Streaming,微批模式在吞吐上不差,但遇到需要精细化事件时间处理和端到端延迟控制的场景,用起来总觉得隔了一层。Flink把“事件时间”“水位线”“状态后端”“checkpoint”这些概念直接融入框架,写出来的作业天然就是为了处理无界乱序流。
还有一层考量是运维成本。日志分析的指标需求变化很频繁,今天要按用户端类型分组,明天要新增一个错误码分类,如果用Java写DataStream API,每次改动都要重新写代码、打包、上线,效率太低。Flink SQL把很大一部分计算逻辑变成了声明式查询,改动一个GROUP BY维度或者加一个过滤条件,改几行SQL就行,这点在业务需求快速迭代的场景里特别值钱。
1.2 我为什么选Flink而不是Spark Streaming和Storm
讨论技术选型时,团队内部也争论过几次。Storm的延迟确实低,但它的消息处理语义是At-Most-Once偏多,做精确统计要靠外部存储配合去重,开发和维护成本高。Spark Streaming的生态成熟,但微批天生的调度延迟摆在那里,虽然Spark Structured Streaming在努力追赶,遇到要精准处理事件时间的需求还是不够顺手。
Flink最打动我的是它的“流处理优先”理念。它把流当作最基础的执行模型,批处理反而是流的一种特化。也就是说,同样的逻辑在实时流和离线批里可以复用,SQL API在这两种模式下能保持一致。我们团队的实际情况是,实时指标跑在Flink上,离线T+1报表也想统一口径,Flink的流批一体让我们能把一套SQL逻辑用在两条链路上,省掉了大量口径对齐的沟通成本。
当然,Flink不是没有门槛。State和Checkpoint的配置、反压的处理、连接器参数调优,这些都要在实际场景里踩过坑才能真正掌握。但我觉得这笔学习成本是值得的,尤其是当业务规模上去之后,Flink的稳定性表现确实让人省心。
2. 整体架构设计与关键选型
2.1 一条日志从产生到报表的完整链路
我搭建的这套架构,整体上分为五个环节:采集、传输、实时计算、存储、查询展示。日志首先由业务应用通过logback的appender写入Kafka,这一层只负责把日志快速搬走,不在业务进程里做太多加工,避免影响业务接口性能。Kafka起到削峰填谷和消息缓冲的作用,实时计算引擎从Kafka拉数据,既可以重放,也可以并发扩展。
Flink作业从Kafka消费日志,在作业内部完成格式解析、脏数据过滤、字段补齐、事件时间窗口聚合,然后把结果写入ClickHouse。明细数据也会被写入ClickHouse的一张日志明细表,方便业务方后面按需查询。最后是查询层,我们用的也是ClickHouse做数据服务,配合一个自研的轻量查询接口,前端报表轮询这个接口就能拿到分钟级趋势数据。
这套链路里每个环节的选型都有明确理由。采集端没有用Filebeat直接怼到Flink,而是统一走Kafka,因为一旦业务实例数量多了,直接让Flink连日志文件会非常复杂,而且不好做多副本备份。存储端用了ClickHouse,是因为日志分析场景下绝大多数查询是“按时间范围+维度分组+聚合统计”,这类分析型查询ClickHouse是吞吐天花板。MySQL在整个链路里没有承担实时日志存储的职责,但我后面发现不少业务方希望把指标结果回传MySQL做关联查询,所以也单独做了一条MySQL同步到ClickHouse的链路,这部分后面会专门讲。
2.2 Kafka分区与消费并发怎么定
Kafka分区数量直接决定了Flink作业的并行度上限,也决定了消费吞吐量。分区太少,Flink的Source并行度上不去,消费能力受限;分区太多,又会让每个分区上数据量偏少,还会增加管理和重平衡的开销。我当时的经验是:先按峰值吞吐和单分区消费能力来估算。假设单分区每秒能稳定消费5MB数据,日志高峰期每秒总量是80MB,那分区数至少要有16个,我给每个Topic留了40%到60%的余量,实际用了24个分区。Flink作业的Source并行度尽量和Kafka分区数保持一致,避免并行度大于分区数导致部分线程空转,也避免并行度小于分区数造成分区处理不均。
分区设计的另一个关键是Key的选取。日志里如果按接口维度做聚合,Kafka的Key可以直接用接口名,这样同一接口的日志会落到同一分区,Flink读取时天然按接口分组,后续做窗口聚合时shuffle成本会小很多。但如果Key字段的基数特别大,比如用user_id做Key,就很容易出现热点分区,某个高频用户的日志会把单一分区打满。这一点在日志分析场景里尤其要小心。
2.3 结果存储为什么用ClickHouse
日志实时分析的结果存储,我前后对比过MySQL、Elasticsearch和ClickHouse。MySQL在数据量几千万以内还凑合,但到了数十亿级别,聚合查询的响应时间完全不可控。Elasticsearch适合做搜索,做深度分页和明细检索很顺手,但高基数维度聚合的性能比较一般,内存占用也高。ClickHouse是列式存储加向量化执行,对这些“时间范围过滤加维度分组计数”的查询场景优势特别明显,压缩比还很高,我这边日志明细表压缩后只有原始磁盘数据的十分之一左右。
很多人有一个误区,以为ClickHouse是实时数据库,用它就需要每一条都立刻写入。其实ClickHouse更适合批量写入,每次插入几百上千行,配合分区裁剪和TTL数据生命周期,性能才会发挥出来。我们在Flink里用ClickHouse JDBC连接器,设置sink的批量大小和flush间隔,而不是逐条写入,实测下来写入吞吐能提高好几倍。存储选型这件事,还是要回到查询模式来反推,先想清楚报表要查什么,再决定数据落在哪里。
3. 接入与清洗:先把数据弄干净
3.1 日志格式统一与Schema设计
实时分析最容易翻车的环节不是引擎,而是数据质量。我接手的时候,各个业务线的日志格式五花八门,有的用竖线分隔,有的用JSON,有的直接在message字段里塞了一整段业务信息,解析逻辑极其痛苦。后来我们统一规定:所有业务日志必须输出为JSON格式,并且强制包含time、level、service、traceId、msg这些通用字段,业务自定义字段统一放在extra对象里。规范定下来之后,Flink端的解析逻辑就变得非常简单,直接用JSON函数提取字段。
Schema设计上要注意字段类型的前后兼容。日志字段加了一个枚举值、改了某个字段含义,在实时链路里是经常发生的事。Flink SQL里如果字段类型定义得过于严格,比如把接口耗时定义成INT,结果线上出现一个小数,数据直接就丢了。我建议对不确定范围的数值字段统一用BIGINT或者DOUBLE,时间字段统一用BIGINT存储毫秒时间戳,展示层再去格式化,这样能避免很多隐式类型转换导致的脏数据。对于日志分析这种场景,宁可字段宽一点,也不要因为类型太严把数据拒之门外。
3.2 脏数据过滤与字段补齐
日志解析之后,第一步就是过滤。大多数日志平台的做法是在Flink作业里加一个WHERE条件,把不是当前业务线的日志、字段缺失的日志、或者明显是测试数据的内容过滤掉。但要注意,过滤不能放在最后,越早过滤越好。我们是在Source之后紧接着做解析和过滤,让进入窗口计算的数据从一开始就是干净的,不干净的数据不参与后续的状态更新和窗口计算,能省掉很多无谓的计算开销。
字段补齐是另一个容易被忽略的点。线上日志有时候会因为客户端版本落后,缺少某些新加的字段,但下游报表又要求这些字段不能为空。我会在Flink SQL里用COALESCE给字段设置默认值,比如平台类型为空就填unknown,耗时为空就填0。还有一个经验是:不要相信客户端上报的时间,客户端时钟经常不准,真正的事件时间最好以日志接入网关接收到消息的时间为准,这个时间会在日志里单独记录为ingest_time,我们用这个字段作为事件时间,能有效减少客户端时钟偏移带来的统计误差。
3.3 乱序数据处理与水位线
日志从客户端产生到进入Kafka,中间可能经过网关、消息队列,延迟从几百毫秒到几分钟不等。如果不处理乱序,窗口计算的结果就会出现漂移。Flink SQL里处理乱序主要靠WATERMARK和窗口时延设置。我常用的配置是:WATERMARK FOR ingest_time AS WITH OFFSET,允许数据晚到30秒,窗口使用TUMBLE或者HOP。这个30秒是结合业务实际情况调的,普通接口日志的延迟基本在10秒内,30秒的余量足够,再大的延迟就直接放到延迟侧输出流里单独处理,避免等得太久拖慢整体实时性。
这里有一个关键心得:水位线设置太短,晚到数据丢失多,分钟级报表会偏低;设置太长,窗口结果迟迟不触发,实时性又受影响。不要追求一个一劳永逸的值,最好根据线上延迟的分位数动态调整。我这边后来做了一套简单监控,统计每条日志从产生到入库的延迟,看P95延迟,再反推水位线应该设多少,效果比拍脑袋强很多。
4. Flink SQL计算与窗口设计
4.1 用SQL写实时聚合的写法
日志分析里最核心的计算就是窗口聚合。举一个最简单的例子,统计每五分钟每个接口的PV和平均响应时间,Flink SQL大概长这样:
CREATE TABLE kafka_source ( service STRING, endpoint STRING, consume_time BIGINT, event_time BIGINT, watermark for event_time as with offset ) WITH ( 'connector' = 'kafka', 'properties.bootstrap.servers' = 'kafka:9092', 'topic' = 'app_log', 'scan.startup.mode' = 'latest-offset', 'format' = 'json' ); CREATE TABLE clickhouse_sink ( service STRING, endpoint STRING, window_start TIMESTAMP(3), pv BIGINT, avg_consume_time DOUBLE ) WITH ( 'connector' = 'jdbc', 'url' = 'jdbc:clickhouse://clickhouse:8123/log_db', 'table-name' = 'endpoint_metrics', 'sink.buffer-flush.max-rows' = '1000', 'sink.buffer-flush.interval' = '2s', 'sink.buffer-flush.max-retries' = '3' ); INSERT INTO clickhouse_sink SELECT service, endpoint, TUMBLE_START(event_time, INTERVAL '5' MINUTE), COUNT(*) AS pv, AVG(consume_time) AS avg_consume_time FROM kafka_source WHERE consume_time IS NOT NULL GROUP BY service, endpoint, TUMBLE(event_time, INTERVAL '5' MINUTE);这段SQL看着简单,但有几个隐含的成本点。GROUP BY的维度越多,Flink需要维护的窗口状态就越多,内存压力就越大。如果维度组合基数特别大,比如按接口加用户ID分组,状态量会爆炸。所以日志类指标我一般只做中低基数的维度聚合,高基数的明细查询交给ClickHouse去做,不在Flink里做全维度的实时明细聚合。
4.2 状态管理:不要让状态无限膨胀
Flink的窗口计算离不开状态,窗口状态在窗口触发之后默认会被清理,但有些场景下状态不会自动清。最典型的是用了聚合函数,且不是窗口聚合,而是无限流上的普通聚合,比如计算累计PV。这种无界聚合的状态会一直增长,如果没有TTL配置,内存和磁盘都可能被拖垮。
我在Flink配置里给状态设置了TTL,比如聚合中间状态保留5分钟,因为我们的指标基本都在分钟级窗口内完成,超过这个时间再来的数据基本没有意义。
state.backend: rocksdb state.backend.rocksdb.ttl: 300s用RocksDB作为状态后端也是规模上来之后的必然选择。堆内存状态在几个GB以上容易出现GC抖动,RocksDB把状态落到磁盘,内存占用可控,代价是读写性能会下降一些,但在日志分析这种吞吐优先的场景完全可以接受。设置状态TTL的时候要留足余量,太短会导致晚到数据无法正确累加,太长又浪费存储。我一般会和窗口时延对齐,比如窗口时延30秒,聚合状态TTL给5分钟,这样既覆盖了乱序窗口的闭合,也避免了无界状态膨胀。
4.3 多维指标与明细库如何取舍
实时分析平台最容易陷入的误区是想用Flink满足所有查询需求,把几十个维度的组合全部预先聚合成结果表。这样做的直接后果是Flink作业状态爆炸、并行度怎么也提不上去、SQL越来越复杂。我更推荐的做法是:常用固定维度组合在Flink里做预聚合,其余查询交给ClickHouse明细表。
ClickHouse在明细表上的聚合能力相当强,只要设计好分区键和排序键,用一条SQL就能把任意维度组合的指标查出来。比如我把订单日志的明细表按照事件时间做月分区、按service字段做一级排序,查询时先通过分区裁剪缩小数据范围,再在ClickHouse里做GROUP BY。这样绝大部分临时分析需求都不需要在Flink里开发新作业。只有那种每天被报表固定调用、对响应时间要求极高、维度组合相对固定的查询,才会在Flink里预聚合。一热一冷分开处理,整个系统的资源利用率和开发效率都能得到比较好的平衡。
5. SpringBoot整合Flink与MySQL同步ClickHouse
5.1 SpringBoot工程里集成Flink要注意什么
用Java开发Flink作业,很多人习惯在SpringBoot工程里直接写Flink代码,把Flink作业当成一个SpringBoot应用来启动。这样做确实方便,能复用Spring的配置和Bean管理,但也会遇到几个非常典型的坑。
第一个坑是依赖冲突。SpringBoot自带的Logback会和Flink的日志框架冲突,导致作业提交时控制台疯狂刷警告,或者干脆把作业给搞挂。我的解决办法是把flink包里的log4j和slf4j相关依赖排除掉,用Maven的exclusion把冲突的依赖剔除干净。第二个坑是类加载机制。Flink在集群上会把用户的Jar包和Flink本身的依赖隔离,如果SpringBoot的fat jar里带了一堆Flink依赖,很容易出现NoClassDefFoundError。我建议开发环境用SpringBoot管理业务Bean,但生产运行用Flink的原生提交方式,让Flink作业以独立Main函数启动,Spring容器只负责提供配置。
第三个坑是序列化问题。Flink内部会对数据类型做序列化,如果用Spring的复杂对象直接作为Flink流的元素类型,性能会非常差。我一般会把日志数据定义成简单的POJO或者直接用JSON格式的字符串在流里传递,业务字段需要解析时再用Flink SQL处理。这几个问题解决之后,SpringBoot整合Flink才能做到既享受Spring的便利,又不影响Flink作业的稳定性。
5.2 MySQL到ClickHouse同步的实现方案
项目里有一部分数据是从MySQL同步到ClickHouse的,主要是把业务系统的配置表、维表数据定期同步到分析平台,供日志数据做关联查询。最直接的做法是用Flink CDC监听MySQL的binlog,实时把变更同步到ClickHouse。但要注意,CDC在初期同步和历史数据回填上都有自己的机制,不是简单建一个source就能跑。
我当时用Flink CDC的MySQL连接器读取业务库的binlog,然后通过JDBC写到ClickHouse。同步任务的核心是启动时的全量快照加增量监听,Flink CDC会先做一次全量扫描,然后无缝切到binlog增量模式。
public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.enableCheckpointing(5000); env.getCheckpointConfig().setCheckpointStorage("hdfs:///flink/checkpoints"); env.getCheckpointConfig().setMinPauseBetweenCheckpoints(500); env.setRestartStrategy(RestartStrategies.fixedDelayRestart(3, Time.seconds(10))); DebeziumSourceFunction<String> source = MySQLSource.<String>builder() .hostname("mysql-host") .port(3306) .databaseList("biz_db") .tableList("biz_db.config_table") .username("flink_cdc") .password("******") .deserializer(new JsonDebeziumDeserializationSchema()) .build(); DataStream<String> stream = env.addSource(source); stream.map(new MySQLToClickHouseRecord()).addSink(clickhouseSink()); env.execute("MySQL CDC to ClickHouse Sync"); }这里有一个容易被忽略的点:如果你需要全量快照同步一张大表,默认CDC机制的并发度在初期可能会比较低,因为快照阶段使用的是单线程读取。解决方法是给source设置scan.incremental.snapshot.chunk.size,让全量阶段也能按照分片并行读取。这个参数调得好,上百GB的配置表同步也能在几分钟内完成。
5.3 同步任务中的字段映射与类型转换
MySQL和ClickHouse的字段类型差异,是同步任务中最常出现报错的地方。比如MySQL的DATETIME类型,在ClickHouse里有DateTime和DateTime64两种对应,精度不同;MySQL的TINYINT经常会被人误当成布尔类型,但在ClickHouse里就是Int8。如果映射的时候不做显式转换,数据同步过去之后查询结果很容易和源库对不上。
我建议在建ClickHouse表的时候,把字段类型和MySQL严格对齐,能直接用相同精度的类型就不用默认推断。对于可能会变化的字段,比如金额、百分比,统一用Decimal类型,并且在小数位数的定义上留一定余量。还有一个来自实践的建议是:同步任务里不要直接在Flink里做复杂的JOIN,CDC流只负责把单表数据搬过去,表之间的关联逻辑交给ClickHouse来做,这样既能保证同步任务的稳定性,也方便后续在分析侧灵活组合。
同步完之后的一致性问题也要考虑。MySQL里的数据改了,ClickHouse里可能因为链路延迟或者写入失败导致不一致。我这边给同步任务加了一个监控:记录每张表的同步位点和同步时间,定期对账,发现不一致就在业务低峰期做一次全量重新同步。实时链路不是写完就完事了,数据一致性追踪同样要投入精力。
6. 常见问题与排查技巧实录
6.1 JDBC连接器异常:现象、原因、处理
Flink作业里的JDBC连接器异常,几乎每个用Flink做过实时分析的人都遇到过。最常见的报错是连接超时、连接池耗尽、以及写入数据量过大导致的背压。我踩过最狠的一次坑,是Flink作业写入ClickHouse时,因为单批次数据量设置得过大,ClickHouse服务端直接报了Too many partitions异常,Flink作业不停重启。
问题根源在于JDBC Sink的flush时机和ClickHouse的分区机制不匹配。ClickHouse每个分区在写入后会生成目录,如果单次insert涉及的分区数量过多,就会把服务端资源打满。解决方案是:把sink.buffer-flush.max-rows和sink.buffer-flush.interval的值调小,让每个批次的数据量更收敛;同时检查表的PARTITION BY表达式,尽量让数据落到少数几个分区。
另一个高频问题是数据库连接空闲超时。MySQL和ClickHouse服务端默认都会回收空闲连接,但如果Flink的JDBC连接池没有及时感知,连接会被服务端断开,作业里就会出现Connection is not available的报错。我给JDBC连接器配置里加了连接存活检查,并且把连接池的idleTimeout设置为小于数据库server端的wait_timeout,问题就基本消失了。连接器异常看起来五花八门,其实排查思路都类似:先看服务端日志,再看Flink JobManager日志,最后检查连接器参数有没有和服务端配置冲突。
6.2 反压和checkpoint卡住
反压是流处理作业里最普遍的健康问题。Flink的反压机制是通过任务节点的背压指标体现的,一旦Source端出现高反压,说明下游计算速度跟不上上游数据的产生速度,Kafka消费就会出现Lag。我处理反压的思路是分层定位:先看是Source慢、算子慢还是Sink慢。最常见的是Sink写入ClickHouse太慢,这时优先优化写入方式,比如把逐条insert改成批量insert,或者给ClickHouse增加写入并发。
Checkpoint卡住则是另一个让人头疼的问题。每次checkpoint的超时时间如果反复失败,作业就会有持续恢复的风险,严重的时候会造成数据重复和延迟叠加。我遇到过的情况是状态太大了,RocksDB的checkpoint写入HDFS耗时过长,于是把checkpoint的超时时间调长,同时增加两个checkpoint之间的最小间隔,让状态快照和业务处理错峰。还有一种情况是算子之间的数据堆积,导致barrier迟迟无法对齐,这时要先处理掉反压问题,再去优化checkpoint参数,因为checkpoint卡住很多时候只是反压的一个次生现象。
6.3 数据倾斜与延迟的优化
日志分析里的数据倾斜,最典型的现象是:某些高频接口或者某些大客户的日志占了单一子任务的绝大部分数据,导致作业整体吞吐卡在那个热点子任务上。我遇到过一次,某个网关服务的错误日志在高峰期突然增长了十几倍,结果所有错误日志都落到同一个Kafka分区,Flink那个子任务的并行度再怎么提高也没用。
处理手段有几个。第一,Kafka消息Key要均匀,不要用基数很低的字段做Key。第二,Flink端在必要的时候可以做两阶段聚合,即先打散Key进行一次部分聚合,再做全量聚合。第三,如果倾斜只出现在某个时间窗口,可以考虑用滚动窗口加局部预聚合来降低热点压力。日志场景里,热点倾斜通常和错误日志聚餐有关,这类问题最好在源头就进行限流或者抽样处理,比如对错误日志单独设置采样率。
数据延迟方面的优化则更细致。除了水位线之外,我会定期查看Kafka消费Lag、Flink算子处理延迟、ClickHouse写入耗时三个指标,任何一个出现明显上升都说明链路上有瓶颈。还有一些细节比如启用minibatch、本地聚合,也能在吞吐上有明显收益,尤其是高基数维度的场景,本地聚合能大幅减少shuffle的数据量。
6.4 监控和报警需要盯哪些指标
实时链路稳定运行的基础是监控。我常用的监控指标分成三层。第一层是作业健康指标:JobManager状态、TaskManager存活数、checkpoint成功率和耗时。第二层是消费链路指标:Kafka消费Lag、每分钟消费条数、解析失败条数、过滤丢弃条数。第三层是数据质量指标:写出到ClickHouse的行数、延迟数据的比例、指标结果和离线对账的偏差。
报警阈值要结合实际业务来设。Kafka Lag不是一有增长就报警,比如凌晨低峰期Lag增长可能是正常波动,但高峰期的持续Lag就需要处理。我一般是按分钟级别监控Lag的斜率,连续三分钟以上持续上升才触发报警,报警内容里附带当前作业的ID和最近一次checkpoint的状态,方便值班人员快速定位。这套监控体系上线之后,很多问题在用户反馈之前就被处理掉了,这也是实时链路能不能长期稳定运行的底层保障。
7. 最后提醒几个能减少折腾的细节
说了这么多,最后再分享几个实际操作中容易被忽视的地方。
第一,每个Flink作业上线前,一定要做一次墨菲测试。所谓墨菲测试,就是人为制造异常,比如把下游ClickHouse表停掉、把Kafka Topic数据清空、把网络断开几分钟,看作业会怎么表现。我曾经以为这些操作都不会有太大影响,直到有一次线上MySQL出了故障,才发现Flink作业里的JDBC Sink在重连失败后反复重启,把下游数据库的连接数彻底打满。提前把故障预案想好,比真正故障发生后再去救火省心得多。
第二,日志分析的结果可靠性离不开对账。我每周会拿Flink的实时统计结果和离线批处理的结果做对比,偏差超过预设阈值就去找原因。很多隐形问题,比如某个字段在历史数据里的格式变了、某些数据被窗口延迟丢弃,都是在对账过程中暴露出来的。永远不要觉得实时结果看起来正常就一定正确。
第三,Flink版本和连接器版本一定要绑定清楚。我见过不少人因为使用了和Flink主版本不匹配的连接器,导致作业能提交但是运行一段时间后出现各种奇怪异常。每次升级Flink版本之前,我都会先把连接器和依赖的兼容性列表对照一遍,避免在半夜被上线后的诡异问题折腾。
实时链路的上手门槛并不高,真正拉开差距的是对细节的把握和长期稳定运行的能力。希望这些踩坑经验能让你在搭建日志实时分析平台的时候少走一些弯路。