ARTICLE DETAIL

资讯详情

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

深入解析 Flink Job 优化技巧:让大数据处理更高效Flink Job 优化全攻略

深入解析 Flink Job 优化技巧:让大数据处理更高效Flink Job 优化全攻略 目录1、 使用 DataGen 造数据1.1 DataStream 的 DataGenerator1.2 SQL 的 DataGenerator2、 算子指定 UUID1提交案例未指定 uid2提交案例指定 uid3、 链路延迟测量4、 开启对象重用5、 细粒度滑动窗口优化1细粒度滑动的影响2解决思路3细粒度的滑动窗口案例4时间分片案例本博客总结为B站尚硅谷大数据Flink2.0调优Flink性能优化视频的笔记总结。尚硅谷https://so.csdn.net/so/search?q%E5%B0%9A%E7%A1%85%E8%B0%B7spm1001.2101.3001.7020大数据Flink2.0调优Flink性能优化https://www.bilibili.com/video/BV1Q5411f76Phttps://www.bilibili.com/video/BV1Q5411f76P1、 使用 DataGen 造数据开发完 Flink 作业压测的方式很简单先在 kafka 中积压数据之后开启 Flink 任务出现反压就是处理瓶颈。相当于水库先积水一下子泄洪。数据可以是自己造的模拟数据也可以是生产中的部分数据。造测试数据的工具DataFactory、datafaker 、DBMonster、Data-Processer 、Nexmark、Jmeter 等。Flink 从 1.11 开始提供了一个内置的 DataGen 连接器主要是用于生成一些随机数用于在没有数据源的时候进行流任务的测试以及性能测试等。1.1 DataStream 的 DataGeneratorimport com.atguigu.flink.tuning.bean.OrderInfo; ​ import com.atguigu.flink.tuning.bean.UserInfo; ​ import org.apache.commons.math3.random.RandomDataGenerator; ​ import org.apache.flink.api.common.typeinfo.Types; ​ import org.apache.flink.configuration.Configuration; ​ import org.apache.flink.configuration.RestOptions; ​ import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator; ​ import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; ​ import org.apache.flink.streaming.api.functions.source.datagen.DataGeneratorSource; ​ import org.apache.flink.streaming.api.functions.source.datagen.RandomGenerator; ​ import org.apache.flink.streaming.api.functions.source.datagen.SequenceGenerator; ​ public class DataStreamDataGenDemo { ​ public static void main(String[] args) throws Exception { ​ Configuration conf new Configuration(); ​ conf.set(RestOptions.ENABLE_FLAMEGRAPH, true); ​ StreamExecutionEnvironment env ​ StreamExecutionEnvironment.createLocalEnvironmentWithWebUI(conf); ​ env.setParallelism(1); ​ env.disableOperatorChaining(); ​ SingleOutputStreamOperatorOrderInfo orderInfoDS env ​ .addSource(new DataGeneratorSource(new RandomGeneratorOrderInfo() ​ { ​ Override ​ public OrderInfo next() { ​ return new OrderInfo( ​ random.nextInt(1, 100000), ​ random.nextLong(1, 1000000), ​ random.nextUniform(1, 1000), ​ System.currentTimeMillis()); ​ } ​ })) ​ .returns(Types.POJO(OrderInfo.class)); ​ SingleOutputStreamOperatorUserInfo userInfoDS env ​ .addSource(new DataGeneratorSourceUserInfo( ​ new SequenceGeneratorUserInfo(1, 1000000) { ​ RandomDataGenerator random new RandomDataGenerator(); ​ Override ​ public UserInfo next() { ​ return new UserInfo( ​ valuesToEmit.peek().intValue(), ​ valuesToEmit.poll().longValue(), ​ random.nextInt(1, 100), ​ random.nextInt(0, 1)); ​ } ​ } ​ )) ​ .returns(Types.POJO(UserInfo.class)); ​ orderInfoDS.print(order); ​ userInfoDS.print(user); ​ env.execute(); ​ } ​ }1.2 SQL 的 DataGeneratorimport org.apache.flink.configuration.Configuration; ​ import org.apache.flink.configuration.RestOptions; ​ import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; ​ import org.apache.flink.table.api.bridge.java.StreamTableEnvironment; ​ public class SQLDataGenDemo { ​ public static void main(String[] args) throws Exception { ​ Configuration conf new Configuration(); ​ conf.set(RestOptions.ENABLE_FLAMEGRAPH, true); ​ StreamExecutionEnvironment env ​ StreamExecutionEnvironment.createLocalEnvironmentWithWebUI(conf); ​ env.setParallelism(1); ​ env.disableOperatorChaining(); ​ ​ StreamTableEnvironment tableEnv StreamTableEnvironment.create(env); ​ String orderSqlCREATE TABLE order_info (\n ​ id INT,\n ​ user_id BIGINT,\n ​ total_amount DOUBLE,\n ​ create_time AS localtimestamp,\n ​ WATERMARK FOR create_time AS create_time\n ​ ) WITH (\n ​ connector datagen,\n ​ rows-per-second20000,\n ​ fields.id.kindsequence,\n ​ fields.id.start1,\n ​ fields.id.end100000000,\n ​ fields.user_id.kindrandom,\n ​ fields.user_id.min1,\n ​ fields.user_id.max1000000,\n ​ fields.total_amount.kindrandom,\n ​ fields.total_amount.min1,\n ​ fields.total_amount.max1000\n ​ ); ​ String userSqlCREATE TABLE user_info (\n ​ id INT,\n ​ user_id BIGINT,\n ​ age INT,\n ​ sex INT\n ​ ) WITH (\n ​ connector datagen,\n ​ rows-per-second20000,\n ​ fields.id.kindsequence,\n ​ fields.id.start1,\n ​ fields.id.end100000000,\n ​ fields.user_id.kindsequence,\n ​ fields.user_id.start1,\n ​ fields.user_id.end1000000,\n ​ fields.age.kindrandom,\n ​ fields.age.min1,\n ​ fields.age.max100,\n ​ fields.sex.kindrandom,\n ​ fields.sex.min0,\n ​ fields.sex.max1\n ​ ); ​ tableEnv.executeSql(orderSql); ​ tableEnv.executeSql(userSql); ​ tableEnv.executeSql(select * from order_info).print(); ​ // tableEnv.executeSql(select * from user_info).print(); ​ } ​ }2、 算子指定 UUID对于有状态的 Flink 应用推荐给每个算子都指定唯一用户 IDUUID。 严格地说仅需要给有状态的算子设置就足够了。但是因为 Flink 的某些内置算子如 window是有状态的而有些是无状态的可能用户不是很清楚哪些内置算子是有状态的哪些不是。所以从实践经验上来说我们建议每个算子都指定上 UUID。默认情况下算子 UID 是根据 JobGraph 自动生成的JobGraph 的更改可能会导致UUID 改变。手动指定算子 UUID 可以让 Flink 有效地将算子的状态从 savepoint 映射到作业修改后拓扑图可能也有改变的正确的算子上。比如替换原来的 Operator 实现、增加新的Operator、删除Operator等等至少我们有可能与Savepoint中存储的Operator状态对应上。这是 savepoint 在 Flink 应用中正常工作的一个基本要素。Flink 算子的 UUID 可以通过 uid(String uid) 方法指定通常也建议指定 name。#算子.uid(指定 uid) .reduce((value1, value2) - Tuple3.of(uv, value2.f1, value1.f2 value2.f2)) .uid(uv-reduce).name(uv-reduce)1提交案例未指定 uidbin/flink run \ ​ -t yarn-per-job \ ​ -d \ ​ -p 5 \ ​ -Drest.flamegraph.enabledtrue \ ​ -Dyarn.application.queuetest \ ​ -Djobmanager.memory.process.size1024mb \ ​ -Dtaskmanager.memory.process.size2048mb \ ​ -Dtaskmanager.numberOfTaskSlots2 \ ​ -c com.atguigu.flink.tuning.UvDemo \ ​ /opt/module/flink-1.13.1/myjar/flink-tuning-1.0-SNAPSHOT.jar触发保存点//直接触发 flink savepoint jobId [targetDirectory] [-yid yarnAppId] #on yarn 模式需要指定-yid 参数 //cancel 触发 flink cancel -s [targetDirectory] jobId [-yid yarnAppId] #on yarn 模式需要指定-yid 参数 bin/flink cancel -s hdfs://hadoop1:8020/flink-tuning/sp 98acff568e8f0827a67ff37648a29d7f -yid application_1640503677810_0017修改代码从 savepoint 恢复bin/flink run \ -t yarn-per-job \ -s hdfs://hadoop1:8020/flink-tuning/sp/savepoint-066c90-6edf948686f6 \ -d \ -p 5 \ -Drest.flamegraph.enabledtrue \ -Dyarn.application.queuetest \ -Djobmanager.memory.process.size1024mb \ -Dtaskmanager.memory.process.size2048mb \ -Dtaskmanager.numberOfTaskSlots2 \ -c com.atguigu.flink.tuning.UvDemo \ /opt/module/flink-1.13.1/myjar/flink-tuning-1.0-SNAPSHOT.jar报错如下Caused by: java.lang.IllegalStateException: Failed to rollback to checkpoint/savepoint hdfs://hadoop1:8020/flink-tuning/sp/savepoint-066c90-6edf948686f6. Cannot map checkpoint/savepoint state for operatorddb598ad156ed281023ba4eebbe487e3 to the new program, because the operator is not available in the new program. If you want to allow to skip this, you can set the --allowNonRestoredState option on the CLI.临时处理在提交命令中添加--allowNonRestoredState short: -n跳过无法恢复的算子。2提交案例指定 uidbin/flink run \ -t yarn-per-job \ -d \ -p 5 \ -Drest.flamegraph.enabledtrue \ -Dyarn.application.queuetest \ -Djobmanager.memory.process.size1024mb \ -Dtaskmanager.memory.process.size2048mb \ -Dtaskmanager.numberOfTaskSlots2 \ -c com.atguigu.flink.tuning.UidDemo \ /opt/module/flink-1.13.1/myjar/flink-tuning-1.0-SNAPSHOT.jar触发保存点//cancel 触发 savepoint bin/flink cancel -s hdfs://hadoop1:8020/flink-tuning/sp 272e5d3321c5c1481cc327f6abe8cf9c-yid application_1640268344567_0033修改代码从保存点恢复bin/flink run \ -t yarn-per-job \ -s hdfs://hadoop1:8020/flink-tuning/sp/savepoint-272e5d-d0c1097d23e0 \ -d \ -p 5 \ -Drest.flamegraph.enabledtrue \ -Dyarn.application.queuetest \ -Djobmanager.memory.process.size1024mb \ -Dtaskmanager.memory.process.size2048mb \ -Dtaskmanager.numberOfTaskSlots2 \ -c com.atguigu.flink.tuning.UidDemo \ /opt/module/flink-1.13.1/myjar/flink-tuning-1.0-SNAPSHOT.jar3、 链路延迟测量对于实时的流式处理系统来说我们需要关注数据输入、计算和输出的及时性所以处理延迟是一个比较重要的监控指标特别是在数据量大或者软硬件条件不佳的环境下。Flink提供了开箱即用的 LatencyMarker 机制来测量链路延迟。开启如下参数metrics.latency.interval: 30000 #默认 0表示禁用单位毫秒监控的粒度分为以下 3 档➢ single每个算子单独统计延迟➢ operator默认值每个下游算子都统计自己与 Source 算子之间的延迟➢ subtask每个下游算子的 sub-task 都统计自己与 Source 算子的 sub-task 之间的延迟。metrics.latency.granularity: operator #默认 operator一般情况下采用默认的 operator 粒度即可这样在Sink 端观察到的 latency metric就是我们最想要的全链路端到端延迟。subtask 粒度太细会增大所有并行度的负担不建议使用。LatencyMarker 不会参与到数据流的用户逻辑中的而是直接被各算子转发并统计。为了让它尽量精确有两点特别需要注意➢ 保证 Flink 集群内所有节点的时区、时间是同步的ProcessingTimeService 产生时间戳最终是靠 System.currentTimeMillis()方法可以用 ntp 等工具来配置。➢ metrics.latency.interval 的时间间隔宜大不宜小一般配置成 3000030 秒左右。一是因为延迟监控的频率可以不用太频繁二是因为 LatencyMarker 的处理也要消耗一定性能。提交案例bin/flink run \ -t yarn-per-job \ -d \ -p 5 \ -Drest.flamegraph.enabledtrue \ -Dyarn.application.queuetest \ -Djobmanager.memory.process.size1024mb \ -Dtaskmanager.memory.process.size2048mb \ -Dtaskmanager.numberOfTaskSlots2 \ -Dmetrics.latency.interval30000 \ -c com.atguigu.flink.tuning.UidDemo \ /opt/module/flink-1.13.1/myjar/flink-tuning-1.0-SNAPSHOT.jar可以通过下面的 metric 查看结果flink_taskmanager_job_latency_source_id_operator_id_operator_subtask_index_latency端到端延迟的 tag 只有 murmur hash 过的算子 ID用 uid()方法设定的并没有算子名称[FLINK-8592] LatencyMetric scope should include operator names - ASF JIRA并且官方暂时不打算解决这个问题所以我们要么用最大值来表示要么将作业中 Sink 算子的 ID 统一化。比如使用了 Prometheus 和 Grafana 来监控效果如下4、 开启对象重用当调用了 enableObjectReuse 方法后Flink 会把中间深拷贝的步骤都省略掉SourceFunction 产生的数据直接作为 MapFunction 的输入可以减少 gc 压力。但需要特别注意的是这个方法不能随便调用必须要确保下游 Function 只有一种或者下游的Function 均不会改变对象内部的值。否则可能会有线程安全的问题。bin/flink run \ -t yarn-per-job \ -d \ -p 5 \ -Drest.flamegraph.enabledtrue \ -Dyarn.application.queuetest \ -Djobmanager.memory.process.size1024mb \ -Dtaskmanager.memory.process.size2048mb \ -Dtaskmanager.numberOfTaskSlots2 \ -Dpipeline.object-reusetrue \ -Dmetrics.latency.interval30000 \ -c com.atguigu.flink.tuning.UidDemo \ /opt/module/flink-1.13.1/myjar/flink-tuning-1.0-SNAPSHOT.jar5、 细粒度滑动窗口优化1细粒度滑动的影响当使用细粒度的滑动窗口窗口长度远远大于滑动步长时重叠的窗口过多一个数据会属于多个窗口性能会急剧下降。我们经常会碰到这种需求以 3 分钟的频率实时计算 App 内各个子模块近 24 小时的PV 和 UV。我们需要用粒度为 1440 / 3 480 的滑动窗口来实现它但是细粒度的滑动窗口会带来性能问题有两点➢ 状态对于一个元素会将其写入对应的(key, window)二元组所圈定的 windowState 状态中。如果粒度为 480那么每个元素到来更新 windowState 时都要遍历 480 个窗口并写入开销是非常大的。在采用 RocksDB 作为状态后端时checkpoint 的瓶颈也尤其明显。➢ 定时器每一个(key, window)二元组都需要注册两个定时器一是触发器注册的定时器用于决定窗口数据何时输出二是 registerCleanupTimer()方法注册的清理定时器用于在窗口彻底过期如 allowedLateness 过期之后及时清理掉窗口的内部状态。细粒度滑动窗口会造成维护的定时器增多内存负担加重。2解决思路DataStreamAPI中自己解决[FLINK-7001] Improve performance of Sliding Time Window with pane optimization - ASF JIRA。我们一般使用滚动窗口在线存储读时聚合的思路作为解决方案1从业务的视角来看往往窗口的长度是可以被步长所整除的可以找到窗口长度和窗口步长的最小公约数作为时间分片一个滚动窗口的长度2每个滚动窗口将其周期内的数据做聚合存到下游状态或打入外部在线存储内存数据库如 RedisLSM-based NoSQL 存储如 HBase3扫描在线存储中对应时间区间可以灵活指定的所有行并将计算结果返回给前端展示。3细粒度的滑动窗口案例提交案例统计最近 1 小时的 uv1 秒更新一次滑动窗口bin/flink run \ -t yarn-per-job \ -d \ -p 5 \ -Drest.flamegraph.enabledtrue \ -Dyarn.application.queuetest \ -Djobmanager.memory.process.size1024mb \ -Dtaskmanager.memory.process.size2048mb \ -Dtaskmanager.numberOfTaskSlots2 \ -c com.atguigu.flink.tuning.SlideWindowDemo \ /opt/module/flink-1.13.1/myjar/flink-tuning-1.0-SNAPSHOT.jar \ --sliding-split false4时间分片案例提交案例统计最近 1 小时的 uv1 秒更新一次滚动窗口状态存储bin/flink run \ -t yarn-per-job \ -d \ -p 5 \ -Drest.flamegraph.enabledtrue \ -Dyarn.application.queuetest \ -Djobmanager.memory.process.size1024mb \ -Dtaskmanager.memory.process.size2048mb \ -Dtaskmanager.numberOfTaskSlots2 \ -c com.atguigu.flink.tuning.SlideWindowDemo \ /opt/module/flink-1.13.1/myjar/flink-tuning-1.0-SNAPSHOT.jar \ --sliding-split trueFlink 1.13 对 SQL 模块的 Window TVF 进行了一系列的性能优化可以自动对滑动窗口进行切片解决细粒度滑动问题。Windowing TVF | Apache Flinkhttps://nightlies.apache.org/flink/flink-docs-release-1.13/docs/dev/table/sql/queries/window-tvf/
返回列表