Hi Zhanghao,

Thanks for driving this FLIP. This is an interesting topic.

Ray has become an important execution layer in the Python data and AI
ecosystem. Besides Ray-native libraries, there are quite a few engines
such as Daft and Modin that use Ray as an execution backend, and Dask
workloads can also run through Dask-on-Ray. Supporting Flink on Ray
would allow PyFlink users to reuse an existing Python-native cluster
and share Ray’s CPU, GPU, and resource scheduling capabilities.

The overall architecture looks reasonable: use Ray actors as
deployment units that supervise Flink JVM processes, while relying on
Ray for resource allocation and keeping scheduling, checkpointing, and
task recovery in the existing Flink runtime.

I have several questions about the design details.

### TaskManager actor creation and teardown

The FLIP could further clarify how TaskManager actor creation and
ownership work across a JobManager failover, for example:

- Whether a JobManager failover causes existing TaskManager actors to
restart. My expectation is that healthy TaskManager actors and their
JVM processes should remain running.
- How a recovered JobManager obtains handles to existing TaskManager actors.
- How duplicate TaskManager actor creation is prevented during failover.
- Whether TaskManager actors are named (if so, how the name will be
derived) and detached from their creator.
- How stale worker-management requests from the previous JobManager
are fenced after a failover?
- How detached TaskManager actors are cleaned up after the application
reaches a terminal state.

The design should also distinguish a JobManager restart from an
application shutdown. Restarting or replacing the JobManager should
preserve healthy TaskManager actors, while application shutdown should
explicitly and reliably terminate all detached TaskManager actors.

### High availability model

Whether the proposed deployment reuses an existing Flink HA
implementation or will introduce a new `RayHighAvailabilityServices`?

Besides, whether JobManager HA supports the following cases:
- one restartable JobManager actor; or
- multiple active/standby JobManager actors.

### GPU resource management

The FLIP could clarify the GPU allocation and isolation granularity.

When a TaskManager actor requests GPU resources, Ray sets
`CUDA_VISIBLE_DEVICES` for the actor. The TaskManager JVM and the
Python UDF workers launched by it may therefore see the same set of
GPUs. If multiple GPU-consuming Python workers run concurrently in the
same TaskManager, they may compete for GPU memory and compute
resources.

Configuring one slot per GPU-enabled TaskManager can reduce the
likelihood of such contention, but it does not by itself guarantee
isolation. A Flink slot may contain multiple operators or operator
chains, and an operator chain may contain multiple GPU-consuming UDFs.

Supporting multiple independently isolated GPU consumers in one
TaskManager would require an additional per-slot or per-Python-worker
device allocation mechanism, including assigning a distinct
`CUDA_VISIBLE_DEVICES` value when each Python worker is launched. The
behavior should be clearly defined. For the initial version, the FLIP
could also document that GPU allocation and visibility are defined at
the TaskManager actor level and that no additional slot- or UDF-level
isolation is provided.

### API and naming

The FLIP currently uses both `init_flink()` and `init_flink_on_ray()`.
I guess it could be simplified as following:

    import pyflink.ray as flink_on_ray
    flink_on_ray.init(...)

The FLIP could also clarify whether `init()` invokes `ray.init()`
automatically, how it behaves when a Ray context already exists, and
whether the corresponding shutdown operation terminates only the Flink
application or the entire Ray cluster.

The actor names could also be made consistent. What about renaming
`RayFlinkJobManagerActor` to `RayJobManagerActor` so that it is
symmetric with `RayTaskManagerActor`.

Thanks,
Dian

On Fri, Sep 18, 2026 at 12:50 PM Zhanghao Chen
<[email protected]> wrote:
>
> Hi all,
> I would like to start a discussion about FLIP-608: Flink on Ray — Ray 
> Resource Backend For Flink [1].
> This FLIP proposes a new Ray resource backend to allow running Flink jobs 
> natively on Ray. The hybrid approach takes the best of both worlds 
> —battle-tested streaming and batch unific computing plus heterogeneous 
> scheduling and the Ray AI eco-system (Data/Serve/Train). One can therefore 
> setup the entire real-time AI workflow in a single cluster with higher 
> resource efficiency and lower maintenance cost.
> The proposal focuses on supporting running Flink jobs on Ray in 
> Application/Session mode by translating Flink resource requirements to Ray. 
> The Flink core mechanisms like RPC, shuffle, checkpointing, failover and slot 
> management are preserved.
> A follow-up FLIP will focus on the interoperability between PyFlink and Ray 
> by supporting conversion between Flink Dataframe/Datastream/Table and Ray 
> Dataset.
> Looking forward to your feedback!
> [1] 
> https://cwiki.apache.org/confluence/spaces/FLINK/pages/449286348/FLIP-608+Flink+on+Ray+%E2%80%94+Ray+Resource+Backend+For+Flink
>
> Best,
> Zhanghao Chen

Reply via email to