Skip to content

feat(plugins): add flyteplugins-nextflow to run Nextflow pipelines on Flyte - #1622

Draft
k1sauce wants to merge 6 commits into
flyteorg:mainfrom
k1sauce:kyle/nextflow-launcher
Draft

k1sauce wants to merge 6 commits into
flyteorg:mainfrom
k1sauce:kyle/nextflow-launcher

Conversation

@k1sauce

@k1sauce k1sauce commented Sep 25, 2026

Copy link
Copy Markdown
Contributor

Summary

Adds flyteplugins-nextflow, which runs a Nextflow pipeline inside a Flyte task. Each Nextflow task runs as a child action of that task, so the whole pipeline is one Flyte run with one action per Nextflow task.

The executor side lives in a separate Nextflow plugin, nf-flyte. That plugin submits each task as a plain container action through the Actions/Run services, and uses the credentials Flyte injects into the task pod. This package is the Flyte side. It builds the image the Nextflow head runs in, and launches Nextflow with everything nf-flyte needs.

import flyte
from flyte.io import Dir
from flyteplugins.nextflow import nextflow_image, run_nextflow

env = flyte.TaskEnvironment(name="nextflow", image=nextflow_image())

@env.task
async def demo() -> Dir:
    return await run_nextflow("nf-core/demo", revision="1.2.0", profile="test", outdir="results")

What's included

  • nextflow_image(): an image containing a JRE, a pinned Nextflow version (default 25.10.6) and the nf-flyte plugin. The plugin comes from the Nextflow registry, or from a local zip for development. The Nextflow runtime and nf-amazon are downloaded when the image is built, not on every run.
  • run_nextflow(...), which is called from any task:
    • writes the Nextflow config that sets process.executor = 'flyte'
    • passes the current run name and action to nf-flyte through environment variables, so task actions are created under the calling action
    • defaults the S3 work directory to a prefix in the run's own storage, so no bucket setup is needed; a fixed work_dir plus resume=True resumes across runs
    • supports revision, profile, params (passed as a params file), config (a file path or config text) and extra_args
    • streams Nextflow's output into the task logs
    • returns outdir as a Dir
    • raises NextflowError with the end of the Nextflow output and of .nextflow.log on failure
    • on task cancellation, sends SIGTERM to Nextflow so it aborts its running task actions
  • examples/nf_core_demo.py and a README

Design notes

  • Nextflow keeps control of retries (errorStrategy), caching (-resume) and file staging. nf-flyte disables Flyte retries and caching on task actions.
  • Process containers don't need Flyte or the AWS CLI. nf-flyte adds an init container that copies a static s5cmd into each task pod for S3 staging. The only image requirement is bash.
  • nextflow_image() doesn't check that a local plugin zip exists. The task module is imported again inside the pod, where the zip isn't present; a missing zip fails the image build instead.

Testing

  • Unit tests (plugins/nextflow/tests) run run_nextflow against a fake nextflow executable. They check the command line, the generated config, params, env vars, the default outdir/work_dir, and error reporting.
  • End to end on a hosted Flyte v2 cluster (AWS/S3): nf-core/demo 1.2.0 with -profile test succeeded. It ran 8 task actions (COWPY, FASTQC ×3, SEQTK_TRIM ×3, MULTIQC) on the pipeline's unmodified biocontainers images, and the results came back as a Dir. A separate test pipeline confirmed S3 staging in and out, and that Nextflow's errorStrategy retries work.

Before merging

  • nf-flyte isn't published to the Nextflow plugin registry yet, and github.com/unionai/nf-flyte isn't public yet. Until then, the default nextflow_image() fails to build. Pass nf_flyte=Path(".../nf-flyte-<version>.zip") instead.
  • Only S3-backed deployments are supported for now: the work dir must be on S3.

Signed-off-by: Kyle Hazen <kyle@union.ai>
@k1sauce
k1sauce force-pushed the kyle/nextflow-launcher branch from 43571ab to c1a7257 Compare September 25, 2026 01:31

@samhita-alla samhita-alla left a comment •

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nice work getting nf-core/demo running end to end. A few notes inline, plus two things here.

Can we hold the merge until nf-flyte is published? CI discovers everything under plugins/ and publishes it, so merging now would release flyteplugins-nextflow with a default image that doesn't build.

The bigger question for me is what Flyte adds while retries, caching and staging all live in Nextflow. Right now it's mostly a pod launcher. I opened unionai/nf-flyte#1 about letting Flyte stage task data and cache on Nextflow's task hash, which would also drop the S3 requirement.

