dl-on-flink
将深度学习训练任务嵌入 Apache Flink 算子,实现数据处理与分布式训练的统一资源调度与故障
加载项目详情…
本应用为开源项目,仅供学习研究,请遵守其开源协议。
将深度学习训练任务嵌入 Apache Flink 算子,实现数据处理与分布式训练的统一资源调度与故障
加载项目详情…
本应用为开源项目,仅供学习研究,请遵守其开源协议。
假设你在一家推荐系统公司工作,线上每天处理数十亿条用户行为数据。你们用 Apache Flink 做实时特征工程——滑动窗口、Session 切分、实时统计特征,数据管道已经跑得非常成熟。
某天,产品同学提了个需求:能不能用深度学习模型来提升推荐效果? 你的第一反应是"当然可以",但真正动手时,问题来了:
Deep Learning on Flink(DL on Flink) 正是为了解决这个"数据处理和深度学习两层皮"的问题而诞生的。它由阿里巴巴 Flink 团队主导开源,核心思路是:让深度学习任务直接跑在 Flink 算子里,由 Flink 来统一管理分布式环境、资源调度和故障恢复。
DL on Flink 的灵感来自一个朴素的工程哲学:如果 Flink 已经能优雅地管理分布式批流处理任务,那为什么不能让它也管一管分布式机器学习训练?
传统的"DataProc"模式(数据处理用 Flink,训练用 TensorFlow/PyTorch 独立集群)存在三个核心痛点:资源碎片化(两套集群需要独立采购和维护)、数据搬运成本(Shuffle 网络 IO 成为瓶颈)、故障处理割裂(训练失败需要人工介入)。DL on Flink 通过将 DL 任务嵌入 Flink JobGraph,让 Flink 的调度器和状态管理能力直接服务于深度学习,实现了真正意义上的"统一资源层"。
项目目前支持 TensorFlow 1.15.x、TensorFlow 2.4.x 和 PyTorch 1.11.x,底层依赖 Flink 1.14.x,Apache License 2.0 开源。
DL on Flink 的架构设计围绕两个核心角色展开:Application Master(AM) 和 Node。
AM 角色运行在 Flink 的 JobManager 侧,负责整个分布式机器学习集群的生命周期管理。它的核心是一个可定制状态机(State Machine):根据不同 ML 框架的特性,状态机逻辑各不相同——TensorFlow 和 PyTorch 的训练流程阶段不同,状态机的状态转移路径也不同。这种设计让框架无关的通用调度逻辑与框架特定的运行逻辑解耦,新增框架只需实现新的状态机即可。
Node 角色运行在 Flink 的 TaskManager 侧,负责启动 Python 进程并准备算法运行环境。每个 Node 持有一个 Runner——同样是抽象接口,TensorFlow Runner 和 PyTorch Runner 的运行逻辑完全不同。这种双层抽象(状态机 + Runner)让整个框架具有极好的可扩展性。
项目采用 Maven 多模块结构,共包含六个核心子项目:
每个子项目均包含 src/(Java 源码)和 python/(Python SDK)两部分,体现了该项目 Java + Python 混合开发的本质。
传统 DataProc 模式中,Flink 处理完特征后需要将数据序列化后网络传输给训练集群。DL on Flink 做了关键优化:数据直接在 Flink 内部流转,Flink Source 读取 HDFS/Kafka 数据,经过特征工程后,直接通过 Flink 内部网络 Shuffle 给 ML 训练 Worker,中间没有网络序列化和反序列化开销。
典型架构中,左侧是传统的 DataProc 模式(DataProc = Flink 处理 + 外部数据传输 + 独立 ML 集群),右侧是 DL on Flink 模式——所有数据流都在 Flink 集群内部完成。
构建分 Java 和 Python 两条线:Java 侧执行 mvn -DskipTests clean install 生成 JAR 包;Python 侧需要先构建 Java 依赖,再执行 pip install 安装各模块的 Python SDK。项目还提供了 tools/build_wheel.sh 脚本用于打包 Python wheels。
Docker 部署方面:项目提供 Dockerfile 基于 flink:1.6-hadoop27 镜像构建,内置 Hadoop 2.8.0 环境。docker/build_cluster/ 目录提供了完整的集群启动脚本(start_hdfs.sh、start_zookeeper.sh、start_flink.sh、start_cluster.sh),但这些脚本需要手动逐个执行,尚未封装为 docker-compose 或 Helm Chart。
整体部署链路较长,需要准备好 Java/Maven/Python/cmake 构建环境,从源码编译 Java 和 Python 两套依赖,启动 HDFS + Zookeeper + Flink 集群,最后提交 Flink ML 作业。整个过程对运维经验要求较高,没有一键部署脚本或 Helm Chart 来简化这个过程。
DL on Flink 的上手门槛较高,主要体现在:需要同时具备 Java(Flink 侧)和 Python(ML 侧)的开发能力;需要理解 Flink 的分布式架构和 JobGraph 机制;需要维护一个完整的 Flink 集群;依赖版本较多(TF/PyTorch/Flink/Java 的版本组合兼容性需要仔细处理)。
但一旦成功部署,它在以下场景下价值显著:大规模特征工程 + 深度学习联合优化(消除数据 Shuffle 瓶颈)、需要强一致性故障恢复的生产级 ML 训练管道、Flink 生态内已有成熟数据处理管道的团队、以及需要统一运维界面的 Flink 用户。
该项目的版本维护已出现明显滞后:Flink 版本停在 1.14.x(当前 Flink 已到 1.20+),TensorFlow 支持最高到 2.4.x(当前 TF 已到 2.18+),PyTorch 停在 1.11.x。这些版本对于追求新特性的团队来说可能不够用。
此外,项目缺乏活跃的社区维护(Star 693,Fork 197),Issues 响应不及时,是一个典型的"内部开源"项目。如果用于生产,需要有自行维护和升级版本的能力。
Deep Learning on Flink 代表了一种将分布式数据处理与分布式深度学习统一编排的有益尝试。它的 AM/Node 双角色架构和状态机抽象在设计上具有可扩展性,Maven 多模块结构也体现了工程的规范性。对于已有成熟 Flink 数据处理管道的团队,它是将 DL 训练纳入统一平台的可行选择。但部署链路较长、版本维护滞后、缺乏一键部署工具等现实问题,意味着它更适合有 Flink 和分布式系统经验的团队,而非初学者或追求"开箱即用"的场景。