Skip to content
Back to skills

Dagster

ASecurity

Dagster is a data pipeline orchestrator built around the concept of software-defined assets. Learn to define assets, ops, jobs, schedules, sensors, and resources for building maintainable data platforms.

  • 142 stars
  • 0 votes
  • 0 copies
  • 4 views
  • Added May 27, 2026
databasespythongobashsqlreactnodegitapidatabase

Works with

  • terminal
  • cli
  • api

Security analysis

A96/100
  • mediumInstalls packages at runtime which could introduce malicious dependencies

Pro scans all 2 files and shows the line behind each finding

Scanned October 4, 2026

npx -y skills add TerminalSkills/skills --skill dagster --agent claude-code

Installs into .claude/skills of the current project.

Are you the author of Dagster?

Add the live security badge to your README. It updates with every re-scan.

Security grade badge for Dagster
[![Security: A — Skills Directory](https://www.skillsdirectory.com/api/skills/terminalskills-dagster/badge)](https://www.skillsdirectory.com/skills/terminalskills-dagster)

More formats (shields.io, HTML) on the badges page. Keep it an A: scan every change in CI with Pro.

Download with Pro
SKILL.md
---
name: dagster
description: |
  Dagster is a data pipeline orchestrator built around the concept of software-defined assets.
  Learn to define assets, ops, jobs, schedules, sensors, and resources for
  building maintainable data platforms.
license: Apache-2.0
compatibility: 'macos, linux, windows; Python 3.10+'
metadata:
  author: terminal-skills
  version: 1.1.0
  category: data-ai
  repository: https://github.com/dagster-io/dagster
  tags:
    - dagster
    - data-pipeline
    - orchestration
    - assets
    - python
---

# Dagster

## Overview

Dagster organizes data pipelines around software-defined assets — declarations of the data artifacts your pipeline produces. Assets track lineage, enable incremental computation, and integrate with the Dagster UI. Version 1.13 (current as of October 2026) introduces **Components**, a scaffolding layer on top of the same asset/job/resource APIs described below; new projects created with the `dg` CLI use that layout, but the classic `dagster` CLI and plain Python definitions in this skill still work unchanged and are the simpler starting point for a single-file pipeline.

## Instructions

### Installation

```bash
# Install Dagster and the UI webserver
pip install dagster dagster-webserver

# Scaffold a project (classic layout: one code location, pyproject.toml, tests)
dagster project scaffold --name my_pipeline
cd my_pipeline
pip install -e ".[dev]"

# Start the dev server (loads Definitions objects found in the package)
dagster dev
# UI at http://localhost:3000
```

Newer projects can instead use `uvx create-dagster@latest project my_pipeline` (or `pip install dagster-dg-cli` then `create-dagster project my_pipeline`), which scaffolds a Components-ready layout and starts the UI with `dg dev` instead of `dagster dev`. Both layouts share the same `dagster` Python package and the asset/resource/schedule APIs below.

### Software-Defined Assets

```python
# my_pipeline/assets.py: Define assets that produce data
import dagster as dg
import pandas as pd
import httpx

@dg.asset(group_name="raw")
def raw_users(context: dg.AssetExecutionContext) -> pd.DataFrame:
    """Fetch raw user data from the billing API."""
    response = httpx.get("https://billing.internal.sharpline.io/v1/users")
    df = pd.DataFrame(response.json())
    context.log.info(f"Fetched {len(df)} users")
    return df

@dg.asset(group_name="raw")
def raw_orders(context: dg.AssetExecutionContext) -> pd.DataFrame:
    """Fetch raw order data from the billing API."""
    response = httpx.get("https://billing.internal.sharpline.io/v1/orders")
    return pd.DataFrame(response.json())

@dg.asset(group_name="analytics", deps=[raw_users, raw_orders])
def revenue_by_user(raw_users: pd.DataFrame, raw_orders: pd.DataFrame) -> pd.DataFrame:
    """Calculate total revenue per user."""
    merged = raw_orders.merge(raw_users, left_on="user_id", right_on="id")
    result = (
        merged.groupby(["user_id", "name"])
        .agg(total_revenue=("amount", "sum"), order_count=("id_x", "count"))
        .reset_index()
    )
    return result
```

`dagster as dg` is the import convention used throughout current Dagster docs and examples (`dg.asset`, `dg.Definitions`, …) — it is unrelated to the separate `dg` *CLI* installed by `dagster-dg-cli`. The older `from dagster import asset, AssetExecutionContext` style still works; both resolve to the same package.

### Resources

```python
# my_pipeline/resources.py: Configurable resources for external systems
import dagster as dg
import sqlalchemy

class DatabaseResource(dg.ConfigurableResource):
    connection_string: str

    def query(self, sql: str) -> list:
        engine = sqlalchemy.create_engine(self.connection_string)
        with engine.connect() as conn:
            result = conn.execute(sqlalchemy.text(sql))
            return [dict(row._mapping) for row in result]

    def execute(self, sql: str):
        engine = sqlalchemy.create_engine(self.connection_string)
        with engine.connect() as conn:
            conn.execute(sqlalchemy.text(sql))
            conn.commit()
```

### Assets with Resources

```python
# my_pipeline/db_assets.py: Assets that use database resources
import dagster as dg
from .resources import DatabaseResource

@dg.asset(group_name="warehouse")
def dim_users(context: dg.AssetExecutionContext, database: DatabaseResource):
    """Load cleaned user dimension table into warehouse."""
    users = database.query("SELECT id, name, email, created_at FROM raw_users")
    context.log.info(f"Loaded {len(users)} users into warehouse")
    return users
```

### Definitions

```python
# my_pipeline/__init__.py: Wire everything together
import dagster as dg
from . import assets, db_assets
from .resources import DatabaseResource
import os

all_assets = dg.load_assets_from_modules([assets, db_assets])

defs = dg.Definitions(
    assets=all_assets,
    resources={
        "database": DatabaseResource(
            connection_string=os.environ["WAREHOUSE_CONNECTION_STRING"]
        ),
    },
)
```

### Schedules and Sensors

```python
# my_pipeline/schedules.py: Time-based and event-based triggers
import os
import dagster as dg

# Job that materializes specific assets
analytics_job = dg.define_asset_job(
    name="analytics_job",
    selection=dg.AssetSelection.groups("analytics"),
)

# Cron schedule
daily_analytics = dg.ScheduleDefinition(
    job=analytics_job,
    cron_schedule="0 6 * * *",  # 6 AM daily
)

# Sensor — trigger on external event
@dg.sensor(job=analytics_job, minimum_interval_seconds=60)
def new_file_sensor(context: dg.SensorEvaluationContext):
    incoming_dir = os.environ["INCOMING_DATA_DIR"]
    new_files = [f for f in os.listdir(incoming_dir) if f.endswith(".csv")]
    if new_files:
        context.log.info(f"Found {len(new_files)} new files")
        yield dg.RunRequest(run_key=new_files[0])
    else:
        yield dg.SkipReason("No new files found")
```

### Partitioned Assets

```python
# my_pipeline/partitioned.py: Time-partitioned assets for incremental processing
import dagster as dg

daily_partitions = dg.DailyPartitionsDefinition(start_date="2026-01-01")

@dg.asset(partitions_def=daily_partitions, group_name="raw")
def daily_events(context: dg.AssetExecutionContext):
    """Fetch events for a specific date partition."""
    date = context.partition_key  # e.g., "2026-09-19"
    context.log.info(f"Processing events for {date}")
    return fetch_events(date)
```

### CLI Reference

```bash
# Development server (classic layout); Components layout uses `dg dev` instead
dagster dev

# Materialize assets
dagster asset materialize --select raw_users,raw_orders

# List assets
dagster asset list

# Run a job
dagster job execute -j analytics_job

# Load and validate all Definitions without running anything
dagster definitions validate
```

## Examples

### Example 1: Build a daily revenue pipeline

Request: "Pull users and orders from our billing API every morning and compute revenue per user."

Write `raw_users`, `raw_orders` and `revenue_by_user` as shown in Software-Defined Assets above, wire them into a `Definitions` object, then schedule them:

```python
import dagster as dg
from .assets import raw_users, raw_orders, revenue_by_user

analytics_job = dg.define_asset_job(name="analytics_job", selection=dg.AssetSelection.all())
daily_schedule = dg.ScheduleDefinition(job=analytics_job, cron_schedule="0 6 * * *")

defs = dg.Definitions(
    assets=[raw_users, raw_orders, revenue_by_user],
    schedules=[daily_schedule],
)
```

Result: `dagster dev` shows a three-node asset graph in the UI; the schedule runs at 06:00 daily, and `dagster asset materialize --select raw_users,raw_orders,revenue_by_user` runs it on demand. Each run's logs and the resulting `revenue_by_user` table are visible in the Runs tab.

### Example 2: React to new files landing in a directory

Request: "A vendor drops CSV exports into a shared folder a few times a day. Kick off the import job automatically when a new file shows up."

Use the `new_file_sensor` from Schedules and Sensors above, pointed at the drop directory via `INCOMING_DATA_DIR`:

```bash
export INCOMING_DATA_DIR="/srv/dagster/incoming"
dagster dev
```

Result: the sensor polls every 60 seconds; when a `.csv` file appears, it yields a `RunRequest` keyed on the filename, so the same file never triggers two runs even if the sensor evaluates again before the run starts. When the folder is empty it yields `SkipReason("No new files found")`, which shows up in the sensor's tick history instead of a silent no-op.

## Guidelines

- `dagster dev` is for local development: it runs an in-memory, single-process webserver and scheduler. For production, deploy with Dagster+ (the managed offering) or self-host `dagster-webserver` plus a daemon process behind a real run launcher and database — do not run `dagster dev` as a production service.
- Keep API calls, file I/O and database writes inside assets or resources, not at import time; Dagster imports your modules to build the asset graph, so top-level side effects run on every CLI invocation and UI reload.
- Prefer `ConfigurableResource` over module-level clients so connection strings and credentials come from `Definitions(resources=...)` (backed by environment variables), not hardcoded values — this also lets tests swap in a fake resource.
- Partition definitions must line up with the data: a `DailyPartitionsDefinition` asset downstream of a non-partitioned asset needs an explicit partition mapping, or materialization fails.
- `dagster definitions validate` loads every asset, job, schedule and sensor without executing them — run it in CI to catch import errors and bad resource wiring before deploying.
- Sensors poll; `minimum_interval_seconds` is a floor, not a guarantee — a slow sensor function or a busy instance can delay the next tick.
- For a new project today, `create-dagster project` (the `dg` CLI) is the path the official docs lead with and scaffolds Components support; `dagster project scaffold` still produces a working classic project and needs no extra packages, which is simpler for a single pipeline like the one above.

Files in this skill

  • SKILL.md5.5 KB
  • _scores.json1.6 KB

Attribution

Is this your skill, or is something wrong with this listing? Request removal or report an issue. Author removals are honored within 72 hours.

Comments

Loading comments…