ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Flink日志实时分析实战:Filebeat+Kafka+ClickHouse全链路落地

Flink日志实时分析实战:Filebeat+Kafka+ClickHouse全链路落地 说真的刚开始接触Flink日志实时分析这个项目时我对流计算的认知也只是停留在“能统计实时指标”这个层面。直到我把一套完整链路从Filebeat采集、Kafka缓冲、Flink清洗聚合、最终落到ClickHouse再被前端大屏实时渲染出来才真正理解日志分析为什么要用Flink、以及用了之后能省下多少原本靠人肉捞日志的时间。这篇不是教材是我自己从0到1落地这套系统的全过程记录包括那些踩了两天才爬出来的坑。这套系统的核心场景是每天大约20亿条业务日志包括Nginx访问日志、微服务应用日志、还有部分客户端上报事件。需求是分钟级甚至秒级看到接口QPS、错误率、P95耗时、慢接口TopN、异常堆栈聚类。我选型的时候把Storm、Spark Streaming和Flink放在一起比了很久最后用Flink 1.17 Kafka ClickHouse这套组合落地。如果你也在纠结日志实时分析怎么做或者正在被Flink的JDBC连接器报错折磨这篇文章里会有你需要的答案。1. 日志实时分析的痛点和为什么是Flink1.1 旧方案哪里不够用做这套系统之前公司日志平台用的是最经典的ELK组合。Filebeat采集日志到KafkaLogstash做过滤最后落到ElasticsearchKibana展示查询。听起来没什么大毛病但真正把业务量跑起来之后问题就来了。首先是清洗能力不够。Logstash对单条日志的处理用正则表达式配置写起来费劲不说遇到堆栈多行日志这种场景Groky解析时间从几十毫秒到几百毫秒不等。高峰期每秒几万条日志进来Logstash消费不过来Kafka的Lag持续增长日志出现“现在能看到10分钟前”的尴尬情况。其次是分析能力受限。ELK强在全文检索和简单的聚合但要做跨多个分钟窗口的滑动统计、计算P95、P99分位数、按照调用链聚合不同服务的耗时写出来的聚合查询一个比一个慢。尤其是P95这种分位数计算ES的percentile聚合在大数据量下性能一言难尽查询秒级返回都算好的。还有一种更原始的做法用定时任务跑离线数仓凌晨算昨天的指标。但做日志平台的都会遇到一个场景——线上出问题时老板和运维同事问的是“现在接口错误率多少”“过去10分钟哪些接口变慢了”这时候昨天的报表没有任何意义。所以日志实时分析真正需要的核心能力不是“能搜日志”而是低延迟秒级响应、能处理乱序事件日志到达的顺序不一定等于发生顺序、能做窗口聚合和状态管理、同时保证作业挂了能恢复不丢数据。1.2 流式计算引擎选型对比我当时把主流的三个流处理框架做了一个对比结果直接影响后续技术选型。维度StormSpark StreamingFlink处理模式原生流式来一条处理一条微批处理默认每2秒一个批次原生流式同时支持批流一体事件时间与乱序处理需要自己实现窗口机制相当繁琐支持但早期版本事件时间支持不成熟原生支持Watermark、窗口、迟到数据状态管理Weak状态存储能力有限基于RDD的状态恢复能力一般原生Keyed State Checkpoint支持增量精确一次语义很难实现靠Structured Streaming才有通过Checkpoint和两阶段提交支持运维成本需要自己处理很多容错细节和YARN生态集成好Flink Web UI JobManager/Worker模式相对成熟实际做日志分析我建议盯住三个词事件时间、乱序容忍、状态恢复。日志数据在采集链路中一定会乱序。比如客户端本地缓冲了一批日志网络抖动导致后发生的日志反而先到Kafka。如果按处理时间统计就会把本该属于10点整窗口的错误算到10点01分。Flink的事件时间模型配合Watermark能比较优雅地把乱序问题控制在可接受范围内。状态恢复对应的是夜里作业挂了、第二天早上才发现的情况。Flink的Checkpoint机制会把Keyed State定期持久化到HDFS或S3作业重启后可以从最近一次Checkpoint恢复消费位点、聚合中间结果都不会丢。这一点对“7x24小时不能断的日志指标”来说太关键了。2. 我落地这套系统的整体链路设计2.1 数据从哪来Filebeat与Kafka的接入细节日志源分为三部分Nginx访问日志、Java服务logback输出的应用日志、客户端上报的埋点事件。统一方案是企业内部所有服务器上都部署了Filebeat 8.x通过Filebeat自带的Kafka输出插件直接写入Kafka集群。Kafka这边根据日志类型拆了几个Topicaccess-logNginx访问日志JSON单行格式app-logJava应用日志带traceId、耗时、异常堆栈等字段event-log客户端埋点事件字段比较灵活每个Topic分区数按目标QPS估算比如app-log平均每秒五六万条我就开了24个分区确保单分区压力不至于太大也给后续Flink并行度留了扩展空间。这里有个实操细节如果Filebeat输出到Kafka带上了timestamp字段最好让Filebeat把原始time_local之类的请求时间字段保留下来Flink做事件时间窗口时要用它而不是用Filebeat的采集时间。采集时间和事件时间在高峰期能差出几十秒用错了指标就偏了。日志格式统一为JSON例如{ log_time: 1736900000123, service_name: order-center, instance_id: 10.10.1.25:8080, trace_id: a9f2c1d0e1f2a3b4, level: ERROR, path: /api/order/create, method: POST, status: 500, cost_ms: 1250, user_id: u_10086, stack_trace: ... }统一JSON有两个好处Flink侧解析不用写正则ClickHouse里直接按字段存储分析的时候筛条件就够了。如果团队里还有正常写多行文本日志的Filebeat的multiline配置要提前处理好否则一条异常日志会被拆成多条统计错误率时直接翻倍。2.2 处理中台与存储选型为什么明细放ClickHouse整个架构里Flink是处理中枢存储层我选了ClickHouse而不是Elasticsearch。原因不复杂日志平台大多数查询是“按服务、按接口、按时间段做聚合统计”而不是全文检索。Elasticsearch擅长的是“搜关键字”但对这种“分组聚合出指标”的场景ClickHouse的列式存储和向量化执行引擎有明显优势。我做了一组实测单节点8C16G对100亿条日志明细做“某接口某小时的avg耗时”聚合ClickHouse大概在1秒内返回Elasticsearch这条查询跑了接近15秒。至于明细日志要不要保留我的选择是保留7天过期用TTL自动清理。ClickHouse的TTL表达式配合MergeTree表引擎可以设定数据过期时间省去了单独写清理脚本的麻烦。存储规划分了两层明细层log_detail表全字段保留按天分区TTL 7天指标层aggr_service_metric和aggr_api_metric表按服务或接口聚合后的分钟级指标保留30天聚合指标虽然也能直接在明细表上用GROUP BY现算但每天20亿条明细全量聚合的代价太高Flink窗口计算完成后落指标表查询侧直接扫指标表速度能快很多。2.3 链路模块一览整体链路从下往上可以分为四段采集层Filebeat收集日志打上基础字段写入Kafka缓冲层Kafka做流量削峰填谷扛住日志洪峰计算层Flink集群消费Kafka做解析、清洗、窗口聚合、状态去重之后写入下游存储展示层ClickHouse存明细和指标配合可视化系统做大盘展示和告警Flink作业本身我没有做成一个大而全的“上帝作业”而是拆成了三个作业log-etl-jobKafka - 解析 - 清洗 - 明细写入ClickHouselog-metric-job从Kafka读原始日志做窗口聚合 - 指标写入ClickHousemysql-sync-job专门把MySQL业务库的部分配置表通过CDC同步到ClickHouse给后续关联分析用拆作业的好处是指标作业挂了不会影响明细入库反之亦然。日志分析最怕的就是一个作业挂了所有数据全断拆开之后隔离性好了很多。3. 核心处理逻辑拆解解析、窗口统计与写入ClickHouse3.1 日志解析与脏数据治理Flink工程我用的Java依赖是flink-streaming-java、flink-connector-kafka、flink-connector-jdbc、flink-clients。版本选了Flink 1.17JDK 11。第一个核心逻辑是从Kafka读原始字符串解析成结构化对象同时把解析不了的脏数据单独放到侧输出流不让它影响主链路。DataStreamString rawStream env.addSource(new FlinkKafkaConsumer( app-log, new SimpleStringSchema(), kafkaProps )); OutputTagString dirtyTag new OutputTagString(dirty-log) {}; SingleOutputStreamOperatorLogEvent logStream rawStream .process(new ProcessFunctionString, LogEvent() { Override public void processElement(String value, Context ctx, CollectorLogEvent out) { try { LogEvent event JsonUtils.parseObject(value, LogEvent.class); if (event.getLogTime() null || event.getServiceName() null) { ctx.output(dirtyTag, value); return; } out.collect(event); } catch (Exception e) { ctx.output(dirtyTag, value); } } }); DataStreamString dirtyStream logStream.getSideOutput(dirtyTag); dirtyStream.addSink(new DirtyLogSink());脏数据单独建一个Sink写到一个_dirty_log目录或者表里。这样做的好处是生产环境排查问题时不会出现“作业没跑但不知道丢的是哪些数”的盲区。我们后来很多线上数据问题都是先翻脏数据流定位到是某个新上线的服务日志格式不对再找对应团队改格式而不是两眼一抹黑。3.2 事件时间、水位线、窗口聚合日志分析指标里最有价值的是窗口聚合结果比如“每分钟每个接口的QPS、错误率、P95耗时”以及“每5分钟某服务的调用次数”。处理这块我的Flink作业设置事件时间语义用forBoundedOutOfOrderness生成Watermark允许日志最多迟到10秒。DataStreamLogEvent withWatermark logStream .assignTimestampsAndWatermarks( WatermarkStrategy.LogEventforBoundedOutOfOrderness(Duration.ofSeconds(10)) .withTimestampAssigner((event, timestamp) - event.getLogTime()) );先解释一下Watermark是什么。你可以把它想成快递站等件的过程站点规定“我等10分钟10分钟前发生的快递如果还没到就不等了直接开始派送”。Watermark就是那个“10分钟”的计时器信号Flink每个窗口根据它决定什么时候触发计算。这个10秒不是随口拍的数字而是结合了我们的日志采集链路实测Filebeat采集延迟一般在1秒内Kafka端到端正常情况下在2秒以内但偶尔网络抖动会导致部分日志延迟5-8秒。设置10秒意味着最多容忍10秒的乱序。窗口延迟也相应多了10秒对“分钟级指标”来说完全可接受。然后做窗口聚合。我用的是一分钟滚动窗口加预聚合而不是把所有原始数据堆在窗口里等触发因为这样状态太大会把内存吃爆。SingleOutputStreamOperatorServiceMetric serviceMetricStream withWatermark .keyBy(LogEvent::getServiceName) .window(TumblingEventTimeWindows.of(Time.minutes(1))) .aggregate(new ServiceMetricAggregate(), new ServiceMetricWindowFunction());ServiceMetricAggregate实现AggregateFunction做增量计算QPS累加、错误数累加、耗时总和累加、同时维护一个近似分位数计算器用于P95估算。这些状态只保留窗口内的中间值窗口触发后就清除不会长期占用内存。窗口触发后还调用了.allowedLateness(Time.minutes(2))允许窗口关闭后2分钟内的迟到数据继续修正结果。这样即使有少量跨窗口的迟到日志指标也不会突然“少一段”。代价是后来的迟到数据会触发窗口的onTimer再次输出结果ClickHouse那边需要做幂等更新所以指标表我直接用ReplacingMergeTree引擎配合更新时间字段做合并。3.3 Sink端写入ClickHouse的配置明细Sink我用的Flink官方JdbcSink连接ClickHouse的JDBC驱动。核心配置是批量写入JdbcExecutionOptions execOptions JdbcExecutionOptions.builder() .withBatchSize(1000) .withBatchInterval(Duration.ofSeconds(5)) .build(); JdbcConnectionOptions connOptions new JdbcConnectionOptions.JdbcConnectionOptionsBuilder() .withUrl(jdbc:clickhouse://clickhouse-host:8123/log_db) .withDriverName(com.clickhouse.jdbc.ClickHouseDriver) .withUsername(default) .build(); sink JdbcSink.sink( INSERT INTO log_detail(service_name, path, method, status, cost_ms, level, log_time) VALUES (?,?,?,?,?,?,?), new JdbcStatementBuilderLogDetail() { Override public void accept(PreparedStatement ps, LogDetail detail) throws SQLException { ps.setString(1, detail.getServiceName()); ps.setString(2, detail.getPath()); ps.setString(3, detail.getMethod()); ps.setInt(4, detail.getStatus()); ps.setLong(5, detail.getCostMs()); ps.setString(6, detail.getLevel()); ps.setTimestamp(7, new Timestamp(detail.getLogTime())); } }, execOptions, connOptions );这里有三个关键点要说明。第一JdbcSink本身只保证了至少一次写入不提供精确一次。如果作业重启可能有少量重复写入。处理方式是明细表加一个ReplacingMergeTree引擎以日志原始唯一ID或trace_id作为排序键后写入的重复数据会覆盖旧数据。这样容忍重复但最终展示时数据是干净的。第二batchSize和batchInterval一定都要设置。只设batchSize会导致低峰期日志量不够1000条时数据一直攒在内存里指标延迟很大只设batchInterval会导致高峰期频繁小批量写入。两个一起设取“哪个先到触发哪个”高峰期按批量低峰期按时间。第三ClickHouse JDBC URL里的/log_db这一段不能省。很多刚开始用的人写jdbc:clickhouse://host:8123没指定库名运行时报Table default.xxx does not exist。日志里明明看到建了log_db库就是写不进。这个坑我在后面“JDBC连接器异常”那节会详细说。4. 落地过程中的两个大坑JDBC连接器异常与MySQL同步ClickHouse4.1 排查过程一个JDBC连接器异常反复拉锯的两天有段时间日志明细入库作业总是运行二十多分钟后自动重启大量数据积压。去JobManager日志里翻看到一段异常java.sql.SQLException: No suitable driver found for jdbc:clickhouse://localhost:8123/log_db Caused by: com.clickhouse.jdbc.exceptions.ClickHouseException: ClickHouse exception, code: 60 DB::Exception: Table default.default does not exist第一次碰到No suitable driver直觉是驱动没有打进fat jar。但检查构建产物com.clickhouse:clickhouse-jdbc确实在里面。后来反复比对才发现坑不在“有没有驱动”而是fat jar里同时存在了多个版本的ClickHouse驱动把jar里META-INF/services/java.sql.Driver的注册信息覆盖掉了。解决办法是构建时排除掉旧版本的驱动依赖只保留一个。我们项目里用的是clickhouse-jdbc:0.4.6通过maven的shade插件把冲突的META-INF/services文件合并而不是覆盖再重新构建。第二个异常Table default.default does not exist更隐蔽。真相是JDBC连接串里的/log_db被某个环节丢掉了导致驱动连接到了ClickHouse的default库而项目表都建在log_db里。这个问题的根子是配置文件里用了变量拼接URLjdbc:clickhouse://${CK_HOST}:${CK_PORT}/${CK_DATABASE}其中CK_DATABASE在配置中心被误设成了空字符串。空字符串拼上去之后URL变成jdbc:clickhouse://host:8123/ClickHouse默认解析成default库。后来我把配置校验加了进去连接串里没有/后面跟库名就直接启动失败宁可启动报错也不能让作业带着脏配置跑下去。这个坑排查过程中我还发现另一个高频问题ClickHouse HTTP接口默认8081端口或8123端口但如果使用HTTPS连接需要用https://并指定ssltrue参数不少人测试环境用HTTP连通了上生产改HTTPS后只改了端口忘了加ssltrue就出现Connection refused。这一类问题在Flink作业日志里长得都一样但根因完全不是同一回事排查时建议先确认环境变量和URL参数。4.2 用Flink CDC把MySQL业务库同步到ClickHouse的配置与类型映射日志分析做了一段时间后业务方提了一个需求想按订单维度关联日志比如“订单创建失败时用户日志上下文是什么样的”。这意味着要把MySQL业务库的订单表实时同步到ClickHouseFlink正好可以做这件事。我这里用的是flink-cdc-connector-mysql这个组件流程是Flink SQL定义一个MySQL CDC Source表再定义一个ClickHouse Sink表最后一条INSERT INTO ... SELECT语句完成同步。MySQL侧需要开binlog并且给同步账号授予REPLICATION SLAVE、REPLICATION CLIENT权限。CDCSource的Flink SQL大致如下CREATE TABLE mysql_orders ( id BIGINT PRIMARY KEY NOT ENFORCED, user_id BIGINT, order_no STRING, status STRING, amount DECIMAL(10, 2), create_time TIMESTAMP(3), update_time TIMESTAMP(3) ) WITH ( connector mysql-cdc, hostname mysql-host, port 3306, username flink_cdc, password ***, database-name trade_db, table-name orders, scan.startup.mode initial );ClickHouse侧建表引擎用ReplacingMergeTree以主键字段作为排序键并在后面加一个版本列CREATE TABLE ck_orders ( id UInt64, user_id UInt64, order_no String, status String, amount Decimal(18, 2), create_time DateTime64(3), update_time DateTime64(3), version UInt64 ) ENGINE ReplacingMergeTree(version) ORDER BY id;其中最坑的是类型映射。MySQL的TIMESTAMP(3)在Flink CDC里被映射为TIMESTAMP(3)但ClickHouse的DateTime精度默认到秒不是毫秒。如果直接建DateTime列写入时会把毫秒截断更新时间的排序就会出问题。正确做法是ClickHouse这边用DateTime64(3)接收。还有金额字段MySQL的DECIMAL(10,2)如果映射成ClickHouse的Decimal(10,2)在写入过程中经常因为精度溢出报错我最终统一映射为Decimal(18,2)留足余量。4.3 同步链路的一致性取舍用Flink CDC同步MySQL到ClickHouse要明确一个事实这个链路不是强一致是最终一致。原因有两点。第一Flink CDC Source和JdbcSink之间的事务机制并不完美。MySQL binlog的位点可以通过Checkpoint保存保证Source端不丢数据但ClickHouse的JDBC Sink不支持标准的两阶段提交写入是“尽力而为”的至少一次。作业重启后可能重复读一段binlog重复写入ClickHouse。第二ClickHouse本身不是事务型数据库单次批量INSERT看似原子但分布式表中多个分片的写入没有全局事务。所以设计上必须接受“可能存在短暂重复和短暂延迟”然后用表引擎兜底。MySQL的UPDATE事件在CDC里表现为先输出一条delete语义的旧数据、再输出一条新数据但JdbcSink不会真的执行DELETE只会把新值写入。配合ReplacingMergeTree以主键排序键做去重并让新写入的数据带一个更大的version最终合并时旧数据会被替代。删除事件则通过增加is_deleted标记位实现软删除查询时过滤掉is_deleted1的行。这套方案已稳定运行两三个月经验是不要幻想“严格精确一次”而是把表设计成能容忍重复并最终收敛这样系统的复杂度和稳定性反而都能接受。5. 性能调优与资源预算跑不动先别急着加机器5.1 背压排查与并行度调整Flink作业跑起来之后第一个要盯的指标就是背压BackPressure。简单理解就是下游处理不过来把压力反向传导到上游最终导致Kafka Lag涨。查看办法是Flink Web UI的BackPressure标签页点进去可以看到每个算子的背压状态OK说明正常HIGH大于50%说明已经在积压了。日志实时分析最常见的背压源头是数据倾斜。我们按service_name做keyBy结果“网关服务”一个key的日志量占了全量日志的35%单个子任务承担了远超其他子任务的数据量。两种解法我都用过加并行度把窗口聚合算子的并行度从4调到8但倾斜key所在子任务依然过载治标不治本。做两阶段聚合先按service_name 随机后缀做一次细粒度预聚合再按service_name汇总。比如把key改成serviceName - random.nextInt(5)把原本压在一个子任务上的压力分散到5个子任务。缺点是预聚合阶段的结果是近似值因为随机后缀分散了精确分组但误差在可控范围内。最后我采用的是两阶段聚合加并行度组合方案网关服务那个倾斜问题直接被压平了。但需要注意两阶段聚合只适用于“可容忍小误差”的QPS、错误数这类指标像精确去重这类不能近似需要另想办法。Sink端并行度则是另一个极端并不是越高越好。JdbcSink并行度过高会导致来自不同子任务的批量写入同时打到ClickHouseClickHouse合并线程压力暴增表现为查询变慢、写入超时。我的经验是日志量在每秒几万条时JdbcSink并行度控制在2到4个就够了再多无益。5.2 状态后端与内存设置Flink状态后端有两种常用选择HashMapStateBackend和EmbeddedRocksDBStateBackend。日志实时分析的指标作业窗口内做的是增量聚合用AggregateFunction状态量并不大但如果做了大量去重、Join状态会迅速膨胀。我的经验是纯窗口聚合和清洗作业用HashMapStateBackend性能最好内存管理简单。涉及大状态去重或长窗口的作业用EmbeddedRocksDBStateBackend状态落盘避免堆内存OOM。JVM堆内存不要无限调大。我遇到过有同事把一个4C8G的TaskManager堆内存调到7G结果GC频繁到每秒停顿好几次。Flink的TaskManager内存里除了堆还要给网络缓冲、管理内存留空间。建议按如下方式起步TaskManager规格总内存JVM Heap托管内存网络内存4C8G开发环境6G4G1G1G8C16G生产环境12G8G2G2G堆内存不要超过总内存的三分之二剩下的留给网络缓冲和管理内存否则作业跑一段时间就会出现奇怪的超时或反压。5.3 参考资源分配以我们的生产规模为例日均日志量约20亿条平均每秒2万到3万条高峰期每秒6万条左右事件大小平均500字节。这套规模的集群配置供参考Kafka3节点每节点4C16G磁盘1.5T SSDTopic总分区约100个Flink2个TaskManager每个8C16G并行度总和16ClickHouse2节点每节点8C32G数据盘SSDFlink并行度分配Source并行度4匹配Kafka分区数解析算子并行度4窗口聚合算子并行度8Sink并行度2。这样一套配置跑下来高峰期Kafka Lag基本在0到几百条之间浮动ClickHouse写入能到每秒几十MB查询P95约500毫秒。整体跑了大半年没有遇到明显的性能瓶颈。如果你刚开始可以按这个规模适当缩半起步但期间要留意背压和Lag不能只凭感觉加机器。6. 顺带聊聊Spring Boot集成Flink的几种姿势6.1 三种常见的集成姿势很多后端团队会问“Spring Boot怎么整合Flink”其实这个问题要先想清楚你要的是“在一个Spring Boot应用里跑Flink作业”还是“用Spring Boot管理Flink作业”两个需求完全不一个量级。姿势一在Spring Boot进程内LocalEnvironment执行作业。适合本地调试、测试环境跑小流量数据。缺点是Flink和Spring Boot依赖容易冲突比如Guava版本而且作业提交状态在进程里是动态的线上做重启或扩缩容都不方便。姿势二用Flink SQL Client纯SQL方式写作业并提交到集群Spring Boot只负责生成SQL文件和调用Client脚本。适合纯SQL链路、逻辑简单、不需要写UDF的场景。姿势三把Flink作业打成fat jarSpring Boot作为控制台通过Flink REST API提交和管理作业。适合生产环境多作业、需要统一管理状态的场景。日志实时分析最终适合姿势三因为我们的作业有Java代码逻辑自定义AggregateFunction、UDF解析器等不可能全用SQL表达同时也需要作业管理界面统一看状态。三种姿势都有人用但如果生产环境只有一个作业、代码又简单用姿势二最省心。6.2 推荐方案独立集群加REST API提交Spring Boot独立部署通过Flink自带的REST API把Jar包提交到独立运行的Flink集群Spring Boot本身不持有任何Flink环境对象所有Flink的版本依赖都只在打包好的task.jar里面。核心逻辑是RestClusterClient如果你不想引入额外依赖直接调HTTP接口也可以。大致流程是先把jar上传到Flink JobManager得到jarId再通过jarId触发run请求并附带mainClass和参数。推荐的Maven依赖只需要flink-rest-client这个不需要打包进fat jarSpring Boot独立托管Configuration conf new Configuration(); conf.setString(rest.address, flink-jobmanager-host); conf.setInteger(rest.port, 8081); conf.setString(rest.client.timeout, 30000); RestClusterClientString client new RestClusterClient(conf, flink-cluster); String jarPath /data/flink-jars/log-etl-job-1.0.0.jar; // 1. 上传jar JarUploadResponse uploadResponse client.uploadJar(new Path(jarPath)).get(); String jarId uploadResponse.getJarId(); // 2. 提交作业 client.runJar(jarId, com.example.LogEtlJob, new String[]{ --env, prod, --checkpoint-dir, s3://flink-checkpoints/log-etl });需要特别提醒一点打fat jar时一定要把Flink自身依赖标记为provided不要把flink-streaming-java这类依赖打进去。否则在提交的时候JobManager和TaskManager的classpath和jar内的Flink类冲突作业大概率失败报类似ClassCastException: org.apache.flink.streaming.api.datastream.DataStream cannot be cast的错。这个问题在Spring Boot工程里格外常见因为Maven默认会把所有依赖打进去。如果做一个相对完整的管理平台可以在此基础上再加作业状态查询、作业停止cancel、Checkpoint历史查看、资源池划分等功能。这就从一个“跑作业的小工具”进化成了真正的“实时计算平台控制台”但核心思想还是Flink集群独立部署Spring Boot只做API壳子。最后分享一个实用经验。我刚开始做这套系统时一上来就想把整个实时数仓都搭出来结果花了两周还在纠结各种组件版本。后来我换了个思路先画一条最小闭环——用一条Nginx访问日志从Filebeat采集进KafkaFlink读出来解析再写到ClickHouse最后前端轮询一个接口查这条日志的最新QPS。整个链路跑通之后我才开始往上叠加错误率、P95、异常聚类这些指标。这个做法帮我少走了很多弯路因为越复杂的系统越应该先把最核心的那条链路踩实后面的扩展反而会顺利很多。
返回列表