Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
110 changes: 110 additions & 0 deletions docs/en/training_guides/assets/pipelines/verify-kfp-features.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,110 @@
"""Verify KFP caching, ParallelFor, and persisted artifact passing."""

import os
import time

from kfp import Client, compiler, dsl


@dsl.component(base_image="docker-mirrors.alauda.cn/library/python:3.12-slim")
def produce_dataset(seed: str, dataset: dsl.Output[dsl.Dataset]) -> None:
"""Write a small dataset to the KFP artifact path."""
from pathlib import Path

output = Path(dataset.path)
output.parent.mkdir(parents=True, exist_ok=True)
output.write_text(f"seed={seed}\n", encoding="utf-8")
dataset.metadata["seed"] = seed


@dsl.component(base_image="docker-mirrors.alauda.cn/library/python:3.12-slim")
def verify_dataset(dataset: dsl.Input[dsl.Dataset], item: int, seed: str) -> str:
"""Read the persisted artifact in one ParallelFor iteration."""
from pathlib import Path

content = Path(dataset.path).read_text(encoding="utf-8")
expected = f"seed={seed}\n"
if content != expected:
raise RuntimeError(f"artifact content {content!r} != {expected!r}")
return f"item={item}:{content.strip()}"


@dsl.pipeline(name="kfp-mechanisms-verification")
def mechanisms_pipeline(seed: str = "alauda-kfp-features") -> None:
"""Produce one artifact and consume it in three parallel tasks."""
producer = produce_dataset(seed=seed)
with dsl.ParallelFor(items=[1, 2, 3], parallelism=3) as item:
verify_dataset(dataset=producer.outputs["dataset"], item=item, seed=seed)


def task_state(task: object) -> str:
"""Return a task state as a stable uppercase string."""
value = getattr(task, "state", "")
return str(value).rsplit(".", 1)[-1].upper()


def run_once(
client: Client, package_path: str, experiment: str, suffix: str, seed: str
) -> object:
"""Submit a cached run and wait for its terminal state."""
result = client.create_run_from_pipeline_package(
pipeline_file=package_path,
arguments={"seed": seed},
run_name=f"kfp-mechanisms-{suffix}-{int(time.time())}",
experiment_name=experiment,
namespace=os.environ["KFP_NAMESPACE"],
enable_caching=True,
)
run = client.wait_for_run_completion(result.run_id, timeout=900, sleep_duration=5)
if str(run.state).upper() != "SUCCEEDED":
raise RuntimeError(f"run {result.run_id} ended in {run.state}")
print(f"run {suffix}: {result.run_id} state={run.state}", flush=True)
return client.get_run(result.run_id)


def main() -> None:
"""Compile the pipeline, run it twice, and verify the three mechanisms."""
package_path = "/tmp/kfp-mechanisms-verification.yaml"
compiler.Compiler().compile(mechanisms_pipeline, package_path=package_path)
client = Client(host=os.environ["KFP_ENDPOINT"], namespace=os.environ["KFP_NAMESPACE"])
experiment = f"kfp-mechanisms-{int(time.time())}"
seed = f"alauda-kfp-features-{time.time_ns()}"

first = run_once(client, package_path, experiment, "first", seed)
first_tasks = first.run_details.task_details or []
for task in first_tasks:
print(
f"first task: name={task.display_name} state={task_state(task)} "
f"outputs={task.outputs or {}}",
flush=True,
)

parallel_tasks = [task for task in first_tasks if task.display_name == "verify-dataset"]
if len(parallel_tasks) != 3 or any(task_state(task) != "SUCCEEDED" for task in parallel_tasks):
raise RuntimeError(
f"ParallelFor expected 3 successful verify-dataset tasks, got {len(parallel_tasks)}"
)

producers = [task for task in first_tasks if task.display_name == "produce-dataset"]
if len(producers) != 1 or task_state(producers[0]) != "SUCCEEDED":
raise RuntimeError(f"expected one produce-dataset task, got {len(producers)}")
print(
"artifact persistence: producer uploaded the Dataset and all 3 consumers read it",
flush=True,
)

second = run_once(client, package_path, experiment, "cached", seed)
second_tasks = second.run_details.task_details or []
for task in second_tasks:
print(f"cached task: name={task.display_name} state={task_state(task)}", flush=True)
skipped = [task for task in second_tasks if task_state(task) == "SKIPPED"]
if not skipped:
raise RuntimeError("second identical run did not report any cached (SKIPPED) task")

print("PASS: ParallelFor created 3 tasks", flush=True)
print("PASS: Dataset artifact was persisted and consumed", flush=True)
print(f"PASS: cache reused {len(skipped)} task execution(s)", flush=True)


