ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

Flink部署模式深度解析:Session、Per-Job、Application选型与实战

Flink部署模式深度解析:Session、Per-Job、Application选型与实战 Flink部署模式这个话题说实话我盯了很久了。刚接手实时数仓那阵子我先后在Standalone集群、YARN Session、K8s Application模式里都踩过一遍坑才算是把“部署模式”这四个字背后的资源分配逻辑、提交方式差异、以及运维侧的实际影响给理顺了。网上讲Flink部署的文章不少但多数要么停在概念层要么只讲一种模式真正能把Session、Per-Job、Application三种主流模式放到实际场景里对比、放到具体案例如MySQL同步ClickHouse里去说的并不多。这篇博文我会把三种部署模式的原理、适用场景、资源平台选型以及真实生产中的参数配置、常见异常排查串起来讲适合刚入门想弄清部署概念的也适合已经在生产环境里跑Flink、想优化现有部署方案的读者。1. 三种部署模式先弄清它们到底解决什么问题1.1 为什么会有三种模式一个Flink作业从你写好代码到真正在集群跑起来中间要回答三个问题资源从哪来、作业怎么提交、生命周期谁来管。三种部署模式本质上就是这三个问题的三种答案组合。Session模式是老大哥最早期的Flink集群基本都是这个形态。它的核心思路是“先开机再干活”预先启动一个常驻的Flink集群包含一个JobManager和若干个TaskManager然后多个作业共用这一套资源。Per-Job模式则走了另一个极端每个作业提交时单独拉起一个完整的集群作业跑完整个集群随之释放。Application模式出现得最晚它是为了解决Per-Job模式里一个尴尬的问题——当使用flink run提交作业时作业的main()方法其实是在客户端本机执行的这意味着客户端必须保持在线还得拥有一份可用的Flink环境这在大规模集群或离线提交场景里非常头疼。Application模式把main()方法挪到了集群内部执行资源、生命周期、提交入口全部统一管理。这三种模式不是谁替代谁的关系而是分别对应了不同的使用场景。我打个比方Session模式像公司里的一间共享会议室大家轮流用成本低但谁都能进可能出现互相干扰Per-Job模式像每个项目租一整层写字楼隔离最好但成本极高Application模式则像云厂商的弹性容器每个应用启动一个独立的环境用完即释放最符合生产环境对隔离和效率的双重需求。1.2 Session模式共享资源的轻量级选择Session模式最典型的用法是flink run提交到已经存在的session集群。我最早学习Flink时就是这么干的启动一个standalone session然后不断提交各种demo作业改代码、重启、再提交整个过程非常流畅因为作业启动时不需要等待集群拉起TaskManager资源已经在那了。这个模式的优点显而易见启动快、资源复用率高、运维成本低。特别是当一个团队有大量短时间运行的探索性作业时共享集群能让整体资源利用率保持在不错的水准。但它有两个让生产环境很不舒服的问题。第一多作业资源隔离差。假设一个低优先级的作业把Slot全部占满另一个核心作业提交时就无Slot可用甚至在内存紧张时JobManager会反复尝试分配资源导致作业长时间卡在RESTARTING状态。第二依赖包冲突风险高。所有作业共用同一个Classloader体系两个作业各自带的不同版本的连接器jar包在加载时往往出现NoSuchMethodError或ClassNotFoundException这问题特别隐蔽排查起来极其痛苦。所以我的建议是Session模式适合开发测试、适合作业量少且资源需求平稳的小团队但如果你的作业是7x24小时跑在线的或者作业之间计算差异很大请尽早绕开它。1.3 Per-Job模式隔离优先的经典方案Per-Job模式在我的理解里更像是一种“作业级别隔离”的思路而不是一个简单的参数开关。在YARN环境下Per-Job模式提交时Flink会为这个作业单独申请一个YARN ApplicationJobManager和TaskManager都在这个Application内部运行作业结束整个Application结束资源立刻归还给YARN。这种隔离带来的好处是实实在在的一个作业的异常退出不会影响其他作业内存溢出、Full GC、连接器崩溃都被隔离在单独的集群里。而且每个作业可以按需申请资源不需要在一个共享集群里跟别人“抢饭吃”。不过Per-Job模式的代价是启动时间变长。因为每次提交都要经历“申请资源 - 启动JobManager - 注册TaskManager - 分发任务”这个过程对于需要频繁启停的作业类型来说这个开销不可忽视。我在一次实时指标作业的开发周期里明显感受到每天改代码重新提交光是等集群拉起和状态恢复就多花了两三分钟。这种体验促使我开始认真考虑Application模式。1.4 Application模式生产环境的主流答案Application模式是我目前在大多数生产项目里的首选也是这篇博文后续两个实际案例的基础。它的核心区别是main()方法在JobManager进程中执行而不是在本地客户端执行。这带来三个关键优势。第一客户端不再持有集群任何状态你提交完命令后就可以关掉笔记本去喝咖啡提交的作业照常运行第二多个作业之间的Classloader天然隔离每个Application都有自己的类加载环境不用再担心jar包冲突第三资源生命周期和作业生命周期完全绑定作业启动时申请资源作业结束后全部释放。用一句话总结三种模式的核心差异Session模式下资源跟着集群走Per-Job模式下资源跟着作业走但业务逻辑还在客户端Application模式下资源、业务逻辑、作业生命周期三者彻底统一。三种模式的选择往往不是单纯比哪个“好”而是看你的运行环境。举个例子如果你就是一个人写几个测试作业Standalone Session够了如果你在一个有完善YARN或K8s平台的大团队那Application模式几乎是把“生产可靠”写在脸上的方案。比较维度Session模式Per-Job模式Application模式资源生命周期集群常驻作业间共享每作业申请独立资源作业结束释放每作业独立资源作业结束释放启动速度快无需等资源慢需拉起整个集群中需拉起JobManager隔离性差作业互相影响好好客户端依赖需要持有一份可用Flink环境需要本地执行main()客户端只负责提交不执行main()适用场景开发测试、探索性分析作业数量少、隔离要求高生产环境长期运行的核心作业Classloader共享容易冲突独立独立1.5 模式选择对日常运维的影响部署模式不只是一个启动命令的区别它会延伸到很多日常运维细节。比如日志排查Session模式下所有作业的日志都混在JobManager日志里你要定位一个作业的异常常常得在几十兆的日志里翻半天Application模式每个作业一个JobManager日志天然按作业隔离排查起来省力得多。再比如资源估算。Session模式下你很难说清楚某个作业到底占用了多少CPU和内存因为TaskManager是被多个作业共享的而Application模式下每个作业的资源申请量是明确写在提交参数里的容量规划、预算核算、甚至故障恢复的资源预留都清晰很多。还有版本升级。Session模式下升级Flink版本你得先停掉所有作业升级完再一个个重启停机时间很难避免Application模式则可以逐个作业迁移老版本作业继续跑新作业用新版本提交实现平滑升级。这些体验都是我实际运维时体会特别深的。2. 部署在哪儿Standalone、YARN 还是 Kubernetes2.1 Standalone适合学习不适合生产Standalone部署是Flink自带的最朴素的模式不依赖任何外部资源管理系统JobManager和TaskManager都是独立的Java进程。这种模式最大的优点是简单解压安装包改几行配置启动脚本集群就起来了。我最初学习Flink时就是在自己电脑上搭了个Standalone集群跑通了第一个WordCount那种“诶能跑起来”的成就感至今记得。但Standalone致命的短板在于没有资源管理能力。CPU、内存、磁盘这些资源全靠手动分配TaskManager死了不会自动重启JobManager挂了没人接管扩容缩容完全靠人肉。生产环境偶尔会见到用Standalone模式跑内部工具的团队但我个人强烈不建议在任何承载核心链路的业务上用它除非你做好了随时修复集群的心理准备。2.2 YARN大数据生态里的老牌主力YARN部署模式是我在实际项目中用得最久的一种方式因为很多公司的技术栈本来就是围绕Hadoop生态构建的数据仓库里的HDFS、Hive、Spark都在YARN上跑Flink接入YARN顺理成章。Flink on YARN可以支持三种部署模式Session和Per-Job的应用场景前面已经说过了。在YARN上跑Flink一个最让我放心的地方是资源管理质量高。YARN经过十余年打磨内存分配、容器调度、失败重试、队列隔离这些机制都比较成熟不需要额外关心基础设施层面的稳定性。而且YARN会把Flink作业的日志统一收集到HDFS上历史作业日志排查非常方便这点在生产排障时价值极大。缺点是如果公司没有完整的大数据平台单独为了Flink去搭一套Hadoop生态性价比很低。YARN本身也是个“大家伙”部署和维护KMS、HDFS、NodeManager并不是一件轻松的事。我后来在容器化改造时发现K8s在不少场景里已经可以替代YARN的职责了。2.3 Kubernetes云原生时代的必然选择Kubernetes部署是这几年的主流趋势Flink社区对它的支持也是越来越完善。目前主流玩法有两种一种是Flink原生Kubernetes模式通过flink run -t kubernetes-application提交Flink会自己调用K8s API创建JobManager和TaskManager的Pod另一种是Flink Kubernetes Operator模式通过声明式YAML描述作业由Operator负责作业的创建、更新和自动恢复。我对K8s模式最大的感受是故障恢复能力强。TaskManager Pod被意外杀掉K8s会自动拉起一个新的Pod作业在Flink checkpoint机制配合下可以实现比较平滑的恢复。配合Horizontal Pod Autoscaler还能做一定程度的弹性伸缩这在夜间低价时段能省不少成本。当然K8s模式也对使用者的容器化知识提出了更高的要求。你要理解Pod、Service、ConfigMap、PVC这些基本概念还要会构建合适的Flink镜像。我见过不少团队在K8s上跑Flink结果因为镜像内jar包版本不对、资源配额没设置好、或者JobManager没有配置持久化存储导致作业一重启就“失忆”。这些坑我会在后面的实操部分详细展开。2.4 资源平台选型速查平台学习成本运维成本故障恢复弹性伸缩适用场景Standalone低中差无本地学习、简单Demo、测试环境YARN中中良一般已建设Hadoop栈的传统大数据团队Kubernetes较高中高优支持配合HPA云原生团队、容器化改造后的业务选型的核心逻辑其实很简单资源平台要和公司的技术栈高度匹配不要为了“新潮”特意上一个K8s也不要为了“稳妥”固守Standalone。我之前带过一个项目团队一直用Standalone跑作业每次集群扩容都要折腾几个小时后来迁移到K8s后扩容几乎变成了改一条配置。但也见到过习惯了YARN的团队硬上K8s之后天天和网络插件的坑搏斗。3. 典型案例用Application模式实现MySQL同步ClickHouse3.1 场景拆解为什么需要实时同步很多公司都会遇到这样一个需求业务数据存在MySQL里但分析报表、实时看板、大屏展示需要更强悍的列式查询能力于是数据仓库的在线分析层选用ClickHouse。传统做法是用定时任务或canal脚本把数据按小时或天同步过去但这样数据的时效性就差很多。业务方跟我说“想看今天的实时销售额”你得给到分钟级甚至秒级的数据。Flink在这个场景里天然有优势它能把MySQL的binlog实时解析成变更流经过简单的清洗转换后通过JDBC连接器写入ClickHouse。整个链路没有中间件部署层面用Application模式资源隔离好、生命周期清晰、恢复也方便。3.2 整体方案设计我的方案长这样Flink应用作为YARN/K8s上的一个独立Application运行数据源是MySQL的orders订单表数据目标是ClickHouse的实时订单表。这里我选用了Flink CDC连接器读取MySQL的binlog配合一个自研的JDBC Sink连接器写入ClickHouse。为什么不用Flink官方自带的flink-connector-jdbc直接写因为ClickHouse官方虽然兼容MySQL协议但JDBC驱动的写入模式偏“一条条执行”批处理性能上不去。我改造后的Sink会在内存中批量攒数据达到批次大小或时间间隔后就执行一次批量插入配合ClickHouse的BatchInsert特性吞吐量能提升好几倍。3.3 核心代码实现先看Maven依赖。这里要注意CDC连接器的版本要和Flink主版本严格对应否则大概率遇到一堆奇怪的序列化异常。properties flink.version1.17.2/flink.version /properties dependencies !-- Flink核心依赖 -- dependency groupIdorg.apache.flink/groupId artifactIdflink-streaming-java/artifactId version${flink.version}/version scopeprovided/scope /dependency !-- Flink CDC连接器读取MySQL binlog -- dependency groupIdcom.ververica/groupId artifactIdflink-connector-mysql-cdc/artifactId version2.4.0/version /dependency !-- ClickHouse JDBC驱动 -- dependency groupIdcom.clickhouse/groupId artifactIdclickhouse-jdbc/artifactId version0.3.2-patch8/version /dependency !-- Flink JDBC连接器自己封装Sink时用 -- dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-jdbc/artifactId version${flink.version}/version /dependency /dependencies然后看数据源部分。用Flink CDC构建MySqlSource这里有一个很重要的参数是debezium.snapshot.mode我推荐设置为initial这样Flink启动时会先全量读取一次当前表数据然后自动切换到增量binlog无缝衔接。MySqlSourceString mySqlSource MySqlSource.Stringbuilder() .hostname(rm-xxx.mysql.rds.aliyuncs.com) .port(3306) .databaseList(analytics) // 要同步的数据库 .tableList(analytics.orders) // 要同步的数据表 .username(flink_user) .password(your_password) .deserializer(new JsonDebeziumDeserializationSchema()) // binlog数据以JSON格式输出 .startupOptions(StartupOptions.initial()) // 先从全量开始然后增量 .build();Sink部分我推荐自己封装一个ClickHouse写入器核心逻辑是维护一个批量缓冲队列配合SinkFunction的invoke方法批量flush。这里我直接放一个精简版思路public class ClickHouseBatchSink extends RichSinkFunctionString { private static final int MAX_BATCH_SIZE 1000; private static final long MAX_BUFFER_TIME_MS 5000; private transient Connection connection; private transient PreparedStatement statement; private transient ListString buffer; private transient long lastFlushTime; Override public void open(Configuration parameters) throws Exception { connection DriverManager.getConnection( jdbc:clickhouse://clickhouse-host:8123/analytics, default, ); statement connection.prepareStatement( INSERT INTO analytics.orders VALUES (?, ?, ?, ?)); buffer new ArrayList(); lastFlushTime System.currentTimeMillis(); } Override public void invoke(String value, Context context) throws Exception { buffer.add(value); if (buffer.size() MAX_BATCH_SIZE || (System.currentTimeMillis() - lastFlushTime) MAX_BUFFER_TIME_MS) { flush(); } } private void flush() throws SQLException { for (String record : buffer) { // 解析JSON并设置statement参数 // statement.addBatch(); } statement.executeBatch(); buffer.clear(); lastFlushTime System.currentTimeMillis(); } Override public void close() throws Exception { if (buffer.size() 0) { flush(); } if (statement ! null) statement.close(); if (connection ! null) connection.close(); } }批量flush参数要重点说明批次大小不要一味贪大。我实测过ClickHouse单次批量写入的吞吐和批次大小是曲线关系约1000条到5000条时性价比最高超过1万条后单条延迟反而上升而且一旦失败重试的代价非常大整个批次的几百KB数据都得重来。所以生产环境我一般设成1000~2000条同时配合5秒的时间窗口能比较好地兼顾吞吐和延迟。主程序里开启Checkpoint并设置成EXACTLY_ONCE因为ClickHouse sink如果不上等幂或事务机制重复写入会导致数据翻倍。这里我把Sink设置成幂等模式靠业务主键去重。StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.enableCheckpointing(60000, CheckpointingMode.EXACTLY_ONCE); env.getCheckpointConfig().setMinPauseBetweenCheckpoints(30000); env.getCheckpointConfig().setCheckpointStorage(hdfs:///flink/checkpoints); DataStreamSourceString stream env.addSource(mySqlSource); stream.addSink(new ClickHouseBatchSink()); env.execute(MySQL-Sync-To-ClickHouse);3.4 部署到K8sApplication模式的完整提交过程这里我用Flink Kubernetes Operator部署这是我认为目前生产体验最好的方式。首先写一个Application的YAML文件核心参数如下apiVersion: flink.apache.org/v1beta1 kind: FlinkApplication metadata: name: mysql-clickhouse-sync namespace: flink-jobs spec: image: registry.example.com/flink-job/mysql-clickhouse-sync:1.0.0 imagePullPolicy: Always flinkVersion: v1_17 serviceAccountName: flink jobManager: replicas: 1 resource: cpu: 1.0 memory: 2g taskManager: replicas: 2 resource: cpu: 1.0 memory: 4g job: jarURI: local:///opt/flink/usrlib/mysql-clickhouse-sync.jar parallelism: 4 entryClass: com.example.flink.MySQLClickHouseJob args: - --checkpoint.interval60000 upgradeMode: stateless restartPolicy: type: OnFailure maxAttempts: 3有几个点我必须提醒你。第一jarURI这里的路径是镜像内的路径不是本地路径所以你在构建镜像时一定要把fatjar拷贝到镜像的/opt/flink/usrlib/目录下。第二jobManager的memory不要给太小JobManager不仅要跑main()方法Application模式还要管理checkpoint元数据和任务调度默认给2g是比较稳的起步值。第三upgradeMode如果作业有状态需求最好设置成last-state并配置persistentVolume来保存checkpoint这样升级Flink作业时能自动恢复状态。命令提交非常简单一条kubectl apply搞定kubectl apply -f mysql-clickhouse-sync.yaml提交后Operator会自动创建Pod并启动JobManager、TaskManager。这时可以用下面的命令确认作业状态kubectl get flinkapplication -n flink-jobs kubectl logs -f deployment/mysql-clickhouse-sync-jobmanager -n flink-jobs我实际部署后观察到的现象是第一次全量同步20万行订单数据用了不到30秒之后增量同步的延迟基本维持在秒级整个链路非常平稳。3.5 同步过程中的坑与排查这个案例上线后我接连遇到三个坑都值得单独说说。第一个坑是时区问题导致时间字段错位。MySQL里order_time存的是CST时区的时间ClickHouse从JDBC驱动读出来后默认按服务器时区解析结果写入ClickHouse后时间比源表早了8个小时。解决办法是给JDBC连接URL追加useTimezonetrueserverTimezoneAsia/Shanghai并在Flink侧设置env.setStreamTimeCharacteristic对应的时区参数。实测加上后时间完全对上了。第二个坑是DDL变更后同步直接崩掉。业务方某天在MySQL订单表加了一列没有通知我们Flink CDC发现元数据不一致直接抛DebeziumException导致作业挂掉。这个问题的解法是在CDC连接器后面加上一次SchemaChangeEvent的处理遇到DDL变更时把消息转发到一个报警通道并暂停写入待人工确认之后再恢复。代码层面可以监听schemaChangeEvents输出到钉钉或企业微信Webhook。第三个坑是网络抖动导致JDBC连接池被耗尽。刚开始Sink用的是每次开启一个新Connection结果ClickHouse偶发慢查询导致连接积累最终报Too many connections。我彻底重构了Sink的连接管理使用Hikari连接池并设置maximumPoolSize10配合空闲回收时间这个问题就消失了。4. 生产环境的坑JDBC连接器异常与Spring Boot整合Flink4.1 JDBC连接器异常的排查思路Flink的JDBC连接器是很多实时任务的“最后一公里”但它也是异常高发区。我在多个项目里遇到过几类典型异常这里以排障实录的形式整理出来。异常一java.sql.SQLTransientConnectionException: HikariPool-1 - Connection is not available, request timed out这类异常最常见出现在ClickHouse、MySQL等外部存储短暂不可用或慢查询时。排查步骤是先看ClickHouse当前活跃连接数SHOW PROCESSLIST确认是否有大量慢查询堆积再检查Hikari连接池的maximumPoolSize是否过小导致并发请求全部等待超时最后检查Flink侧是否有背压导致Sink写入速度跟不上连带着把连接池挤爆。我之前一个项目就是下游ClickHouse导入线程数不够高峰期几十个查询挤在一个连接上几分钟就把连接池拖垮了。异常二Caused by: java.lang.NoClassDefFoundError: com/mysql/cj/jdbc/Driver这种异常本质上是运行时缺类多半是打包没打进去或者依赖冲突。我建议排查时先确认MySQL驱动和Flink SQL客户端是否已包含在镜像的/opt/flink/lib目录中再确认fatjar里的META-INF/services文件是否因为shade操作被合并覆盖了。有一个很实用的经验是不要把所有jar包都打成uber-jar而是把连接器相关jar放到镜像的lib目录业务逻辑打成瘦jar这样类加载会清晰很多。异常三JDBC update failed with SQLException: Data truncation: Data too long for column这类字段长度不匹配的问题其实很常见。源MySQL字段是varchar(255)目标ClickHouse列定义成了String看起来没问题但MySQL写入时可能包含超长文本ClickHouse的严格模式直接拒绝写入。解决办法是明确字段映射并在Flink侧做长度截断或类型转换。异常四Connection reset by peer/SocketTimeoutException这类网络层异常多数是防火墙或负载均衡空闲超时导致的。K8s环境下Service默认的idle timeout往往很短长连接长时间无数据后会被中间设备切断。解法是给JDBC URL配置connectTimeout和socketTimeout同时让连接池定期执行SELECT 1进行健康检查。4.2 Spring Boot整合Flink什么该整合什么不该整合“Spring Boot整合Flink”这个话题我近期被问得特别多很多人一上来就想把Flink作业直接塞进Spring Boot应用里比如在SpringBootApplication的main方法里env.execute()。这样做开发调试很爽但一上生产就会发现问题Spring Boot进程一旦启动就常驻TaskManager无法跟作业生命周期联动资源释放不彻底而且Spring Boot里跑一个Flink作业本质上还是本地模式local execution并没有真正提交到集群。那Spring Boot能做什么我认为最合理的结合方式是把Spring Boot当作Flink作业的配置中心、任务编排中心和运维控制台而不是运行环境。举个例子你可以在Spring Boot里写一个REST接口接收SQL或作业参数把配置写入数据库或消息队列然后由Flink作业启动时从同一个配置源读取。这样好处是权责清晰Spring Boot负责“管作业”Flink集群负责“跑作业”。我在一个内部数据平台项目里就是这么做的前端页面配置同步任务后端Spring Boot把任务配置存到MySQLFlink作业启动时根据任务名引用配置真正实现“配置与执行分离”。开发新任务时只需要在页面上填表名和同步策略完全不用改Flink代码。4.3 两个值得记住的生产教训第一个教训是关于日志与异常处理。Flink Application模式下stdout和stderr日志默认都会输出到JobManager容器的控制台如果你没有配置日志持久化到HDFS或云存储那么作业挂掉后你会失去所有日志。我建议在所有Flink作业中加入日志输出配置例如使用log4j.properties把org.apache.flink.streaming和com.yourjob包下的日志输出到独立文件并挂载到PVC。这个配置虽然小但在事故复盘时价值巨大。第二个教训是不要贪图大并行度盲目开Slot。K8s Application模式下TaskManager的并行度和CPU配额是两回事。我见过一个团队把TaskManager的taskmanager.numberOfTaskSlots设为16而Pod CPU只申请了1核结果高并发场景下GC严重、吞吐骤降。合理分配是每个Slot对应0.5~1核CPU内存根据实际堆需求调整。一个经典的起步配置是TaskManager 1个、4核CPU、8GB内存、4个Slot每个Slot大约1核2GB。还有个细节K8s上TaskManager的CPU配额我一般会设成cpu: 2.0而不是留给调度器自动分配。因为Flink的计算比较依赖CPU配额不足会导致频繁线程阻塞和检查点超时。5. 我自己的实操习惯和测试建议讲完模式和案例我想把自己最近用下来比较顺手的几个习惯分享出来。我现在的新项目默认都走Application模式资源平台优先K8s只有在碰到公司已有成熟YARN平台、且Flink作业量不大时才考虑退回YARN。对我个人而言Application模式的Classloader隔离和资源生命周期管理带来的收益已经远远超过了我学习K8s的那点成本。开发和测试阶段我建议你本地先跑Session模式因为调试快、改代码反复提交很高效。一旦代码稳定立刻切到Application模式生产提交。为了不让两套模式切换时出现配置混乱我在每个作业里都会强制从环境变量或配置文件读取运行模式确保本地和提交后的行为一致。最后再提一个容易被忽视的小经验提交Application模式作业时尽量把-Dexecution.checkpointing.interval这类参数放到作业代码里而不是只放在命令行。因为命令行参数很容易被运维脚本覆盖或遗忘而代码里的默认值至少有兜底。我的习惯是在StreamExecutionEnvironment里显式设置checkpoint间隔、模式、超时和最小间隔命令行的参数仅作为临时覆盖手段。关于这个项目我个人在实际部署和运维中的最大体会是部署模式本身不是技术难点难的是你愿不愿意花时间去理解每个模式背后的资源模型并且在项目初期就把部署方案当成一个正式设计项来做。很多Flink作业线上出问题根因根本不是代码逻辑而是部署时资源预估错误或模式选错。这篇文章把我踩过的坑和验证过的配置都写了希望能让你少走几个月的弯路。
返回列表