Use this skill when a Dagster project shells out or runs containers (non-dbt): PipesSubprocessClient / open_pipes_session, dagster-shell (execute_shell_command / create_shell_command_op), SSHResource (dagster-ssh) running remote commands, k8s_job_op / PipesK8sClient, PipesECSClient, or PipesDatabricksClient. Triggers: any op running a subprocess/shell command that is not dbt, any SSHResource remote command, any Kubernetes/ECS Pipes execution. Note: dbt is handled by dbt-core-dagster-to-orchestra.
66
80%
Does it follow best practices?
Run evals on this skill
Adds up to 20 points to the overall score
View guide
Passed
No findings from the security scan
Fix and improve this skill with Tessl
tessl review fix ./skills/migrate-to-orchestra/skills/dagster-shell-ssh-to-orchestra/SKILL.mdDagster runs external and shell workloads through several mechanisms: Dagster Pipes (PipesSubprocessClient, PipesK8sClient, PipesECSClient, PipesDatabricksClient, PipesGlueClient, PipesLambdaClient), dagster-shell (execute_shell_command), and SSHResource (dagster-ssh). Orchestra has no local shell execution model; map each to a purpose-built integration:
LINUX_SSH + LINUX_SSH_EXECUTE_COMMAND or WINDOWS_SSH + WINDOWS_SSH_COMMANDPYTHON + PYTHON_EXECUTE_SCRIPT (the better pattern)AWS_ECS, GKE, AZURE_CONTAINER_APPS, or AWS_EKS/AKSPipesGlueClient -> AWS_GLUE, PipesLambdaClient -> AWS_LAMBDA, PipesDatabricksClient -> DATABRICKS)Is it dbt?
-> YES: use dbt-core-dagster-to-orchestra
-> NO: continue
Does it run a Python script via subprocess/Pipes?
-> YES: PYTHON + PYTHON_EXECUTE_SCRIPT
Does it SSH to a remote server (SSHResource)?
-> Linux: LINUX_SSH + LINUX_SSH_EXECUTE_COMMAND
-> Windows: WINDOWS_SSH + WINDOWS_SSH_COMMAND
Is it a Kubernetes job (PipesK8sClient / k8s_job_op)?
-> AWS: AWS_EKS + AWS_EKS_RUN_JOB
-> GCP: GKE + GKE_RUN_JOB
-> Azure: AZURE_KUBERNETES_SERVICE + AKS_RUN_JOB
Is it a container (PipesECSClient / Cloud Run)?
-> AWS: AWS_ECS + AWS_ECS_RUN_TASK
-> GCP: GCP_CLOUD_RUN + GCP_CLOUD_RUN_EXECUTE_JOB
Managed-compute Pipes?
-> PipesGlueClient -> AWS_GLUE_RUN_JOB
-> PipesLambdaClient -> AWS_LAMBDA_EXECUTE_ASYNC_FUNCTION
-> PipesDatabricksClient -> DATABRICKS_RUN_WORKFLOW
Pure local glue (echo, mkdir)?
-> Assess whether it can be dropped or folded into an adjacent task# Dagster
@asset
def process(context, pipes_subprocess_client: PipesSubprocessClient):
return pipes_subprocess_client.run(
command=["python", "scripts/etl.py", "--env", "prod"],
context=context.op_execution_context,
).get_results()# Orchestra
task-001:
integration: PYTHON
integration_job: PYTHON_EXECUTE_SCRIPT
name: process
connection: my_python_git_conn_12345
parameters:
command: 'python scripts/etl.py --env prod'
python_version: '3.12'
package_manager: PIP
depends_on: []
condition: null
tags: []# Dagster
@op
def run_remote(context, ssh: SSHResource):
ssh.execute_remote_command("cd /data && ./process.sh --date 2024-01-01")# Orchestra
task-001:
integration: LINUX_SSH
integration_job: LINUX_SSH_EXECUTE_COMMAND
name: run_remote
connection: my_linux_server_12345
parameters:
command: 'cd /data && ./process.sh --date 2024-01-01'
depends_on: []
condition: null
tags: []The SSH user must have read/write access to /tmp/ on the target server.
# Dagster
@asset
def k8s_job(context, pipes_k8s_client: PipesK8sClient):
return pipes_k8s_client.run(
context=context.op_execution_context,
image="my-registry/etl:latest",
namespace="data-team",
).get_results()# AWS EKS
task-001:
integration: AWS_EKS
integration_job: AWS_EKS_RUN_JOB
name: k8s_job
connection: aws_eks_prod_12345
parameters:
job_manifest: |
apiVersion: batch/v1
kind: Job
metadata:
name: my-etl-job
spec:
template:
spec:
containers:
- name: etl
image: my-registry/etl:latest
restartPolicy: Never
depends_on: []task-001:
integration: AWS_ECS
integration_job: AWS_ECS_RUN_TASK
name: run_ecs_task
connection: aws_default_12345
parameters:
cluster: my-cluster
task_definition: my-task
subnet_ids: subnet-abc123,subnet-def456
security_group_ids: sg-abc123
assign_public_ip: false
depends_on: []from dagster import op, asset, job, Definitions
from dagster_ssh import SSHResource
from dagster_pipes import PipesSubprocessClient
@op
def fetch_data(context, ssh: SSHResource):
ssh.execute_remote_command("cd /data && python fetch.py --date 2024-01-01")
@asset
def process(context, pipes_subprocess_client: PipesSubprocessClient):
return pipes_subprocess_client.run(
command=["python", "/opt/scripts/process.py"], context=context.op_execution_context,
).get_results()
@job
def data_pipeline():
fetch_data()
defs = Definitions(jobs=[data_pipeline], assets=[process],
resources={"ssh": SSHResource(remote_host="data_server"),
"pipes_subprocess_client": PipesSubprocessClient()})version: v1
name: data-pipeline
pipeline:
stage-fetch:
tasks:
fetch-data:
integration: LINUX_SSH
integration_job: LINUX_SSH_EXECUTE_COMMAND
name: fetch_data
connection: data_server_12345
parameters:
command: 'cd /data && python fetch.py --date 2024-01-01'
depends_on: []
condition: null
tags: []
depends_on: []
stage-process:
tasks:
process-data:
integration: PYTHON
integration_job: PYTHON_EXECUTE_SCRIPT
name: process
connection: my_python_git_conn_12345
parameters:
command: 'python scripts/process.py'
python_version: '3.12'
package_manager: PIP
depends_on: []
condition: null
tags: []
depends_on: [stage-fetch]dbt-core-dagster-to-orchestra; DBT_CORE_EXECUTE is dedicated.PipesSubprocessClient (Python) -> PYTHON — the command list ["python","scripts/x.py"] becomes command: 'python scripts/x.py'.SSHResource.execute_remote_command -> LINUX_SSH — parameters.command.command, not a list; join the Pipes command list.PipesK8sClient / k8s_job_op -> EKS/GKE/AKS — Orchestra uses existing jobs/manifests, not arbitrary image launches.PipesECSClient -> AWS_ECS; PipesDatabricksClient -> DATABRICKS_RUN_WORKFLOW.set_outputs./tmp/ on the target server.alerts:
- name: on-failure
statuses: [FAILED]
destinations:
- integration: SLACK
destination: '#data-alerts'3a29fe4
If you maintain this skill, you can claim it as your own. Once claimed, you can manage eval scenarios, bundle related skills, attach documentation or rules, and ensure cross-agent compatibility.