if __name__ == "__main__":
main()
2 changes: 2 additions & 0 deletions docs/en/training_guides/index.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -20,4 +20,6 @@ End-to-end recipes for fine-tuning and pretraining LLMs on Alauda AI.
| Interactive exploration, custom scripts, VolcanoJob submission | Workbench Notebook | [Fine-tuning LLMs using Workbench](./fine-tuning-using-notebooks.mdx) |
| Full-parameter SFT / pretraining on Ascend NPU | Workbench `PyTorch CANN` / `MindSpore CANN` | [Fine-tune and Pretrain on Ascend NPU](./fine-tune-and-pretrain-llms-on-ascend-npu.mdx) |
| Track a KFP run's parameters and metrics in MLflow | KFP component + MLflow SDK | [Kubeflow Pipeline + MLflow Integration](../mlflow/how_to/pipelines-mlflow-integration.mdx) |
| Train tabular or time-series models with reusable AutoGluon assets | Managed KFP pipelines or composable components | [Use Reusable Kubeflow Pipeline Components](./reusable-pipeline-components.mdx) |
| Use KFP caching, parallel loops, and persistent typed artifacts | KFP 2.16.1 execution mechanisms | [Kubeflow Pipelines Execution and Storage Behavior](./kfp-execution-and-storage.mdx) |
| Daily fine-tune → evaluate → compare loop with MLflow + TrustyAI | KFP Recurring Run + MLflow Model Registry + `LMEvalJob` | [Daily Fine-Tuning Pipeline with MLflow and TrustyAI](./fine-tuning-pipeline-with-mlflow-trustyai.mdx) |
224 changes: 224 additions & 0 deletions docs/en/training_guides/kfp-execution-and-storage.mdx
Original file line number Diff line number Diff line change
@@ -0,0 +1,224 @@
---
weight: 59
---

# Kubeflow Pipelines Execution and Storage Behavior

Kubeflow Pipelines 2.16.1 installed by Alauda `kfp-operator` v26.3.1 or the
Data Science Pipelines Operator (DSPO) v2.15.2 supports task caching,
`dsl.ParallelFor`, typed KFP artifacts, and S3-compatible artifact persistence
on amd64 and arm64 clusters. These mechanisms are KFP features; the operator
configures the API server, workflow controller, cache server, and artifact
store that implement them.

## Feature summary

| Mechanism | KFP 2.16.1 | Behavior |
|---|---|---|
| Task caching | Supported | An identical task execution can be reused and is reported as `SKIPPED` in the repeated run. |
| `dsl.ParallelFor` | Supported | Each loop item becomes a separate task; `parallelism` limits concurrent iterations. |
| Typed artifacts | Supported | `dsl.Output[dsl.Dataset]` and `dsl.Input[dsl.Dataset]` create an explicit producer-to-consumer dependency. |
| Artifact persistence | Supported | The launcher uploads artifact files and metadata to the configured S3-compatible store and downloads them for consumers. |

The verification pipeline linked below was run with DSPO on an x86 cluster. It
produced one `Dataset`, read it from three `ParallelFor` tasks, and then
submitted the identical pipeline again. All three consumers succeeded, and
four executor tasks in the second run were reported as `SKIPPED` because their
cached results were reused. The dataset object remained present in S3 after
the workflow pods completed.

:::note
The managed AutoGluon pipelines intentionally disable caching for their main
tasks. Their source object key can refer to mutable data, and each run must
publish fresh progress and model artifacts. Caching remains available to your
own pipelines.
:::

## Cache a task safely

Caching is enabled for a run by default. It can also be set when submitting a
compiled package:

```python
client.create_run_from_pipeline_package(
pipeline_file="pipeline.yaml",
arguments={"dataset_version": "2026-07-23"},
enable_caching=True,
)
```

Disable caching for a task with external side effects or mutable inputs:

```python
task = publish_report(...)
task.set_caching_options(False)
```

KFP computes a cache key from the component and its declared inputs. It cannot
detect that the contents behind an unchanged S3 object key, URL, database
query, or image tag changed. Pass an immutable object version, digest, or
explicit version parameter when fresh external content must invalidate the
cache.

## Run tasks with `ParallelFor`

Use `dsl.ParallelFor` when the same component can process independent items:

```python
from kfp import dsl


@dsl.pipeline(name="parallel-evaluation")
def evaluation_pipeline():
dataset = produce_dataset()
with dsl.ParallelFor(
items=["accuracy", "precision", "recall"],
parallelism=3,
) as metric:
evaluate(dataset=dataset.outputs["dataset"], metric=metric)
```

The controller creates one `evaluate` task for each item. Set `parallelism`
below the item count when namespace quota, node capacity, or an external API
cannot sustain every iteration at once.

## Persist and pass a typed artifact

Write files to the path supplied by the output artifact. Do not invent a local
path and expect it to survive the producer pod:

```python
from kfp import dsl


@dsl.component(base_image="python:3.12-slim")
def produce_dataset(dataset: dsl.Output[dsl.Dataset]):
from pathlib import Path

output = Path(dataset.path)
output.parent.mkdir(parents=True, exist_ok=True)
output.write_text("value\n42\n", encoding="utf-8")
dataset.metadata["rows"] = 1


@dsl.component(base_image="python:3.12-slim")
def consume_dataset(dataset: dsl.Input[dsl.Dataset]):
from pathlib import Path

print(Path(dataset.path).read_text(encoding="utf-8"))
```

