ARTICLE DETAIL

资讯详情

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

Kafka + Flink SQL实战:Java实时统计项目完整指南

Kafka + Flink SQL实战:Java实时统计项目完整指南 很多Java团队的第一套实时统计都是从Kafka加Flink SQL这个组合开始的。原因很简单数据已经在Kafka里躺着实时统计这件事如果一个个用Java写消费者、算窗口、管状态写到后面你会发现你其实是在重新造轮子。而Flink SQL把这块几乎全部封装好了你只需要定义几张表写一条INSERT INTO语句一个能实时输出的统计任务就跑起来了。这篇文章我会用Java工程的方式带你完整搭一遍从本地起Kafka、造模拟数据到用Flink SQL创建Kafka源表、配置Watermark、做滚动窗口聚合再把结果打到控制台或者MySQL。也会把我实际踩过的几个坑比如窗口不输出、Kafka报Error while fetching metadata这类问题放在最后统一说。适合刚接触Flink SQL、又不想一上来啃官方文档的Java开发也适合想快速验证实时统计方案的架构师。1. 实时统计项目选型与整体数据流1.1 这个项目到底在解决什么问题先说场景。我们最常见的实时统计需求长这样系统里不断产生用户行为日志、订单流水、支付回调业务方希望看到最近一分钟、最近五分钟的访问量、订单量、成交金额最好还能按商品、按渠道、按用户维度拆开看。这类需求听上去不难但真要用Java手写你会面临几个绕不开的问题。第一是乱序数据怎么处理。同一个用户先点A商品再点B商品但由于网络、日志采集等环节两条消息到达Kafka的顺序可能是反的。按处理时间统计肯定不准按事件时间统计就需要自己维护水位、缓存迟到数据。第二是窗口的状态怎么存。每分钟一个窗口每个窗口里要保存中间结果还要定期清理光这一块代码量就不少。第三是故障恢复。任务挂了重启要从上次的offset继续消费并且不能重复计数这里的一致性处理又是一大摊子事。Flink SQL把这些全部做进去了。你在DDL里声明一下事件时间和Watermark窗口聚合的中间状态、触发时机、迟到数据处理框架自动帮你管。你只需要把精力放在业务SQL怎么写上。这套方案目前也是国内绝大多数实时数仓、实时大屏项目的标准玩法不是某个小圈子里的偏门技术。1.2 为什么是 Flink SQL 而不是 DataStream API我经常被问到既然大数据领域有DataStream API能写Java代码做实时统计为什么还要绕一层SQL我的看法是DataStream API确实更灵活能处理各种定制化逻辑但如果你做的不是特别偏门的实时流处理SQL 90%的场景都覆盖了而且开发效率高一个量级。我给你一个直观对比维度Flink SQLDataStream API备注开发成本几条DDL加一条INSERT即可需要写Source、Transformation、Sink全链路代码SQL可以把整个Pipeline压缩到几十行窗口计算内置滚动、滑动、会话窗口语法简单需要自己调用window等算子再处理trigger、evictor窗口越复杂SQL优势越大状态管理框架自动管理设置TTL即可需要自己定义StateDescriptor并管理生命周期SQL省去大量状态细节乱序处理Watermark策略一句话配置手动assignTimestampsAndWatermarksSQL更贴合“声明式”思路复用能力SQL结果能直接被BI、报表工具消费一般需要额外开发接口团队协作成本低实际项目里还有一个隐藏优势Flink SQL的表结构和Kafka中的消息格式是一一对应的改动字段时只要改SQL脚本不需要重新编译Java代码。这对需求频繁变化的实时统计项目来说太重要了。当然DataStream API也不是完全没有出场机会。比如要做复杂的CEP复杂事件检测、要对数据做非常规的清洗、要和外部存储做复杂交互还是得写流式代码。绝大多数实时统计场景用Flink SQL就够了。1.3 一条数据从Kafka到结果的完整流转我们这个项目选的技术组件很扎实数据源是Kafka计算引擎是Flink结果可以输出到控制台、MySQL或者Elasticsearch。为了后面讲起来不混乱先把数据流理清楚。整个流程是业务系统或者模拟程序产生JSON格式的日志消息发送到Kafka的user_behavior这个topic。Flink SQL任务作为Kafka消费者按topic读取消息将JSON解析成一张逻辑表。然后SQL引擎根据我们声明的Watermark处理乱序数据按事件时间切分钟级滚动窗口对窗口内的数据做COUNT、SUM等聚合计算。最后把每个窗口的统计结果写出去。我建议第一次做的时候结果先打到print sink也就是直接打印到控制台和日志里确认整个链路没问题之后再改成MySQL或者ES。这样做的好处是一旦没有输出或者数据不对你最先怀疑的应该是Kafka源表配置、时间字段、Watermark这些上游环节而不是存储端的问题。实际排查起来会快很多。2. 环境准备与依赖搭建2.1 本地开发环境清单开始写代码之前先把环境准备好。不要在一个不熟悉的环境里边踩环境坑边写业务那样你会分不清到底是自己SQL写错了还是Flink启动不了。我自己常用的本地开发环境是这样一套JDK 8或者JDK 11都可以跑Flink建议至少8u201以上版本IDEA或者你熟悉的任何Java IDEMaven 3.6以上Kafka我建议用2.8以上版本新版本对KRaft模式支持更好不强制依赖ZooKeeperFlink严格来说本地开发不需要安装Flink因为我们在IDEA里把Flink作为依赖引入后直接跑main方法它会自动启动一个本地嵌入式的Flink运行时。如果你想用SQL Client或者提交到集群那才需要单独下载Flink发行包。版本匹配是新手最容易忽略的坑。Flink核心版本和Flink连接器版本是分开发布的比如Flink 1.17.1对应的flink-connector-kafka版本是3.0.2-1.17。你如果直接把Flink版本号拿来当连接器版本号Maven会报错或者拉到错误依赖。我建议直接用一套已知稳定的组合Flink 1.17.1 Kafka 3.4.0 Java 8这套组合我实测过本地和简单集群场景都没问题。如果你所在的团队有统一版本以团队规定为准只要记得连接器版本独立这件事就行。2.2 Maven依赖配置与版本匹配细节Maven工程里的依赖是第一个容易踩坑的地方。Flink SQL项目通常需要引入下面几个核心依赖flink-table-api-java-bridge、flink-table-planner-loader、flink-clients、flink-connector-kafka、flink-json。我给你一个可以直接用的pom片段。properties flink.version1.17.1/flink.version kafka.version3.4.0/kafka.version maven.compiler.source1.8/maven.compiler.source maven.compiler.target1.8/maven.compiler.target /properties dependencies dependency groupIdorg.apache.flink/groupId artifactIdflink-table-api-java-bridge/artifactId version${flink.version}/version /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-table-planner-loader/artifactId version${flink.version}/version /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-clients/artifactId version${flink.version}/version /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-kafka/artifactId version3.0.2-1.17/version /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-json/artifactId version${flink.version}/version /dependency /dependencies这里有两个细节你们要注意。第一flink-table-planner-loader是必须的。只引入flink-table-api-java-bridge而不引入planner运行时会报找不到执行计划的工厂类这是个非常典型的错误。我见过不少人在这一步卡了很久其实只是少了一个依赖。第二flink-connector-kafka的版本号不是${flink.version}而是连接器自己独立的版本号。以Flink 1.17.x为例对应的是3.0.2-1.17。如果你用的是Flink 1.16那就找2.4.1-1.16这种对应版本。偷懒的办法是改成引入flink-sql-connector-kafka版本号和Flink主版本同步比如1.17.1但这个包会把Kafka客户端等依赖shade进去体积更大适合提交到集群的场景。在IDEA本地写Java代码调试用flink-connector-kafka更轻量。2.3 本地Kafka快速启动与常用命令Kafka我习惯用二进制包方式启动因为可控、方便看日志。如果你是Windows系统注意把命令换成bin\windows\下的.bat文件。新版Kafka支持KRaft模式可以不依赖ZooKeeper直接启动命令如下# 生成集群ID KAFKA_CLUSTER_ID$(bin/kafka-storage.sh random-uuid) # 格式化存储目录 bin/kafka-storage.sh format -t $KAFKA_CLUSTER_ID -c config/kraft/server.properties # 启动Kafka bin/kafka-server-start.sh config/kraft/server.properties如果你是旧版本的Kafka那就走传统路线先起ZooKeeper再起Kafka。命令分别是bin/zookeeper-server-start.sh config/zookeeper.properties和bin/kafka-server-start.sh config/server.properties。逻辑一样就是多一个依赖进程。Kafka起来后创建本项目需要的topic。分区数我建议至少3个这样后边测试Flink并行度和Kafka分区的对应关系时感受更明显。bin/kafka-topics.sh --create \ --topic user_behavior \ --bootstrap-server localhost:9092 \ --partitions 3 \ --replication-factor 1然后验证一下topic是否创建成功bin/kafka-topics.sh --describe --topic user_behavior --bootstrap-server localhost:9092Kafka启动没问题之后我们后边灌测试数据会用到kafka-console-producer.sh排查问题会用到kafka-console-consumer.sh这三个命令基本能覆盖整个项目的数据准备工作。3. 核心实现用Java代码嵌入Flink SQL3.1 创建执行环境与Checkpoint配置我们的程序入口是一个标准的Java main方法。第一步不是写业务而是创建Flink的执行环境和Table环境。import org.apache.flink.configuration.Configuration; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.table.api.EnvironmentSettings; import org.apache.flink.table.api.bridge.java.StreamTableEnvironment; public class KafkaFlinkSqlJob { public static void main(String[] args) { // 本地调试时可以设置并行度避免打印结果过于分散 Configuration configuration new Configuration(); configuration.setString(taskmanager.numberOfTaskSlots, 4); StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(configuration); // 开启Checkpoint生产环境必须本地调试也建议开 env.enableCheckpointing(10000); EnvironmentSettings settings EnvironmentSettings.newInstance() .inStreamingMode() .build(); StreamTableEnvironment tableEnv StreamTableEnvironment.create(env, settings); } }这里重点说一下为什么用StreamTableEnvironment而不是TableEnvironment。StreamTableEnvironment是流批一体的Table API入口能同时处理流式和批式任务而纯TableEnvironment更多是批处理场景使用。我们做实时统计数据源是Kafka这种无界流必须用StreamTableEnvironment。Checkpoint在本地调试时容易被忽略但它对Kafka来源特别重要。Flink读Kafka时会周期性向Kafka提交offset但这个提交并不是实时消费了就立刻提交而是等Checkpoint完成之后由一个回调去提交。如果你不开Checkpoint任务停止后重启可能会从头消费或者重复消费数据统计就会不准。本地随手加一行env.enableCheckpointing(10000)能省掉后面很多莫名其妙的重复消费问题。3.2 创建Kafka源表字段、Watermark与连接参数执行环境创建好之后接下来就是Flink SQL的重头戏——用DDL创建一张映射到Kafka topic的逻辑表。这张表建得好不好直接决定后面统计SQL能不能跑对。CREATE TABLE user_behavior ( user_id STRING, action STRING, amount DOUBLE, ts TIMESTAMP(3), WATERMARK FOR ts AS ts - INTERVAL 5 SECOND ) WITH ( connector kafka, topic user_behavior, properties.bootstrap.servers localhost:9092, properties.group.id flink-group-001, scan.startup.mode earliest-offset, format json, json.ignore-parse-errors true );在Java代码里执行的话直接用tableEnv.executeSql(...)把这条DDL跑一遍即可。我来逐项解释几个关键设计。字段类型映射上user_id是STRINGaction是STRINGamount是DOUBLEts是TIMESTAMP(3)。这里TIMESTAMP(3)的3代表毫秒精度。JSON里的时间字段如果是字符串格式得能被Flink解析默认期望的是yyyy-MM-dd HH:mm:ss这种格式。如果你的时间字段是Unix时间戳比如1710000000000这种13位毫秒值建议先把字段定义成BIGINT再在SQL里用TO_TIMESTAMP_LTZ函数转换成时间类型Watermark再建在这个转换后的字段上。否则直接拿BIGINT字段定义Watermark会报类型错误。Watermark这行是整个DDL的灵魂。WATERMARK FOR ts AS ts - INTERVAL 5 SECOND的意思是允许事件时间最多乱序5秒Flink会持续跟踪已经接收到的最大ts值然后减去5秒作为当前水位。窗口触发条件就是watermark推进到窗口结束时间。这个5秒不是拍脑袋定的它取决于上游消息乱序程度。乱序比较少1秒、2秒都行乱序严重可能需要10秒甚至更多。但Watermark设置越大窗口结果输出越延迟生产环境需要根据业务容忍度来调。WITH参数里scan.startup.mode earliest-offset表示从头开始消费这个设置对测试非常友好。如果你改成latest-offsetFlink只消费启动之后到达的数据而你启动之前Kafka里已经存在的数据全部被忽略这会导致很多新手死活看不到输出。实际调试期我建议直接用earliest-offset等确认结果正确了再根据需求调整。json.ignore-parse-errors true也很实用。生产环境里Kafka topic中偶尔会有脏数据设置了它任务不会因为个别消息解析失败而崩溃。不过代价是脏数据会被静默丢弃所以线上使用时要结合监控告警看丢弃了多少数据。3.3 定义结果表与核心统计SQL源表建好之后再建一张结果表。本地调试阶段最省事的Sink是print它会把每一条计算结果直接打到TaskManager的标准输出在IDEA的控制台就能直接看到。CREATE TABLE result_sink ( window_start TIMESTAMP(3), window_end TIMESTAMP(3), action STRING, cnt BIGINT, total_amount DOUBLE ) WITH ( connector print );print sink对字段名没有严格约束不管字段叫什么它都会以I开头的变更日志格式打印出来。这里定义五个字段是为了让最终结果清晰对应窗口起止时间、行为类型、次数和金额。接下来是核心的统计SQL。我按分钟统计每种action的发生次数和金额合计。Flink 1.14之后推荐使用Table Valued Function的窗口写法这也符合官方对未来窗口扩展的规划INSERT INTO result_sink SELECT window_start, window_end, action, COUNT(*) AS cnt, SUM(amount) AS total_amount FROM TABLE( TUMBLE(TABLE user_behavior, DESCRIPTOR(ts), INTERVAL 1 MINUTE) ) GROUP BY window_start, window_end, action;这条SQL的语义可以拆成三步理解。第一步TUMBLE(TABLE user_behavior, DESCRIPTOR(ts), INTERVAL 1 MINUTE)表示按事件时间ts将数据划分成1分钟的滚动窗口。第二步GROUP BY window_start, window_end, action表示每个行为类型独立统计。第三步COUNT(*)统计条数SUM(amount)统计金额。老版本写法是TUMBLE_START(ts, INTERVAL 1 MINUTE)配合GROUP BY TUMBLE(ts, INTERVAL 1 MINUTE), action。这种写法在新版本中也能跑但官方推荐用TVF方式因为TVF能让窗口开始时间、结束时间直接变成可查询字段语义更清晰。我建议新项目直接学TVF写法网上很多文章还在用老语法你看到时注意分辨。3.4 执行作业与提交方式SQL都定义好之后用tableEnv.executeSql(...)执行DDL和DML语句。注意执行顺序不能乱必须先建源表、再建结果表、最后执行INSERT。所有语句准备好后调用env.execute(kafka-flink-sql-job)真正启动任务。完整Java代码跑起来之后你会发现Flink在IDEA里直接启动了一个MiniCluster任务就在本地线程里运行。这种模式非常适合调试日志直接输出到控制台断点也能正常打断。项目打包提交到集群的方式稍微不同。打jar包时要注意别漏依赖。我习惯用maven-shade-plugin把项目依赖打成一个fat jarplugin groupIdorg.apache.maven.plugins/groupId artifactIdmaven-shade-plugin/artifactId version3.4.1/version executions execution phasepackage/phase goals goalshade/goal /goals configuration transformers transformer implementationorg.apache.maven.plugins.shade.resource.ManifestResourceTransformer mainClassKafkaFlinkSqlJob/mainClass /transformer /transformers /configuration /execution /executions /plugin打包后执行flink run命令flink run -c KafkaFlinkSqlJob target/kafka-flink-sql-1.0-SNAPSHOT.jar如果只是快速验证SQL也可以不写Java代码直接启动Flink的SQL Client把刚才那段DDL和INSERT粘贴进去执行。SQL Client本质上就是把我们Java代码里做的那些事用命令行方式暴露出来非常适合做原型验证。4. 实操过程与验证调试4.1 准备测试数据往Kafka灌数据程序写好之后最激动也最容易出问题的环节来了往Kafka里发数据。很多人习惯一上来就写一个生产者程序其实完全没必要Kafka自带的kafka-console-producer命令足够用来拆数据链路。我通常按下面这个顺序造数据。先启动一个生产者bin/kafka-console-producer.sh --topic user_behavior --bootstrap-server localhost:9092然后逐条输入JSON消息。注意每输入一行按一次回车Kafka的producer客户端会把每行当成一条独立的消息发送。{user_id:1001,action:click,amount:1.0,ts:2024-06-01 10:00:01} {user_id:1002,action:click,amount:1.0,ts:2024-06-01 10:00:05} {user_id:1003,action:pay,amount:199.0,ts:2024-06-01 10:00:10} {user_id:1004,action:pay,amount:88.5,ts:2024-06-01 10:00:30} {user_id:1005,action:click,amount:1.0,ts:2024-06-01 10:00:50}这些数据分布在10:00:00到10:00:59之间正好是同一个分钟窗口。如果Flink按1分钟滚动窗口计算这一批消息应该输出两条结果click出现3次共3.0元pay出现2次共287.5元。发送完数据后建议先用kafka-console-consumer快速确认消息真的进了Kafka再跑Flink程序避免出现“Flink没读到数据”和“数据压根没进Kafka”分不清楚的情况。bin/kafka-console-consumer.sh \ --topic user_behavior \ --bootstrap-server localhost:9092 \ --from-beginning4.2 运行结果演示与逐行解读测试数据发好之后启动我们的Java程序。如果一切正常你会在控制台看到类似下面这样的输出I[2024-06-01T10:00:00, 2024-06-01T10:01:00, click, 3, 3.0] I[2024-06-01T10:00:00, 2024-06-01T10:01:00, pay, 2, 287.5]I表示这是一条插入操作记录。方括号里按顺序对应我们定义的结果表字段窗口开始时间是10:00:00窗口结束时间是10:01:00click这个动作出现了3次总金额3.0元pay出现了2次总金额287.5元。这里我故意把第一条消息的时间定在10:00:01就是为了让窗口边界看起来更直观窗口属于左闭右开区间10:00:00开始10:01:00之前的数据都会落入这个窗口。很多人第一次跑可能看不到这个输出不用慌大概率是Watermark还没推进到窗口触发点。Flink的滚动窗口并不是源源不断输出结果的而是要等到watermark越过窗口结束时间才会触发计算。我们设置Watermark为最大事件时间减5秒所以要等到Flink看到时间戳在10:01:05之后的数据才会触发10:00到10:01这个窗口的计算。这种机制保证了窗口结果不会因为少量乱序数据而反复更新但也意味着窗口结果天然有最多5秒的“延迟可见性”。如果你用处理时间替代事件时间输出会快很多但数据统计的是Flink处理消息的时间而不是消息本身携带的业务时间在实时统计场景里通常不可接受。4.3 从控制台输出升级为MySQL结果落库print输出验证通过之后就可以把结果接到真正的存储里了。最常见的是写MySQL用来支撑实时报表或者大屏项目。首先在MySQL里建一张结果表CREATE TABLE action_stat ( window_start DATETIME(3) NOT NULL, window_end DATETIME(3) NOT NULL, action VARCHAR(32) NOT NULL, cnt BIGINT NOT NULL, total_amount DOUBLE NOT NULL, PRIMARY KEY (window_start, window_end, action) );然后Flink侧把结果表的connector改成jdbcCREATE TABLE mysql_sink ( window_start TIMESTAMP(3), window_end TIMESTAMP(3), action STRING, cnt BIGINT, total_amount DOUBLE, PRIMARY KEY (window_start, window_end, action) NOT ENFORCED ) WITH ( connector jdbc, url jdbc:mysql://localhost:3306/rtdash?useSSLfalseserverTimezoneAsia/Shanghai, table-name action_stat, username root, password 123456 );JDBC sink默认使用upsert语义当主键冲突时会执行更新而不是插入。这里PRIMARY KEY后边的NOT ENFORCED是Flink特有的语法意思是字段上声明了主键但Flink不会强制去数据库里校验唯一性它只是告诉Flink这个表有主键方便做更新删除操作。MySQL表主键和Flink表主键必须一一对应。对应的依赖也需要加进pomdependency groupIdorg.apache.flink/groupId artifactIdflink-connector-jdbc/artifactId version${flink.version}/version /dependency写MySQL时的执行方式与print一样把目标表从result_sink换成mysql_sinkINSERT语句不用动。实际项目里如果要做实时大屏通常会把结果写成JSON格式再发回Kafka或直接写Elasticsearch再由前端查询展示原理完全一样只是换一个Sink连接器而已。5. 常见问题与排查技巧实录5.1 从Kafka拉数据报错Error while fetching metadata这个报错几乎每个做Kafka开发的人都遇到过原文大概长这样org.apache.kafka.common.errors.TimeoutException: Topic user_behavior not present in metadata after 60000 ms.外加一句Error while fetching metadata with correlation id。看到这个不要慌原因无非下面几种。Kafka没启动是最常规的原因。先检查一下Kafka进程是否还在再用kafka-console-consumer手动消费一次topic如果手动消费也失败那就不用怀疑Flink了先把Kafka侧的连通性问题解决掉。bootstrap.servers配置不对也很常见。本地调试时是localhost:9092如果Flink任务跑在Docker容器里就要写成宿主机的IP或者Kafka容器的服务名。你如果Docker部署的Kafka还要特别注意advertised.listeners配置broker注册到集群的地址如果写的是localhost外部客户端拿到这个地址后依然连不上。这个问题的典型表现是kafka-console-producer在本机能发消息但Flink在容器里消费报metadata超时。还有一种是topic不存在Kafka自动创建又没开启。检查auto.create.topics.enable配置手动创建topic是最保险的办法。排查顺序我建议这样先看topic能不能用console消费再用客户端连接bootstrap.servers测试元数据获取最后才回头看Flink的WITH参数。5.2 作业跑起来但控制台没有任何输出这是实时统计里最让人头疼的问题。Kafka里有数据Flink任务也一直running但控制台就是静悄悄。我总结下来百分之八十的情况出在三个地方。第一个是消费位点。如果scan.startup.mode配置成latest-offsetFlink只会读启动之后的数据。你先发数据再启动Flink把历史数据全错过了当然没输出。测试期统一改成earliest-offset或者先启动Flink再发数据。第二个是Watermark没有推进。窗口聚合的结果是等Watermark越过窗口结束时间才触发的你的数据如果时间字段设置得不对比如ts字段字符串格式不能被解析Watermark就一直不涨。可以在源表DDL里把ts类型改成STRING先看原始数据或者临时把Watermark改成ts也就是延迟0秒验证是不是Watermark的问题。第三个是数据压根就没解析成功。比如JSON里字段名和DDL对不上导致user_id、amount都是null。这种情况下Flink不会报错SQL执行时null参与聚合也不会崩溃count可能还是照常统计但sum的结果就会很奇怪。所以排查时不要只盯着有没有输出要把输出内容也过一遍。我的调试习惯是先确认Kafka里有符合格式的数据再开print sink用最小的窗口和最小的Watermark把链路打通最后再一步步加上真实业务参数。链路通没通是第一步数据准不准是第二步两个问题不要混在一起排查。5.3 结果数据不准确金额为0、行数不对如果程序有输出但数据明显不对问题基本都在数据解析和事件时间语义上。金额一直为0优先检查amount字段的JSON类型。如果上游传的是字符串199.0而DDL定义的是DOUBLEFlink的JSON格式在严格模式下会解析失败。虽然设置了json.ignore-parse-errors true会让任务不挂但这条数据会被丢弃统计结果自然就少了。排查时先关掉这个参数让解析错误暴露出来比对着结果瞎猜效率高得多。行数偏多或偏少多半是乱序和迟到数据问题。我们的Watermark只允许5秒乱序超过5秒的迟到数据默认会被丢弃。如果测试时故意发一条10:00:55的数据再发一条10:01:10的数据理论上10:00窗口能正常出结果但那条55秒的数据如果在10:01:06之后才到就会被丢掉。生产环境如果业务上完全不能接受丢数据可以考虑把迟到数据发送到侧输出流里用另外一套计算逻辑补偿但这属于进阶用法第一次做实时统计先不要引入这么复杂的语义。另一个导致行数不对的常见原因是任务重启导致重复消费。前面说过不开Checkpoint时Flink不会提交offset每次重启都从头消费。本地开发过程中程序重启个十几次太正常了统计结果里那批旧数据就会被反复计算。这个问题不解决你会陷入“为什么我代码改了一行结果翻倍了”的迷茫。开启Checkpoint是最直接的解决办法。5.4 依赖与版本冲突排查本地运行最常见的一类异常是NoClassDefFoundError。比如跑SQL时提示找不到org.apache.flink.table.api内部类十有八九是flink-table-planner-loader没有引入。再比如报和Jackson相关的错误就要看看项目里是不是有其他版本的Jackson被Maven解析进来了。排查依赖冲突的标准动作是跑一下mvn dependency:tree看看有没有同一个groupId下多个版本的包。Flink的shaded机制已经解决了很多类冲突但JDBC驱动、Kafka客户端这些第三方库还是可能和项目里的其他组件打架。我吃过一次亏项目里原本有一个老版本的kafka-clients 2.2.0Flink连接器用的是kafka-clients 3.4.0Maven仲裁选择了老的结果运行时报了一堆找不到方法的异常。最后把旧依赖排除掉只在pom里显式声明了3.4.0才解决。如果你不想处理这些破事有两个省心的替代方案。第一直接用flink-sql-connector-kafka替换flink-connector-kafka它自带shade依赖冲突面小很多。第二提交作业到集群时尽量利用Flink发行包自带的lib目录把连接器jar直接丢进去避免和应用jar打包到一起。5.5 关于Kafka消费延迟与消费组调整很多人在排查实时统计问题时还会顺手看一下Kafka消费延迟。用kafka-consumer-groups命令可以查到lagbin/kafka-consumer-groups.sh \ --bootstrap-server localhost:9092 \ --describe \ --group flink-group-001这个命令会列出当前消费组在各个分区上的当前offset、log end offset和lag。如果你发现lag越来越大说明Flink消费速度跟不上生产速度优先考虑调大并行度。但要注意Flink读Kafka的并行度最多等于Kafka分区数一个分区不可能被两个并行子任务同时消费。所以topic分区数如果只有3个并行度调成8也没用该升的是Kafka分区数。关于group.id我有一个特别想提醒的坑多个Flink任务不要共用同一个group.id。Flink通过group.id管理offset如果两个任务共用会发生频繁rebalance表现为消费一会儿停一会儿统计结果稀稀拉拉。更危险的是一个任务重启后可能抢走另一个任务的offset导致两边都消费错乱。给每个任务一个独立的group.id是基本素养。有人会在网上搜“Kafka如何延迟30分钟消费”如果你也有这个需求最优雅的做法不是去控制Kafka消费节奏而是在Flink SQL的Watermark和窗口设计上做文章或者通过事件时间过滤实现类似效果。直接修改消费者组offset延迟也是可以的但可维护性很差生产环境不建议这么搞。6. 最后再分享几个保命习惯这套Kafka加Flink SQL的方案你照着前面的步骤跑通一遍之后基本就掌握核心玩法了。最后我再聊几个自己从实战里沉淀出来的习惯它们救过我很多次。第一个习惯是调试期永远先跑通最小链路。我第一次做的时候直接就把窗口设成1小时Watermark设成30秒结果数据发进去半天没输出我一度以为是Kafka连不上。后来改成10秒窗口、5秒延迟整个链路几秒钟就有结果排查效率立刻上来了。窗口参数设小只是为了验证链路跑通后再改回业务要求的参数这是最稳妥的节奏。第二个习惯是Watermark从宽到严慢慢收。先设置一个比较大的乱序容忍度比如30秒确保数据基本都能进窗口、结果稳定可预期。等业务理解了迟到的现象后再根据真实数据乱序程度慢慢收紧到5秒、3秒。一上来就设1秒可能漏掉很多业务上该统计的数据到时候对账对不上非常痛苦。第三个习惯是任何写入外部存储之前先打印out到控制台。不要嫌麻烦print sink是Flink给你的免费礼物。控制台输出是一个纯副作用操作不依赖外部系统看到输出就证明Flink内部逻辑正确然后再接MySQL、接ES。外部系统一旦出问题你至少能分清楚到底是Flink算得不对还是存储侧写不进去。第四个习惯是保存一份可直接执行的SQL脚本。即使我们是用Java代码嵌入SQL也最好把DDL和DML独立成sql文件放进项目里。因为业务变更时很可能只是加一个字段、改一个窗口大小这时候直接用SQL Client验证一下效率极高根本不用重新编译Java代码。Flink SQL Gateway这类工具也在支持这种工作方式未来你团队里的数据分析师可能不需要写Java也能查实时数据。如果你把这套跑通了后面可以做很多扩展换上不同的Kafka topic、做多维度GROUP BY、接上Canal做数据库实时同步统计、把结果同时写到Kafka和ES支撑大屏。实时统计这扇门用Flink SQL敲开之后你会发现后面其实是一片大工地可以盖的东西太多了。
返回列表