kestra-io / kestra-io/plugin-oci

[Plugin] OCI — Data Flow (Spark-as-a-Service)

Open
#8 0 comments 0 reactions 0 assignees View on GitHub
area/plugin
Dominant language
Java
Stars
0
Forks
0
Avg merge
18h 3m
Merged PRs (30d)
1

Description

## Summary

The OCI Data Flow sub-plugin for `plugin-oci` enables Kestra flows to submit, monitor, and manage Spark jobs on OCI Data Flow — Oracle's managed Spark-as-a-service. Data engineering teams can embed large-scale distributed processing steps directly in their Kestra pipelines, with full visibility into job state and the ability to react to completion or failure without polling infrastructure.

## Motivation

Teams running PySpark or Scala Spark workloads on OCI Data Flow currently submit jobs via OCI CLI scripts or the console and poll status manually. Kestra can't natively wait for a Data Flow run to complete and branch on success or failure. A native `SubmitRun` task with built-in wait-for-completion logic closes that gap and lets data pipelines treat Spark jobs as first-class tasks with typed outputs (exit code, logs URI, outputs location).

## Context

Part of the OCI Plugin Suite EPIC: https://github.com/kestra-io/plugin-oci/issues/2
Reference: `plugin-databricks` job run submission pattern.

## API Reference

- **Official docs**: https://docs.oracle.com/en-us/iaas/api/#/en/dataflow/latest/
- **Authentication**: Config-file (`~/.oci/config`), instance principal, or `SimpleAuthenticationDetailsProvider`
- **Base URL pattern**: `https://dataflow.{region}.oci.oraclecloud.com/20200129/`
- **SDK**: OCI Java SDK v3.87.0 via BOM

## Gradle Dependencies

Add to `build.gradle`:

```groovy
// OCI Java SDK BOM
implementation platform("com.oracle.oci.sdk:oci-java-sdk-bom:3.87.0")
// Data Flow
implementation "com.oracle.oci.sdk:oci-java-sdk-dataflow"
implementation "com.oracle.oci.sdk:oci-java-sdk-common"
```

## Plugin Structure

- **Repository**: `plugin-oci`
- **Namespace**: `io.kestra.plugin.oci.dataflow`
- **Sub-plugins**: `dataflow` (applications, runs)

## Suggested Tasks

1. `CreateRun` — submit a Data Flow application run with parameters, wait for completion
2. `GetRun` — fetch run details and emit lifecycle state, log URI as outputs
3. `ListRuns` — list runs for an application with optional state filter
4. `DeleteRun` — cancel and delete a run
5. `CreateApplication` — register a new Data Flow application (JAR/Python + config)
6. `ListApplications` — list available applications in a compartment
7. Add `RunCompletedTrigger` — poll until a run reaches SUCCEEDED or FAILED state
8. Write unit + integration tests

## YAML Examples

### Example 1 — Submit a Spark job and wait for completion

```yaml
id: run_spark_job
namespace: company.data

inputs:
- id: input_path
type: STRING
- id: output_path
type: STRING

tasks:
- id: submit_run
type: io.kestra.plugin.oci.dataflow.CreateRun
region: eu-frankfurt-1
tenancyOcid: "{{ secret('OCI_TENANCY_OCID') }}"
userId: "{{ secret('OCI_USER_OCID') }}"
fingerprint: "{{ secret('OCI_FINGERPRINT') }}"
privateKey: "{{ secret('OCI_PRIVATE_KEY') }}"
applicationId: "{{ secret('OCI_DATAFLOW_APP_OCID') }}"
compartmentId: "{{ secret('OCI_COMPARTMENT_OCID') }}"
displayName: kestra-etl-{{ execution.id }}
arguments:
- "--input={{ inputs.input_path }}"
- "--output={{ inputs.output_path }}"
waitForCompletion: true

- id: log_result
type: io.kestra.plugin.core.log.Log
message: "Run {{ outputs.submit_run.runId }} finished with state {{ outputs.submit_run.lifecycleState }}"
```

### Example 2 — List all runs for an application

```yaml
id: list_dataflow_runs
namespace: company.data

tasks:
- id: list_runs
type: io.kestra.plugin.oci.dataflow.ListRuns
region: eu-frankfurt-1
tenancyOcid: "{{ secret('OCI_TENANCY_OCID') }}"
userId: "{{ secret('OCI_USER_OCID') }}"
fingerprint: "{{ secret('OCI_FINGERPRINT') }}"
privateKey: "{{ secret('OCI_PRIVATE_KEY') }}"
applicationId: "{{ secret('OCI_DATAFLOW_APP_OCID') }}"
lifecycleState: SUCCEEDED

- id: log_count
type: io.kestra.plugin.core.log.Log
message: "{{ outputs.list_runs.count }} successful runs found"
```

### Example 3 — React when a Data Flow run completes

```yaml
id: on_spark_run_complete
namespace: company.data

triggers:
- id: run_completed
type: io.kestra.plugin.oci.dataflow.RunCompletedTrigger
region: eu-frankfurt-1
tenancyOcid: "{{ secret('OCI_TENANCY_OCID') }}"
userId: "{{ secret('OCI_USER_OCID') }}"
fingerprint: "{{ secret('OCI_FINGERPRINT') }}"
privateKey: "{{ secret('OCI_PRIVATE_KEY') }}"
runId: "{{ secret('OCI_DATAFLOW_RUN_OCID') }}"
interval: PT1M

tasks:
- id: notify
type: io.kestra.plugin.core.log.Log
message: "Spark run {{ trigger.runId }} completed — state: {{ trigger.lifecycleState }}, logs: {{ trigger.logsBucketUri }}"
```

## Acceptance Criteria

- [ ] `CreateRun`, `GetRun`, `ListRuns`, `DeleteRun`, `CreateApplication`, `ListApplications` tasks implemented
- [ ] `CreateRun` supports `waitForCompletion: true` with configurable poll interval
- [ ] `RunCompletedTrigger` polling trigger implemented
- [ ] Run outputs include `lifecycleState`, `logsBucketUri`, `outputsBucketUri`
- [ ] All `Property` fields support Kestra expression language
- [ ] Unit + integration tests pass (`./gradlew test`)
- [ ] `package-info.java` with `@PluginSubGroup(category = PluginSubGroup.PluginCategory.CLOUD)`
- [ ] Build passes with `./gradlew build`

Contributor guide

No contributing guide indexed for this repository

Research direction

Start by reading the plugin-databricks job run submission pattern and the existing build.gradle in plugin-oci. The issue lists the CreateRun, GetRun, ListRuns, DeleteRun, CreateApplication, ListApplications, and RunCompletedTrigger entry points, plus ./gradlew test and ./gradlew build. Done means all listed tasks, polling behavior, outputs, expression support, package-info.java, and tests satisfy the acceptance criteria.

Written by the indexing model from the issue text.

Assessment

Tech stack
java, spark
Domain
cloud, data-engineering, distributed-systems
Issue type
Feature
Difficulty
5/5
Estimated time
Over a week
Activity status
Active
Clarity
Mostly clear
Newbie friendliness
30/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.