Smaller follow-ups that could fit here: render -with-report/-with-timeline/-with-dag into the task's Flyte report, and in local mode run Nextflow with its own local or docker executor instead of raising.

Comment thread plugins/nextflow/src/flyteplugins/nextflow/_run.py
Comment thread plugins/nextflow/src/flyteplugins/nextflow/_run.py Outdated
Comment thread plugins/nextflow/src/flyteplugins/nextflow/_run.py Outdated
Comment thread plugins/nextflow/src/flyteplugins/nextflow/_run.py Outdated
Comment thread plugins/nextflow/src/flyteplugins/nextflow/_run.py Outdated
Comment thread plugins/nextflow/src/flyteplugins/nextflow/_run.py Outdated
Comment thread plugins/nextflow/src/flyteplugins/nextflow/_run.py Outdated
Comment thread plugins/nextflow/src/flyteplugins/nextflow/_image.py Outdated
…w fixes

- Keep resume state in the work dir: cloud cache (or NXF_CACHE_DIR for local
  work dirs), a session id derived from the work dir, history lookup off and
  an explicit -name. A bare -resume was ignored, since .nextflow/ lived in the
  pod's temp launch dir.
- Default work dir and relative outdir to the task's raw data prefix, minus
  the per-attempt segment, and pass -resume on retries so a retried head
  picks up where the failed attempt stopped.
- Parent task actions under tctx.task_action.
- Forward SIGTERM to Nextflow while it runs (the runtime installs no handler,
  so the CancelledError path never ran on abort); keep the shutdown wait
  inside the pod grace period and keep relaying output while it stops.
- Read Nextflow output in chunks, so lines over 64 KiB don't fail the task.
- report=True renders Nextflow's report, timeline and DAG into the task's
  Flyte report.
- Local runs use Nextflow's own executors instead of raising.
- Pre-install nf-cloudcache in nextflow_image().

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Signed-off-by: Kyle Hazen <kyle@union.ai>
@k1sauce

k1sauce commented Sep 25, 2026

Copy link
Copy Markdown
Contributor Author

Thanks for the careful review. Most of the inline points are fixed in the latest push, and the resume changes are verified on a cluster. Details are in the threads. On the broader points:

  • Holding the merge: agreed. I'll keep this as a draft until nf-flyte is on the registry.
  • What Flyte adds: fair question, and I agree with the direction in unionai/nf-flyte#1. That's the bigger change, so I'd keep it there rather than grow this PR. This PR stays a thin launcher, and it shouldn't need to change much when staging and caching move into Flyte.
  • Reports: run_nextflow(report=True) now adds -with-report, -with-timeline and -with-dag, and renders each into its own tab of the task's Flyte report (the task needs report=True). Verified on union-internal: the report has Nextflow report, Timeline and DAG tabs, each with Nextflow's full document.
  • Local mode: in a local run, run_nextflow no longer raises. It skips nf-flyte and uses Nextflow's own executors (local, or docker through profile/config).

k1sauce and others added 4 commits September 25, 2026 13:02
nf-flyte now stages task data through Flyte's copilot (unionai/nf-flyte#3), so
task pods no longer access the work dir and it no longer has to be on S3. On a
cluster the work dir must be an object store location Flyte can read; the
default, under the task's raw data prefix, always is.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Signed-off-by: Kyle Hazen <kyle@union.ai>
The repo's docstring check rejects Sphinx field lists (:param:, :return:);
convert nextflow_image() and run_nextflow() to Args: / Returns:.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Signed-off-by: Kyle Hazen <kyle@union.ai>
@samhita-alla

Copy link
Copy Markdown
Contributor

@k1sauce one more case: the head pod dying while tasks are running. Your retry test failed the head after all 8 tasks had finished, so nothing was in flight.

The 409 attach in FlyteClient.enqueue is meant for this, but I don't think it fires. On resume, Nextflow gives an unfinished task dir a new hash (it bumps tries in checkCachedOrLaunchTask), and the action name comes from task.hash. So the in-flight task gets a new action name and runs a second time next to the old one. Could the name come from cacheKey() plus failCount instead? cacheKey() doesn't change across sessions, so a retried head would hit the 409 and attach.

Can we add the pod-kill case as an example: start nf-core/demo, delete the head pod during FASTQC, and check that attempt 2 finishes without resubmitting the in-flight tasks? With that in, I think the PR is done.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants