ARTICLE DETAIL

资讯详情

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

Flink SQL生产实战:从DataStream重构到ClickHouse同步

Flink SQL生产实战:从DataStream重构到ClickHouse同步 年初接手实时数仓项目时我把自己之前用DataStream API写的一套实时指标任务全部重构了一遍。重构完成后代码量从两千多行降到四百多行开发周期从两周缩到三天。让我下定决心做这件事的正是Flink SQL。这篇不是官方文档复读机而是我在生产环境里把Flink SQL从选型、建模、同步到ClickHouse再到排查JDBC连接器异常的一整套实战记录。如果你已经会用DataStream API但始终没把Flink SQL用起来或者刚接触大数据实时计算、想找一份能落地参照的案例这篇文章值得你花二十分钟读完。1. 先用SQL去写流处理到底图什么很多人一听到“Flink SQL”第一反应是“这不就是把流数据查一下吗哪有DataStream API灵活”。这个想法不算错但生产级的灵活并不体现在“所有事情都手写”上。我见过太多团队用DataStream API写了大量窗口聚合、双流关联、状态管理最后线上出问题排查一堆时间找不出原因。SQL的引入本质上是把“你关心什么”和“怎么实现”分开。1.1 从DataStream API到SQL阈值在哪先说清楚边界。如果任务只是简单的读Kafka、过滤、写MySQLDataStream API一点问题没有代码也短。但一旦涉及多流关联、滚动窗口、会话窗口、维表补全、迟到数据处理你要手动处理的细节就会指数级上升timers怎么注册、状态该存在哪个key、watermark怎么在算子和窗口之间传递、迟到数据能不能被兜住。这些逻辑在SQL里就是几句声明Flink的查询优化器会负责生成执行计划、调整join顺序、选择合适的state backend。我自己的判断标准是任务里只要出现两个join以上或者有按时间维度的聚合那用DataStream API手写就划不来了。SQL虽然不能在每个细节上精准控制但稳定性和可维护性要远高于手写状态代码。实时任务最贵的不是CPU而是“半夜报警后花两小时查逻辑”。1.2 Flink SQL能替代什么替代不了什么能替代的场景非常明确实时数据清洗和格式转换JSON解析、类型转换、异常值过滤。实时指标计算按分钟、小时进行PV、UV、GMV统计。双流关联订单流和支付流、订单流和发货流。CDC数据同步把MySQL binlog转成流再写入其他存储。维表关联用户维度、商品维度补全用lookup join搞定。替代不了的场景也有需要自定义state结构的复杂风控规则比如一台设备上五分钟内关联了多个账号需要用图结构去存。强依赖自定义窗口语义比如窗口不是按时间闭区间而是要按业务里的“一单完结”事件触发。对背压和checkpoint有极致控制要求的低延迟场景。不过这些场景占比很小。我所在的业务线实时任务里大概80%都在SQL能力范围内。剩下20%才拆出来用DataStream API。1.3 实时数仓分层里的SQL位置Flink SQL最常见的落地方式是实时数仓分层。OSS层日志先进KafkaSQL任务消费Kafka做解析落到Hive或者IcebergDWD层用SQL做清洗、维表关联、双流JOINDWS层用SQL做分钟级、小时级聚合结果写ClickHouse或Doris。这条链路里除了最底层的数据接入偶尔要写SourceFunction剩下全是SQL。举个例子我们有个“用户点击行为实时分析”的需求原链路是DataStream API读取点击日志 → 做IP解析 → 开三分钟窗口 → 去重 → 聚合成指标。这一套下来要写四百多行Scala。换成SQL之后建一张点击日志源表、一张IP维表、一张目标结果表中间加一条CREATE VIEW和一条INSERT INTO就结束了。代码少了粒度还更清晰后面同事接手也更容易理解。2. 搭一套能跑的Flink SQL开发环境Flink SQL看起来只是写几句SQL但真要落到自己的项目里环境搭建和依赖选择是第一个坑。版本选错、依赖缺包、作业提交方式不合适都会让你在SQL还没跑起来前就消耗掉热情。2.1 版本不要乱选先对齐生态我推荐直接从Flink 1.17或1.18起步Java环境用8或11都行但尽量用11。原因很简单官方连接器和Table API在1.13之后才进入稳定期1.15之后语法和配置项与现代用法比较接近再老的版本很多特性缺不说社区资料也少。这里列一份我常用的依赖清单你直接在Maven里配好依赖作用flink-table-api-javaJava版的Table API和SQL APIflink-table-planner-loader内置查询优化器作业运行时需要flink-table-runtimeSQL执行时的运行时组件flink-connector-kafkaKafka Source/Sink连接器flink-connector-jdbcJDBC Sink连接器mysql-connector-javaMySQL驱动clickhouse-jdbcClickHouse驱动有一件事特别提醒Flink的flink-table-planner和flink-table-runtime如果单独引入版本必须和Flink主版本完全一致。我遇到过很多次本地能跑、提交集群就报NoClassDefFoundError十有八九是planner-loader和表的运行时依赖没对齐。2.2 代码写SQL还是SQL客户端写SQL我本人习惯先用Flink SQL客户端做语法验证再把SQL以字符串形式固化到Java代码里。这样既能快速调试又能把最终SQL纳入版本控制。如果只是临时分析直接启动sql-client.sh写CREATE TABLE和INSERT INTO试试就行。但生产任务建议都写成Java代码因为你多半要配合自定义UDF、运行时参数和状态恢复机制。用Java代码写SQL有一个好处可以用StreamTableEnvironment.executeSql()一条一条提交逻辑清晰也方便做多级ETL的编排。2.3 一个最小可跑的Java任务下面这个是我每次新建Flink SQL项目时的起点读Kafka、打印结果验证环境和依赖没问题后再往里面加业务StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); EnvironmentSettings settings EnvironmentSettings.inStreamingMode(); StreamTableEnvironment tEnv StreamTableEnvironment.create(env, settings); tEnv.executeSql( CREATE TABLE source_event (\n user_id BIGINT,\n event_type STRING,\n event_time BIGINT,\n ts AS TO_TIMESTAMP_LTZ(event_time, 3),\n WATERMARK FOR ts AS ts - INTERVAL 5 SECOND\n ) WITH (\n connector kafka,\n topic user-event,\n properties.bootstrap.servers localhost:9092,\n properties.group.id flink-sql-demo,\n scan.startup.mode latest-offset,\n format json\n ) ); tEnv.executeSql( CREATE TABLE print_sink (\n user_id BIGINT,\n event_type STRING,\n cnt BIGINT\n ) WITH (\n connector print\n ) ); tEnv.executeSql( INSERT INTO print_sink\n SELECT user_id, event_type, COUNT(*)\n FROM source_event\n GROUP BY user_id, event_type ); env.execute(flink-sql-demo);这段代码跑起来后你会看到控制台每秒输出聚合结果。等这一步通了再往上加真实业务就不容易产生“环境问题”。3. 把Kafka数据接进来动态表与连续查询第一次真正理解Flink SQL和普通数据库SQL的区别是在排查一个聚合任务为什么一直输出不完整时。普通SQL查的是静态的数据集查完就结束了Flink SQL查的是一张永远在增长的“动态表”查询也永远不会停止除非作业被取消。这个本质差异决定了建表、写查询、调状态的整个思路。3.1 流和表到底怎么映射Kafka里一条一条消息本质上是一个无限增长的流。Flink SQL把它看作动态表插入一条消息相当于往表里加一行删除或者更新一条也对应表的变更。动态表有三种模式Append-only只插入不更新。适合日志、事件流。Retract既可以插入也可以删除。适合COUNT、SUM这类会产生回撤的聚合。Upsert按主键更新。适合维表、订单状态表。这点非常关键。比如你执行SELECT user_id, COUNT(*) FROM source_event GROUP BY user_id因为输入流不断有新数据同一个user_id的COUNT(*)会不断变大Flink SQL就需要通过Retract消息把上一轮的旧结果作废再输出新结果。如果你用print sink调试会看到先输出一个“-1”再输出“2”这代表撤回和新增。3.2 建表语句里的格式与时间字段处理Kafka里的数据一般是JSON字符串。建表时除了指定connector和topic最容易被忽略的是时间字段。比如原始日志里事件时间是毫秒级时间戳CREATE TABLE user_event ( user_id BIGINT, event_type STRING, event_time BIGINT, ts AS TO_TIMESTAMP_LTZ(event_time, 3), WATERMARK FOR ts AS ts - INTERVAL 5 SECOND ) WITH (...)这里ts AS TO_TIMESTAMP_LTZ(event_time, 3)是把BIGINT类型的事件时间转成TIMESTAMP_LTZWATERMARK则告诉Flink最多容忍5秒乱序。没有这一步窗口聚合基本跑不出准确结果。还有一个小提示Kafka JSON格式解析时尽量加json.ignore-parse-errors true。生产环境里上游同事偶尔会发一条脏数据如果不开忽略解析错误整个作业会因为一条坏消息直接失败然后反复重启。3.3 连续查询不是“跑一次就结束”普通SQL对静态表查一次返回结果就完了。Flink SQL的连续查询会一直运行只要输入表有数据进来就会重新计算并更新结果。这就意味着状态会持续累积。拿GROUP BY为例Flink必须把每个key的当前存量一直保存在状态里。如果条件不合适比如你按用户ID分组用户量上亿状态就会非常大。好的一点是Flink有状态生存时间配置你可以在建表或者作业配置里设置table.exec.state.ttl超过时间没更新的key自动清理防止状态无限增长。我在生产环境常用的设置是SET table.exec.state.ttl 2 h;这里的2小时需要结合业务合理设置。如果窗口聚合是1小时周期那2小时TTL足够如果业务需要精确去重且去重周期为7天TTL就必须不小于7天。4. 生产最常用的三类Flink SQL场景过滤、聚合、双流JOINFlink SQL真正产生价值的地方还是那些在实时数仓里高频出现的场景。我总结下来就三类清洗过滤、窗口聚合、流与流的关联。把这三类吃透你就能覆盖绝大多数实时指标需求。4.1 清洗与字段治理实时清洗不是只做WHERE这么简单。你经常要处理类型不一致、字段缺失、枚举值混乱甚至同一个日志里有新旧两版格式并存。一种常见写法是用CASE WHEN做标准化INSERT INTO dwd_event SELECT user_id, LOWER(COALESCE(event_type, unknown)) AS event_type, TO_TIMESTAMP_LTZ(event_time, 3) AS event_ts, CASE WHEN platform ios OR platform IOS THEN iOS WHEN platform android THEN Android ELSE other END AS platform FROM source_event WHERE user_id IS NOT NULL AND event_time 0;注意COALESCE只能把NULL转成默认值如果是空字符串还需要单独处理。另外如果源Kafka的JSON字段和表字段类型对不上Flink SQL默认会直接报错这时候可以在WITH里加formatjson配合json.fail-on-missing-fieldfalse允许缺失字段能让你省下很多折腾解析的时间。4.2 窗口聚合的语法陷阱窗口聚合是实时指标最常见的操作。Flink SQL里最常用的滚动窗口是TUMBLESELECT TUMBLE_ROWTIME(ts, INTERVAL 5 MINUTE) AS window_time, event_type, COUNT(*) AS cnt FROM source_event GROUP BY TUMBLE(ts, INTERVAL 5 MINUTE), event_type;这里有个很容易踩的语法坑GROUP BY里必须同时带上TUMBLE(...)否则Flink不知道数据要按哪个时间列开窗。窗口函数返回的window_time有特殊的类型写TUMBLE_ROWTIME或TUMBLE_END都可以但别把两者混用否则结果表里的时间字段含义会不一致。事件时间的乱序处理也值得单独说。如果你在源表上定义了5秒的watermark那么在窗口末尾5秒内到达的数据会先触发窗口计算之后又来一批5秒前的事件Flink不会重新更新已关闭的窗口结果。想要处理更晚的数据需要额外用侧输出流或ALLOWED LATENESS语法但生产环境里为了简单我一般把watermark的容忍范围设成日志到达延迟的90分位值。4.3 双流JOIN怎么避免状态爆炸实时双流JOIN是很多人头痛的问题。两个Kafka流直接JOIN如果没有任何时间边界Flink就必须把所有历史数据都放在状态里数据量一大直接OOM。所以99%的生产任务都会用interval join只关联某个时间窗口内的数据SELECT o.order_id, o.user_id, p.pay_amount FROM orders o JOIN payments p ON o.order_id p.order_id AND p.pay_time BETWEEN o.order_time AND o.order_time INTERVAL 10 MINUTE;这个BETWEEN ... AND ...是核心它告诉Flink只需要为每个订单保留10分钟内的支付记录状态超过之后自动清理。很多实时数仓项目里的“近N分钟活跃用户”“实时转化漏斗”底层其实都是这种interval join。有一点要提醒如果两个流的order_id分布极度不均匀比如某个头部用户占了大量订单JOIN算子仍然会倾斜。处理办法是给订单ID加盐拆散计算或者把热点key单独分流这个后面会详细说。4.4 用UDF补齐SQL覆盖不到的逻辑Flink SQL内建函数很多但不是一切都能覆盖。比如我们业务里要根据IP解析城市、根据UA解析设备这些就得用UDF。做法很简单把函数注册进Table EnvironmenttEnv.createTemporaryFunction(ip_to_city, IpToCityUdf.class);然后在SQL里直接用SELECT ip_to_city(ip) FROM source_event。UDF能极大拓展SQL的边界但注意不要在UDF里做一些阻塞式的远程RPC调用否则会把整个TaskManager线程卡住。如果需要查维表优先用lookup join而不是UDF里调接口。5. 从MySQL同步到ClickHouse一套Sink链路实战实时数仓里最有代表性的落地案例之一就是“MySQL数据实时同步到ClickHouse”。我们的订单表在MySQL运营看板在ClickHouse中间需要毫秒级延迟地搬数据。这套链路用Flink SQL写起来非常清爽。5.1 同步需求拆解订单表有三张订单主表、订单明细表、订单状态流水表。ClickHouse里对应的表结构不一样需要按维度宽表建模。数据同步不能只是简单复制还要把状态流水合并进去按主键更新最新状态。这就涉及三类操作读取MySQL变更、按订单ID关联、写入ClickHouse并更新。5.2 读取MySQL的两条路线路线一是直接用Flink的JDBC连接器轮询查询CREATE TABLE mysql_orders ( id BIGINT PRIMARY KEY, user_id BIGINT, amount DECIMAL(10, 2), status STRING, update_time TIMESTAMP(3) ) WITH ( connector jdbc, url jdbc:mysql://localhost:3306/business, table-name orders, username root, password 123456 );这能做但JDBC source是扫描查询不是binlog级实时变更捕获会有分钟级延迟而且每次扫描全表成本很高。路线二是用Flink CDCCREATE TABLE mysql_orders ( id BIGINT PRIMARY KEY, user_id BIGINT, amount DECIMAL(10, 2), status STRING, update_time TIMESTAMP(3) ) WITH ( connector mysql-cdc, hostname localhost, port 3306, username flinkuser, password 123456, database-name business, table-name orders );CDC方式会直接消费binlog数据是真正的实时变更还能识别数据的新增、更新和删除。我们最终选的是CDC。这里提醒一下mysql-cdc是独立连接器需要在pom里额外引入flink-connector-mysql-cdc不是Flink官方自带。5.3 Sink到ClickHouse用JDBC还是ClickHouse连接器Flink官方其实没有提供ClickHouse连接器社区里有一个flink-clickhouse-connector但成熟度参差不齐。我评估之后选了最稳的JDBC连接器加上ClickHouse官方JDBC驱动。JDBC Sink的建表语句长这样CREATE TABLE clickhouse_order_sink ( id BIGINT, user_id BIGINT, amount DECIMAL(10, 2), status STRING, update_time TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED ) WITH ( connector jdbc, url jdbc:clickhouse://localhost:8123/default, table-name order_total, username default, password , sink.buffer-flush.max-rows 1000, sink.buffer-flush.interval 2s, sink.buffer-flush.max-retries 3 );重点参数是sink.buffer-flush.max-rows和sink.buffer-flush.interval。ClickHouse对批量写入的吞吐要求比较高如果每一条来一条就写一条性能会非常差。我们设置成攒到1000条或2秒刷一次单并发写入速度提升非常明显。ClickHouse实现更新我使用的是ReplacingMergeTree表引擎加复合主键这样即使Flink发送了重复数据ClickHouse也能通过后台合并去重。如果你想用JDBC sink实现真正的精确更新需要自己处理幂等逻辑JDBC连接器本身不保证exactly-once。5.4 完整同步SQL整个同步任务一共分三步定义CDC源表、定义ClickHouse结果表、执行插入INSERT INTO clickhouse_order_sink SELECT id, user_id, amount, status, update_time FROM mysql_orders WHERE status cancel;就这么简单。任务跑起来后MySQL里改一行订单数据ClickHouse里大概一秒内就能查到更新。如果还要把订单明细和订单主表关联成宽表再建一个中间视图把两个源表做interval join之后统一写入ClickHouse即可。5.5 同步延迟和精确性这套链路做不到的是真正意义上的“exactly-once写入ClickHouse”。原因有两个一是JDBC sink对ClickHouse的支持没有内置两阶段提交二是ReplacingMergeTree本身是最终一致。我们通过两个办法把误差缩小到可接受范围写入前在SQL里对主键做DISTINCT降重。ClickHouse表设置ReplacingMergeTree并指定版本字段让重复数据后台自动合并。如果你需要严格不重复建议在ClickHouse写入端再挂一个去重引擎或用外部字典。实时数仓领域里能做到最终一致已经很好了别为了绝对一致把自己困在无法上线的死胡同里。6. 踩坑实录JDBC连接器异常和并行度陷阱看了不少文章大家都爱讲“Flink SQL多方便”但实际生产环境里连接器异常和并行度设置才是真正让人头皮发麻的地方。我把自己线上遇到过的几个高频问题整理出来希望能帮你们少走弯路。6.1 ClassNotFoundException和NoClassDefFoundErrorSQL作业本地可以跑一上集群就报ClassNotFoundException: com.mysql.cj.jdbc.Driver这种是最常见也是最容易解决的。原因是你的Flink集群/lib目录下没有MySQL驱动包只有本地IDE里带了依赖。解决方法是把驱动jar包直接放到TaskManager的lib目录或者提交作业时用-j /path/to/mysql-connector-java.jar显式指定。注意同一个连接器对应的driver版本要一致我曾经在集群lib里放了MySQL 8.x驱动但代码里引的是5.x结果报Communications link failure定位了好几个小时。ClickHouse因为同样用JDBC也要保证clickhouse-jdbc.jar在每台TaskManager上可见。如果用了Flink on YARN部署最容易的处理方式是放在HDFS上然后提交时通过配置把jar分发到所有TaskManager。6.2 Connection is not availablerequest timed out这个错误很吓人我第一反应是ClickHouse数据库挂了结果排查下来发现是连接池被耗尽。Flink SQL的JDBC Sink默认会为每个并行子任务维护一个连接池每个连接池的默认大小有限。如果Sink并行度是32ClickHouse能承受的连接数是50那么所有TaskManager同时写的时候就会有几批任务拿不到连接。解决办法有两个方向调低Sink并行度比如从32降到8减少同时打开的连接数。调整用户侧连接池和ClickHousemax_connections参数。Flink SQL里不能直接给单表设置并行度但你可以在整个作业设置execution.parallelism或者给Sink表单独包装成一条INSERT INTO然后通过PARALLELISM提示来控制。比如INSERT INTO clickhouse_order_sink /* OPTIONS(sink.parallelism8) */ SELECT ...;如果发现仍然不够再配合调大sink.buffer-flush.interval让写入更批量、更平缓。6.3 并行度和Kafka分区数不匹配导致的数据空转Kafka topic有12个分区Flink SQL任务并行度设成了24。Kafka Source最多只有12个并行子任务读数据剩下12个子任务闲置。这个问题不致命但会出现一种“看着有24个并发实际吞吐只有12个分区的能力”的错觉。反过来如果并行度小于分区数则会同时出现数据倾斜和状态膨胀。一般我的做法是Kafka Source并行度等于min(partition数, max(4, 分区数))聚合和Sink并行度再独立设置。遇到热点key就用SQL的DISTRIBUTE BY或GROUP BY次数扩展分两步做聚合。6.4 时区问题导致的日期漂移这个绝对算隐蔽坑。Flink SQL默认的本地时区是UTC如果你的Kafka数据里时间戳是“2025-01-01 12:00:00”这样的本地时间直接转成TIMESTAMP并在窗口聚合你会发现窗口边界全都偏了8小时。解决办法是在作业启动参数里明确设置时区SET table.local-time-zone Asia/Shanghai;或者在Java代码里创建Environment时配置Configuration。时区问题不解决所有按天、按小时聚合的指标都会对不上。这是我见过最容易被忽视、但影响最严重的一个“配置参数”。6.5 Spring Boot整合Flink时的类加载器坑因为业务平台是Java后端很多团队想把Flink SQL任务直接嵌到Spring Boot进程里跑不用独立提交到集群。这在技术上可行但坑非常多。最典型的问题是类加载器冲突Spring Boot用FatJar方式加载类Flink的Table环境中很多类是懒加载一旦StreamTableEnvironment在Spring Bean初始化时被创建后续的动态类加载很可能找不到Flink的planner类导致TableException或ClassCastException。我的经验是如果只是要管理SQL作业不要直接在Spring Boot的Tomcat类加载环境里去跑Flink集群。可以把Spring Boot应用当成一个“SQL任务提交器”只负责把SQL文件内容组装成提交命令真正执行还是交给独立的Flink集群。如果必须在进程内运行建议单独写一个FlinkSqlRunner类在子线程里用隔离类加载器启动不要用Spring注解自动注入StreamTableEnvironment。7. 上线前必看的State与性能调优手段Flink SQL任务在开发环境跑通很容易上线之后能不能扛住高峰期数据量才是考验。以下这些调优点都是我根据线上实例反复调参总结出来的。7.1 Checkpoint配置不能省状态是Flink SQL任务的心脏checkpoint是保证状态不丢的前提。我在生产环境里至少会开execution.checkpointing.interval: 60s execution.checkpointing.mode: EXACTLY_ONCE state.backend: rocksdb state.checkpoints.dir: hdfs:///flink/checkpoints间隔不建议太短比如10秒一次在高吞吐场景下会导致频繁的barrier对齐反而降低吞吐。60秒是我常用的区间。如果任务状态很大把RocksDB启成增量模式能显著减少checkpoint时长state.backend.rocksdb.incremental: true7.2 资源并行度估算公式我给团队定的一个粗估公式Source并行度约等于Kafka分区数聚合节点并行度等于QPS * 窗口内key的复杂度 / 单核心吞吐但更简单的做法是先按Kafka分区数的1.5倍起再观察背压Sink并行度则依据下游写入阈值来定。例如我们的订单同步任务Kafka分区是24MySQL CDC任务并行度24够用。ClickHouse Sink如果并行度也是24每秒能处理超过1万条。如果ClickHouse出现写入瓶颈不是一味加并行度而是先调大批量刷新参数。7.3 MiniBatch和LocalAggregate是聚合任务的救星默认Flink SQL聚合是一条一条计算每次来一条数据都要访问状态。数据量大时状态访问会成为瓶颈。开启MiniBatch后Flink会在聚合算子内先攒一批数据再统一访问状态减少重复读写。table.exec.mini-batch.enabled: true table.exec.mini-batch.size: 20000 table.exec.mini-batch.latency: 2s同时打开本地聚合table.optimizer.agg-phase-strategy: TWO_PHASE这个组合能极大缓解高基数聚合的性能问题。我们在计算用户行为指标时开启前单作业CPU达到70%开启后降到40%延迟也稳定多了。7.4 数据倾斜的通用解法Flink SQL聚合最怕某一个key的数据量特别大。比如实时大屏按省份统计点击量北京一个省的数据可能占到40%。这时候不管并行度开到多少所有数据都压到同一个key的聚合算子其他key的空闲。一种常用的SQL层解法是“两阶段无Key聚合”先给数据带一个随机后缀partition到一个临时的局部聚合再把后缀去掉做全局聚合。Flink SQL里可以这样写INSERT INTO result SELECT province, SUM(cnt) AS total FROM ( SELECT province, CONCAT(province, _, MOD(CAST(ROWNUM AS INT), 10)) AS random_key, COUNT(*) AS cnt FROM source_table GROUP BY province, CONCAT(province, _, MOD(CAST(ROWNUM AS INT), 10)) ) GROUP BY province;这个做法牺牲了一点精确度但对大部分指标误差很小。如果必须完全准确就得在DataStream层自己写keyby的加盐逻辑。7.5 监控指标比日志更重要上线不是结束而是另一种开始。Flink SQL作业我一般盯四个指标currentInputWatermark判断数据源是否断流、numRecordsInPerSecond判断吞吐是否达标、numBufferedRecords判断Sink是否反压、lastCheckpointSize判断状态增长是否异常。这四个里任何一个出现明显异常我都会第一时间看SQL的执行计划而不是先去翻日志。举个例子有一次lastCheckpointSize突然涨了3倍日志里一切正常。后来检查发现是源表某个上游业务字段脏数据太多导致GROUP BY的key基数飙升。要不是提前盯住了checkpoint大小等到状态过大再恢复作业时整条链路可能已经停了一个小时。最后再分享一个小技巧这套Flink SQL链路在我们线上稳定跑了近半年最大并行度到48每天处理亿级事件。踩了这么多坑以后我最大的体会是Flink SQL不是“低配版流处理”而是一种边界明显的工程选择。把合适任务交给SQL把复杂状态剥离给自定义算子你会省下大量时间。最后补一个比JDBC连接器更隐蔽的坑如果线上任务用了GROUP BY做大窗口聚合并且状态一直增长别只想着加资源先去查状态里的key是不是因为某个字段值本身就一直在变。一个字段写错状态增长的速度会远远超出你的预期。
返回列表