Passing `producer.outputs["dataset"]` to a consumer establishes the dependency.
The KFP launcher uploads the producer's file to the artifact store and
materializes it at `dataset.path` in the consumer pod.

This is different from a KFP workspace PVC:

| Storage | Best for | Lifetime |
|---|---|---|
| Typed KFP artifact | Declared component inputs and outputs, lineage, results that must remain after the run | Persisted in the configured object store according to its retention policy |
| Workspace PVC | Large intermediate files shared by tasks in one run | Created for the run; lifecycle follows KFP workspace handling |
| Container filesystem | Temporary scratch data inside one task | Lost when the pod is removed |

## Run the verification pipeline

Download
[`verify-kfp-features.py`](https://github.com/alauda/aml-docs/tree/master/docs/en/training_guides/assets/pipelines/verify-kfp-features.py)
and run it from a Workbench or another environment that can reach the KFP API:

```bash
python -m pip install 'kfp==2.16.1'

export KFP_ENDPOINT='http://<kfp-api-service>:8888'
export KFP_NAMESPACE='my-pipeline-namespace'
python verify-kfp-features.py
```

The script compiles one producer and three parallel consumers, submits it
twice with the same parameter, and fails unless all of the following are true:

- the first run succeeds;
- exactly three `verify-dataset` tasks succeed;
- every consumer reads the producer's Dataset content; and
- the second run reports cached tasks as `SKIPPED`.

The final output contains one `PASS` line for each mechanism. An administrator
can independently confirm persistence by listing the run prefix in the KFP
artifact bucket after the pods finish.

## Deploy on an arm64 cluster \{#arm64-deployment\}

The Alauda KFP 2.16.1 releases package the pipeline service and runtime images
for both `linux/amd64` and `linux/arm64`. The optional
`kfp-metadata-writer` image remains amd64-only. DSPO does not deploy this
component; Alauda `kfp-operator` must disable it on arm64.

### `kfp-operator`

With Alauda `kfp-operator` v26.3.1, install the operator normally from
OperatorHub and set the optional metadata writer to `false` in the
`KubeflowPipelines` custom resource before reconciliation:

```yaml
apiVersion: operator.alauda.io/v1
kind: KubeflowPipelines
metadata:
name: kubeflowpipelines
spec:
global:
images:
kfpMetadataWriter:
enabled: false
```

Keep the remaining `spec.global.images` values from the release sample. The
`support_arm` image field only controls packaging/image relocation; it does not
disable the Deployment. Verify that the writer is absent after reconciliation:

```bash
kubectl -n kubeflow get deployment metadata-writer
# Error from server (NotFound) is expected
```

### Data Science Pipelines Operator (DSPO)

The V2 reconciler in DSPO v2.15.2 does not deploy `kfp-metadata-writer`. No
arm64 override is required. Create the DSPA with the normal V2 settings and do
not add an MLMD writer field:

```yaml
apiVersion: datasciencepipelinesapplications.opendatahub.io/v1
kind: DataSciencePipelinesApplication
metadata:
name: sample
namespace: data-science-project
spec:
dspVersion: v2
mlmd:
deploy: true
```

Verify the namespace contains MLMD gRPC/Envoy resources but no
`metadata-writer` Deployment. DSPO and `kfp-operator` are mutually exclusive;
choose one installation model per cluster.

## Platform compatibility and database backends

The following results apply to the Alauda builds aligned with KFP 2.16.1:

| Installation | amd64 | arm64 | MySQL | PostgreSQL |
|---|---|---|---|---|
| DSPO v2.15.2 | Supported; smoke and mechanism tests pass | Supported; DSPA readiness and KFP v2 SDK smoke pass | Supported | Not supported |
| Alauda `kfp-operator` v26.3.1 | Supported | Supported when `kfpMetadataWriter.enabled` is `false`; KFP API and Argo workflow execution verified | Supported | Not supported |

All required runtime images in these releases have amd64 and arm64 manifests,
except `kfp-metadata-writer`. Disabling that component is required only for an
arm64 Alauda `kfp-operator` installation; it is not part of a DSPO deployment.

Both operators currently expose only MySQL connection settings. DSPO performs
its external database health check with the MySQL driver and renders
`DB_DRIVER_NAME=mysql`. The `kfp-operator` chart renders `dbType=mysql`,
`DBDriverName=mysql`, and `DBCONFIG_MYSQLCONFIG_*`. Pointing either operator at
a PostgreSQL endpoint does not switch its driver and does not produce a usable
KFP installation. Use MySQL until a PostgreSQL driver selection and the
upstream PostgreSQL deployment patches are wired into the operator interface.
Loading