ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

基于Flink构建日志实时分析系统:架构、实战与踩坑总结

基于Flink构建日志实时分析系统:架构、实战与踩坑总结 说到日志实时分析我最深的记忆来自一次差点翻车的线上排查。业务量在一周内翻了三倍日志从每天几百G涨到了几个T。以前出问题我们习惯登录服务器用 grep 一条条追虽然慢但至少能定位后来那套打法彻底失灵了——一条慢接口的调用链路要横跨五个应用、翻十几个文件才能拼出来用户反馈的问题定位到具体原因往往已经过了半小时。痛过之后我们下决心从零搭一套基于Flink的日志实时分析系统把日志采集、清洗、指标计算、告警触发、维度数据同步全部纳入一条当天可见的实时链路。这篇东西就是那次实战的完整复盘从选型、架构、开发到疑难问题排查尽量把决策过程和细节讲清楚适合正在规划日志平台、准备上手Flink或者想把Flink任务托管到现有Java服务里的同学参考。1. 项目背景与整体架构设计1.1 痛点分析与技术选型思路积压的问题很直接日志量暴涨后传统ELK方案开始吃不消。Elasticsearch索引膨胀查询变慢集群节点加了又加成本直线上升。更难受的是我们对实时性的要求——线上出问题时要能在一两分钟内看到异常趋势而ELK从采集到可查询通常有分钟级延迟加上业务高峰期间的写入背书延迟翻倍是常事。当时在技术选型上我们把候选框缩小到了Spark Streaming、Storm和Flink三选一。Storm确实是老牌的纯流处理框架但它的实时性虽然不错状态管理和精确一次处理都要自己额外折腾而且运维成本不低。Spark Streaming用微批模拟流吞吐高但延时普遍在秒级甚至十秒级做日志告警这种高频小窗口场景并不顺手。Flink赢在三点真正的流式计算模型毫秒级延迟内置状态管理、窗口机制和精确一次语义生态里的连接器特别全Kafka、JDBC、CDC、ClickHouse都有成熟封装。日志分析这件事本质上就是“流式事件处理”Flink的思路和场景天然匹配。我不是说Flink在任何场景下都优于其他框架。如果你的需求是跑小时级离线报表Spark挺好用如果你只是简单做日志检索那还是ES主场。但在“实时指标计算复杂事件告警跨系统数据同步”这个组合面前Flink当时确实是最合适的选择。1.2 整体数据链路设计整个架构的设计目标很朴素日志从产生到被查询和触发告警端到端控制在两分钟以内所有离线手工脚本能省则省凡是能改成实时流的都改成实时流。整体链路是这样一个走向业务应用产生JSON格式日志 → Filebeat采集 → Kafka消息队列 → Flink做实时ETL与指标计算 → ClickHouse/告警服务/ES。Kafka放中间做缓冲核心原因有两条。第一日志洪峰不可控大促、线上事故都会让流量瞬间翻几倍没有消息队列缓冲下游直接被冲垮。第二Kafka把生产者和消费端解耦日志格式调整、消费逻辑变更都不需要整个链条停机Flink消费位点可以随时重置重跑数据也方便。Flink这一层承担了所有“需要状态”的计算指标统计、去重、窗口聚合、CEP规则匹配。计算完的结果有两类出口指标类数据写ClickHouse供查询平台和看板使用告警类数据直接推到告警服务再走飞书和短信通知。这里要特别强调一点ClickHouse在我们的架构里不是替代ES而是替代了传统OLAP分析引擎。ES保留原始日志的全文检索能力ClickHouse承接结构化指标类数据的聚合分析。两者分工明确没有混杂。1.3 集群环境与版本选型版本选型上我们用的Flink 1.13Kafka 2.5ClickHouse 21.8JDK 1.8。Flink版本我特意没有用当时最新的1.14因为公司内部依赖的Flink连接器和部分自研代码还没有充分验证兼容性稳妥优先。如果现在新起项目可以直接选Flink 1.16以上版本语法和内存模型差别不大但社区修复的坑会少一些。部署方式是Flink on YARN。原因很简单公司集群本来就跑着离线任务YARN现成的资源管理和队列机制可以直接复用。我们用YARN Session模式跑流任务多个日志Job共享一个Flink集群资源利用率高关键任务单独分配TaskManager槽位避免互相抢资源。现在回头看如果当时任务量再大几十个Flink on Kubernetes可能更好管理但对中规模日志团队来说Flink on YARN已经够用了。运行环境的参数在搭建时踩了不少坑后面第6章会专门讲。2. 日志采集与数据接入2.1 日志格式规范化这个部分看起来不如写Flink代码酷但整个项目的成败很大程度上取决于日志格式规范不规范。我们的业务应用历史上有各种日志风格有的打JSON、有的用空格分隔、有的用自定义分隔符解析起来非常痛苦。这次统一要求所有核心系统输出JSON格式日志包含公共字段时间戳、服务名、实例IP、TraceID、用户ID、接口路径、状态码、响应耗时、请求参数摘要、错误堆栈摘要。Java服务里用Logback很容易配出JSON输出加上一个LogstashEncoder就能搞定。appender nameJSON_FILE classch.qos.logback.core.rolling.RollingFileAppender encoder classch.qos.logback.classic.encoder.JsonEncoder/ /appender日志打到本地文件后Filebeat负责采集并推送到Kafka。日志规范化的收益在Flink解析阶段体现得最明显。如果日志是JSON格式Flink端只要一个ObjectMapper就能完成反序列化字段不齐的脏数据直接过滤掉如果是纯文本自定义格式你可能得写一堆正则和条件分支解析那种代码维护起来非常想摔键盘。2.2 Kafka主题设计与分区策略Topic设计我们按业务线划分app-access-log、app-order-log、app-gateway-log而不是一个大而全的all-log. 这么做的原因很直接——日志量级和应用性质完全不同网关日志和业务日志的消费速度、计算逻辑、保留周期都不一样混在一个Topic里会导致非热点日志永远被饿死。分区数规划是新手容易忽略的隐蔽点。Kafka分区数是Flink处理并行度的天花板假设一个Topic只有3个分区Flink消费这个Topic并行度设为5实际只有3个并行子任务在拉数据另外2个闲置。分区数的设定参考服务端日志写入吞吐我们核心访问日志Topic的峰值写入在6万条/秒左右单分区写入能扛2万条/秒规划了8个分区。给未来半年的增长预了留余量又多加了两个分区最终设了10个。2.3 Flink消费端的核心参数配置在Flink消费Kafka很多人直接抄网上的demo落下一堆隐患。我分享一份我们线上验证过的核心配置。Properties kafkaProps new Properties(); kafkaProps.setProperty(bootstrap.servers, kafka-01:9092,kafka-02:9092); kafkaProps.setProperty(group.id, flink-log-analysis); // 关闭自动提交offset完全由Flink检查点管理 kafkaProps.setProperty(enable.auto.commit, false); DataStreamSourceString sourceStream env.addSource( new FlinkKafkaConsumer(app-access-log, new SimpleStringSchema(), kafkaProps)) .setStartFromLatest();注意enable.auto.commitfalse这一条配合Flink的检查点机制才能保证一致性。Flink检查点成功时会把Kafka的offset也存下来任务重启后自动从上次成功的位置继续消费不会丢数据也不会重复消费。如果你保留自动提交一旦处理逻辑先于offset提交出现问题重启后会丢一批数据。消费端的并行度设置我们踩过一次冤枉坑。当时TaskManager给了4个槽位Kafka分区数10个但消费并行度默认是1结果数据积压在Kafka里下游延迟飙到十几分钟。后来把消费并行度设成8核心逻辑算完端到端延迟降到了秒级。建议是消费并行度尽量接近Kafka分区数略小于或等于留出分区数余量以便扩容场景。Watermark策略用forMonotonousTimestamps还是forBoundedOutOfOrderness取决于业务能不能接受乱序数据。我们的访问日志会产生多跳后端的日志可能比网关先到所以用了forBoundedOutOfOrderness乱序容忍度设置为10秒兼顾延迟和准确。WatermarkStrategyAccessLog watermarkStrategy WatermarkStrategy .AccessLogforBoundedOutOfOrderness(Duration.ofSeconds(10)) .withTimestampAssigner((event, timestamp) - event.getTimestamp());乱序容忍度设太小上游稍有一点延迟波动就会导致大量迟到数据被丢弃设太大窗口计算延迟增加实时性被拉低。10秒是我们测试下来比较平衡的值。3. 核心处理逻辑实现3.1 实时ETL清洗过滤与字段规整日志流进Flink后的第一层处理是实时ETL这块写起来不难但在细节上有很多讲究。我们把ETL拆成了三个算子链路先做反序列化再做数据脱敏和字段补全最后做过滤和质量监控。反序列化阶段用JSONObject.parseObject把原始字符串解析成结构化对象解析失败的非JSON日志不直接丢弃而是旁路写到一个parse-error-log的Kafka Topic里用来反馈日志格式问题。这一步很关键如果只是把解析失败的数据扔掉你永远不会知道某些服务发过来的日志格式已经变了。SingleOutputStreamOperatorAccessLog logStream sourceStream .process(new ProcessFunctionString, AccessLog() { Override public void processElement(String value, Context ctx, CollectorAccessLog out) throws Exception { try { AccessLog log JSONObject.parseObject(value, AccessLog.class); if (log.getTimestamp() 0) { out.collect(log); } } catch (Exception e) { // 解析失败的数据旁路输出 ctx.output(parseErrorTag, value); } } });脱敏和字段补全里最典型的是IP转地理位置。我们加载了一份精简版IP库到Flink的RichFlatMapFunction的open()方法里每条带IP的日志直接查内存映射输出省份和运营商字段。这样做比每次请求远程接口快几个数量级。IP库的精度不用太追求省份级别就够业务分析用了。3.2 PV/UV与响应耗时统计日志分析最核心的一类指标就是PV/UV和接口耗时。Flink的Window API天然适合这类场景代码写起来很简洁。我们的PV统计的粒度是每分钟每个服务每个接口。用TumblingEventTimeWindows开1分钟窗口配合聚合函数完成计数。DataStreamPvUvResult pvResult logStream .keyBy(log - log.getServiceName() | log.getPath()) .window(TumblingEventTimeWindows.of(Time.minutes(1))) .aggregate(new CountAggregate(), new PvWindowResult());UV计算的难点在于去重。最简单粗暴的方式是窗口内把所有用户ID放进SetState但用户量一大状态膨胀会非常严重。后来我们改用HyperLogLog去重牺牲一点精度换内存的指数级下降。我们在实测中100万级用户的UV统计误差控制在1%以内状态内存从原来的几百MB降到了几十MB。做日志指标统计精度到这种程度业务上是完全能接受的。响应耗时统计用的是ProcessWindowFunction在每个窗口结束的时候输出P95、P99和平均耗时。这里我特别提醒一个点计算P99的耗时列表不要全部存下来用一个固定大小的采样窗口就够了否则状态数据会失控。Mermaid不适用这里就不画图了。3.3 异常检测与告警CEP的实战运用业务方最初的需求很简单“接口报错多了要能马上看到”。后来需求升级成“同一用户短时间内连续报错5次以上立刻告警”。这个场景用Flink CEP实现非常合适。第3.3节我们实现了基于顺序匹配的复杂事件规则。PatternAccessLog, ? pattern Pattern .AccessLogbegin(firstError) .where(evt - evt.getStatusCode() 500) .timesOrMore(5) .greedy() .within(Time.minutes(5)); DataStreamAlertMessage alertStream CEP.pattern( logStream.keyBy(AccessLog::getUserId), pattern) .process(new PatternProcessFunctionAccessLog, AlertMessage() { Override public void processMatch(MapString, ListAccessLog match, Context ctx, CollectorAlertMessage out) throws Exception { String userId match.get(firstError).get(0).getUserId(); out.collect(new AlertMessage(userId, 短时间内连续错误, System.currentTimeMillis())); } });注意CEP的模式匹配是在keyBy之后的流上执行也就是只追踪同一个用户的事件序列。一开始我们把模式直接套在全局流上结果就是所有用户的所有错误日志混在一起匹配结果乱七八糟。后来想明白这是典型的“先分组、再匹配”问题加上keyBy之后就正常了。CEP还有一个容易踩的坑是timesOrMore和greedy()的组合语义。greedy()会让匹配尽可能多地吞并事件在连续错误日志场景下它会把同一用户的N条错误全归到一个匹配里避免重复告警。如果去掉greedy可能会因为相邻事件同时满足多个匹配路径而产生重复告警这部分我调了很长时间才理清楚。告警推送我们直接调用公司统一告警服务的HTTP接口Flink端用AsyncIO异步调用避免同步阻塞导致的反压。当时测试发现同步调用在高峰时段会把下游告警服务打挂因为大量日志同时触发了告警异步IO可以把并发请求数控制在一个合理范围。4. MySQL数据同步到ClickHouse实战4.1 为什么需要MySQL同步ClickHouse日志实时分析经常要与维度数据关联。最典型的场景访问日志里只有用户ID业务看板要展示用户所属的会员等级、所在城市、注册时间这些信息都在MySQL的用户表里。如果在Flink计算时每条日志实时查一次MySQL数据库必然被压垮。我们的做法是把维度数据实时同步到ClickHouse里Flink计算出来的指标结果直接和ClickHouse里的维度表关联。ClickHouse大宽表的查询性能极好几十亿行的表做GROUP BY聚合也就是几十毫秒的级别。这个方案的好处是完全不侵入业务系统的数据库同时维度数据的变更可以在秒级同步到分析系统。4.2 基于Flink CDC的同步实现MySQL同步ClickHouse我们选择Flink CDC基于Debezium因为它能同时处理“存量全量数据”和“增量变更数据”。传统方案是先跑一个离线导入再订阅binlog持续解析两套逻辑维护成本非常高Flink CDC把这两个动作统一到一个Job里启动时自动加一个全局锁并做全量快照之后切到增量监听模式。用Flink CDC读取MySQL的表数据代码大概是这个样子MySqlSourceString mySqlSource MySqlSource.Stringbuilder() .hostname(mysql-master) .port(3306) .databaseList(shop) .tableList(shop.dim_user) .username(flinkuser) .password(******) .deserializer(new JsonDebeziumDeserializationSchema()) .startupOptions(StartupOptions.initial()) .build();schema里的JSON数据包含了操作的类型c表示创建u表示更新d表示删除和数据快照。下游解析后按主键写入表。写入ClickHouse时我特意选择了JDBCAppender而不是ClickHouse官方提供的ClickHouseSink原因是对已有Flink版本和连接器的兼容性考虑。JDBC连接器写ClickHouse的批量参数需要仔细调。下面这些参数是我们最终线上稳定运行的值sink.buffer-flush.max-rows设为2000sink.buffer-flush.interval设为5s也就是攒够2000条或者超过5秒就刷一次。JdbcExecutionOptions execOptions JdbcExecutionOptions.builder() .withBatchSize(2000) .withBatchInterval(Duration.ofSeconds(5)) .withMaxRetries(3) .build(); JdbcConnectionOptions connOptions JdbcConnectionOptions.builder() .withUrl(jdbc:clickhouse://clickhouse:8123/shop) .withDriverName(com.clickhouse.jdbc.ClickHouseDriver) .build();这组参数是在吞吐量和延迟之间反复试出来的。批次太小大量时间花在网络往返上写入吞吐上不去批次太大单次写入内存峰值会飙升而且ClickHouse端合并压力加大。4.3 JDBC连接器常见异常排查实录这部分内容当时花了大半天才定位也是热搜里最常提到的问题之一。线上报过最多的异常是通信链路断掉错误大概长这样java.sql.SQLTransientConnectionException: HikariPool-1 - Connection is not available, request timed out一开始以为是连接池配小了把最大连接数调大还是偶发。后来仔细排查发现是Flink作业长时间空闲时JDBC连接器与ClickHouse服务端的空闲连接被服务端关闭但连接池不知道导致下次写入时连接失效。解决的方式是给ClickHouse的JDBC URL加两个参数socket_timeout60000和connect_timeout30000同时在Hikari连接池设置connectionTimeout30000和maxLifetime1800000让连接池自动淘汰过期连接。另一个高频异常是批量写入失败Caused by: java.sql.BatchUpdateException: Code: 50. DB::Exception: Table shop.dim_user doesnt exist.别笑这个我们真遇到过一次但不是表不存在而是ClickHouse的集群副本还没建好查错的副本返回了Missing Table。出现这类问题优先确认ClickHouse集群里表是否同步到所有副本而不是急着改Flink代码。最隐蔽的一个异常是驱动类冲突项目里同时引入了ru.yandex.clickhouse.ClickHouseDriver老驱动和com.clickhouse.jdbc.ClickHouseDriver新驱动Flink作业在提交时会出现NoClassDefFoundError。这个没什么好解法就是彻底清除旧驱动依赖统一到一个驱动版本上。5. Spring Boot整合Flink实践5.1 为什么要在一个应用进程里跑Flink有网友问过我们Flink不是有独立的集群吗为什么还要在Spring Boot里整合理由有几个尤其是在我们这种中大型公司里很现实。第一独立Flink集群的运维权限通常掌握在数平团队手里业务团队想提交一个新任务要提工单、排队、等待批准一套流程走下来半天过去了。但线上日志异常今天等不到明天业务想要的是一个能自己灵活控制的任务生命周期。第二Spring Boot应用本身就是一个常驻的Java进程具备资源环境我们直接把这个进程当作任务的启动器和托管容器。把Flink的JobManager、TaskManager申请逻辑抽出来封装成管理接口一个简单的HTTP请求就能提交或停止任务这极大降低了操作门槛。第三有些公司根本没有独立的Flink集群要用Flink只能依赖云厂商提供的托管服务价格贵、定制性差。在这种情况下把Flink嵌入到已有的Spring Boot网关服务里做成一个轻量级任务调度中心是更务实的玩法。5.2 整合实战与依赖管理整合Spring Boot和Flink的核心是把Flink的local执行环境嵌入到Spring Boot的Spring容器里任务启动挂在ApplicationRunner上。为了让任务可以被远程控制我们又封装了一个内置Controller支持启停任务列表。Component public class FlinkJobRunner implements ApplicationRunner { private final MapString, FlinkJob jobRegistry new ConcurrentHashMap(); Override public void run(ApplicationArguments args) throws Exception { jobRegistry.put(log-pv-job, new LogPvFlinkJob()); jobRegistry.put(mysql-sync-job, new MySqlToClickhouseJob()); for (FlinkJob job : jobRegistry.values()) { job.submit(); } } }任务类统一接口内部用StreamExecutionEnvironment.execute()提交。整合过程中踩坑不计其数最值得说的有三个。第一是依赖冲突。Spring Boot默认引入的spring-boot-starter-web包含大量的第三方依赖Flink自身也带了一套Akka、Netty和Hadoop相关的jar。直接放一起ClassLoader经常加载错版本导致任务提交后起不来。后来我们的工程统一改成了provided作用域把Flink的核心依赖排除掉运行时再单独把Flink jar放到FLINK_HOME/lib里。这个做法有一个潜在问题就是本地启动环境也要配好依赖但线上可靠性大大提升了。dependency groupIdorg.apache.flink/groupId artifactIdflink-streaming-java_2.12/artifactId version1.13.6/version scopeprovided/scope /dependency第二个是FatJar的打包问题。用maven-shade-plugin打包时如果不配置ServicesResourceTransformerFlink在运行时因为找不到META-INF/services里的服务实现会各种报错常见的是Could not initialize class org.apache.flink.runtime.util.EnvironmentInformation。解决方案是在shade插件里加transformer implementationorg.apache.maven.plugins.shade.resource.ServicesResourceTransformer/这个细节我当时盯着报错看了很久周末加班排查才发现网上一个冷门帖子提到了这个坑。第三个坑是本地调试并行度比线上低的问题。Spring Boot进程里的Flink任务默认并行度为1线上数据量大必须手动指定TaskManager数量和槽位数。但本地和集群的并行度不一致会导致测出的Watermark延迟、窗口结果偏差很大调试时很多问题测不出来。我们最终在本地用和线上相同的并行度配置来跑小批量数据虽然会占资源但保证了验证效果的一致性。6. 常见问题排查与性能调优6.1 背压与反压排查思路实现完一整套流程后第一次压测就遇上了反压在Flink UI的指标页可以看到任务边的“背压”颜色变红。背压的本质是生产速率大于消费速率数据在算子之间积压。最常见也最容易误判的原因是数据倾斜。我们的日志流按服务名做了keyBy一些核心服务的数据量比其他服务大了50倍导致负责核心服务分区的子任务成为瓶颈其他子任务闲置。当时Flink UI上并没能直接看到哪个key的量最大只能靠临时打印或者用Counter埋点来统计各分区数据量。后来我们做了一个两步策略先按天数打散key计算完再按真实key聚合。虽然多了一层flatMap和keyBy但负载瞬间均衡下来背压缓解非常明显。logStream .map(log - new DayKeyedLog(log, RandomUtils.nextInt(0, 50))) .keyBy(DayKeyedLog::getDayPrefix) .window(...)背压还有一个常见诱因是AsyncIO等待时间过长一般是下游HTTP服务响应慢导致异步等待队列堆积。排查时先在UI上找到具体卡住的算子再结合inPoolUsage和outPoolUsage判断问题方向——如果inPoolUsage接近100%说明数据进不来很可能是上游Kafka消费速率不够如果outPoolUsage接近100%问题在下游算子处理速度。6.2 检查点失败处理与状态配置检查点是Flink重启用的一致性保障但它本身也是经常出问题的点。我们线上遇到过一次大范围“检查点超时失败”具体表现是作业偶尔重启但业务方反馈数据出现延迟。分析后原因有两个。一是检查点生成的间隔太短默认1分钟一次但我们的状态比较大每次检查点执行要30多秒如果刚好赶上数据洪峰这次检查点还没做完下次又要开始前后重叠导致超时。解决方式是调大检查点间隔并限制并发检查点数量。env.enableCheckpointing(120000); env.getCheckpointConfig().setMinPauseBetweenCheckpoints(30000); env.getCheckpointConfig().setMaxConcurrentCheckpoints(1); env.getCheckpointConfig().setCheckpointTimeout(60000); env.getCheckpointConfig().setExternalizedCheckpointCleanup( ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION);第二个原因是状态后端超容。我们把RocksDB作为状态后端给了TaskManager 4GB的堆外内存但用户维度的宽表状态持续膨胀最终撑爆了RocksDB的存储路径。最好用的办法是给状态TTL把内存中不再活跃的用户状态及时清理掉。StateTtlConfig ttlConfig StateTtlConfig .newBuilder(Time.hours(24)) .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite) .setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired) .build(); ValueStateDescriptorUserDim userDimDescriptor new ValueStateDescriptor(user-dim, UserDim.class); userDimDescriptor.enableTimeToLive(ttlConfig);TLL设24小时业务上足够用了旧用户的状态最多保留一天内存压力大大缓解。6.3 内存调优与资源规划Flink 1.13之后新内存模型对很多人来说是黑箱尤其是TaskManager上JVM堆内和堆外内存的分配不透明导致明明集群内存很大任务却频繁OOM。我们当时给单个TaskManager分配的是8GB内存运行mysql-sync-job和log-pv-job两个作业。结果总内存看着够用但Flink的JVM堆内内存只有3GB左右状态数据超过堆内存上限后直接Full GC作业频繁抖动。后来我们把TaskManager的配置拆开分析taskmanager.memory.process.size8GB、taskmanager.memory.jvm-overhead.fraction0.15、taskmanager.memory.managed.fraction0.3、taskmanager.memory.framework.heap.size256mb。其中托管内存主要给RocksDB使用框架堆内存是固定给Flink内部框架的不能挪用。真正留给用户代码的堆内存反而没有配置全局总大小那么富余。这里要给个经验性建议不要凭感觉配8GB或者16GB先用默认比例跑再在UI的“Memory”页看实际使用占比。如果堆内利用率长期高于85%就适当增加TaskManager进程总内存或者降低托管内存比例给堆内让出一部分空间。如果堆内利用率很低但任务依然OOM那问题大概率在RocksDB的托管内存上。内存调优没有银弹一定是在监控基础上逐步调整出来的。6.4 容灾与任务重启策略最后补充一下任务的容灾设置。Flink作业在遇到检查点失败、任务异常退出时默认不会自动重启这在大规模日志场景下不可接受。我们在环境里统一配置了固定延迟重启策略。env.setRestartStrategy(RestartStrategies.fixedDelayRestart( 3, Time.seconds(60) ));加上启动时的-yarnD参数把Flink作业提交到独立的YARN容器里作业挂掉后YARN可以快速拉起新的容器配合RocksDB状态恢复实际故障重启时间能控制在两分钟内。这个是用过多次的保命措施强烈建议大家一上来就配好。最后再聊几条实操心得这些配置和代码大家拿过去就能用但有几句体己话还是想多说一点。第一日志实时分析的技术难点其实不在Flink本身而在数据质量治理。你如果把日志格式控制好了Flink侧的解析逻辑会非常薄反过来日志脏乱差Flink代码越写越复杂延迟越高最终谁也救不了。我们项目后期很大一部分精力放在推进各业务线日志规范上那是整个实时平台平稳运行的根基。第二不要把Flink当黑盒用。出问题首先要看Flink自带的UI指标尤其是“数据量”“检查点”“背压”这几页90%的问题都能在UI上找到线索。UI看不懂再上日志和断点避免盲人摸象式乱试。第三性能和稳定性要按项目周期持续投入。第一天搭好平台并不难难的是在业务量翻十倍、翻二十倍之后依然稳定。我们用的是持续监控定期容量评估的做法每月看一遍各任务的数据吞吐和资源使用趋势提前扩容或调优避免在业务高峰当天手忙脚乱。这个项目的日志实时链路到现在已经稳定运行大半年了偶尔还会有新问题冒出来比如新版本的连接器行为和参数改了、ClickHouse副本扩容后写入分布变化了但整套体系和排查思路是稳的。后面我打算把Kafka到Flink这段的跨集群容灾再梳理一版等稳定了再分享出来。
返回列表