Repository navigation
feat(pipelines): a Python SDK for running notebooks #505
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: feat/sdk-notebooks
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
Rebuilt onto the current feat/sdk-notebooks tip. Squashes the original commits of PR #505: the deepnote-sdk Python package. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
- Loading branch information
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,77 @@ | ||
| name: CD PyPI Python SDK | ||
|
|
||
| on: | ||
| push: | ||
| tags: | ||
| - 'python-sdk-v*' | ||
|
|
||
| permissions: | ||
| contents: read | ||
|
|
||
| jobs: | ||
| build: | ||
| name: Build sdist and wheel | ||
| runs-on: ubuntu-latest | ||
| timeout-minutes: 10 | ||
|
|
||
| steps: | ||
| - name: Checkout code | ||
| uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6 | ||
|
|
||
| - name: Setup Python | ||
| uses: actions/setup-python@ece7cb06caefa5fff74198d8649806c4678c61a1 # v6 | ||
| with: | ||
| python-version: '3.11' | ||
|
|
||
| - name: Check that the tag matches the package version | ||
| working-directory: packages/pipelines/python | ||
| run: | | ||
| version="$(python -c 'import tomllib, pathlib; print(tomllib.loads(pathlib.Path("pyproject.toml").read_text())["project"]["version"])')" | ||
| expected="python-sdk-v${version}" | ||
| if [[ "${GITHUB_REF_NAME}" != "${expected}" ]]; then | ||
| echo "Tag ${GITHUB_REF_NAME} does not match pyproject.toml version ${version} (expected ${expected})." >&2 | ||
| exit 1 | ||
| fi | ||
|
|
||
| - name: Install Python build tools | ||
| run: python -m pip install --upgrade pip build twine | ||
|
|
||
| - name: Build sdist and wheel | ||
| run: python -m build packages/pipelines/python --outdir packages/pipelines/python/dist | ||
|
|
||
| - name: Validate package metadata | ||
| run: python -m twine check packages/pipelines/python/dist/* | ||
|
|
||
| - name: Smoke test wheel | ||
| run: | | ||
| python -m pip install --force-reinstall packages/pipelines/python/dist/*.whl | ||
| python -c 'import deepnote, deepnote.sync; print(deepnote.Deepnote, deepnote.sync.Deepnote)' | ||
|
|
||
| - name: Upload distribution artifacts | ||
| uses: actions/upload-artifact@ea165f8d65b6e75b540449e92b4886f43607fa02 # v4 | ||
| with: | ||
| name: python-sdk-dist | ||
| path: packages/pipelines/python/dist/* | ||
| if-no-files-found: error | ||
|
|
||
| publish: | ||
| name: Publish to PyPI | ||
| needs: build | ||
| runs-on: ubuntu-latest | ||
| timeout-minutes: 10 | ||
| environment: release | ||
| permissions: | ||
| contents: read | ||
| id-token: write | ||
|
|
||
| steps: | ||
| - name: Download distribution artifacts | ||
| uses: actions/download-artifact@d3f86a106a0bac45b974a628896c90dbdf5c8093 # v4 | ||
| with: | ||
| name: python-sdk-dist | ||
| path: dist | ||
|
|
||
| - name: Publish to PyPI | ||
| uses: pypa/gh-action-pypi-publish@ed0c53931b1dc9bd32cbe73a98c7f6766f8a527e # release/v1 | ||
| with: | ||
| packages-dir: dist |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,26 @@ | ||
| # A pipeline in Python | ||
|
|
||
| The same fan-out-and-gate as [`../script`](../script) and [`../sdk`](../sdk), in Python. | ||
|
|
||
| ```bash | ||
| python -m pip install deepnote-sdk | ||
| DEEPNOTE_TOKEN=… NA_NOTEBOOK_ID=… EU_NOTEBOOK_ID=… APAC_NOTEBOOK_ID=… \ | ||
| python3 examples/pipelines/python/pipeline.py | ||
| ``` | ||
|
|
||
| The point of the example is what is missing from it. There is no workflow object, no step registry, | ||
| no DAG to declare: `asyncio.gather` fans out, a comprehension gates, `await` sequences. Python | ||
| executes the pipeline, and the SDK only makes each remote operation awaitable, typed, and named. | ||
|
|
||
| Two details worth reading in the source: | ||
|
|
||
| - **A failed notebook is an outcome, not an SDK error.** `gather(..., return_exceptions=True)` lets | ||
| every region finish, and `DeepnoteRunError` carries the whole result, so the handler can print what | ||
| the failing block actually said. | ||
| - **`outputs.last_json()` names no block id.** Deepnote reassigns block ids when it creates a | ||
| notebook from a file, so a binding that names one is fragile in exactly that case. | ||
|
|
||
| And the trade: because coordination lives in this process, it is not durable. Kill the script | ||
| mid-run and the notebook runs continue in Deepnote — they are detached, and their ids are printed — | ||
| but nothing aggregates them and no gate fires. Pick one back up with `deepnote.runs.get(id)`, or use | ||
| one of the durable options in the [package README](../../../packages/pipelines/python/README.md). | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,85 @@ | ||
| """The fan-out-and-gate pipeline in Python, with no pipeline API at all. | ||
|
|
||
| Every remote operation is awaitable, so the pipeline is the language: `asyncio.gather` fans out, a | ||
| comprehension gates, `await` sequences. Nothing here schedules, persists, replays, or supervises, | ||
| which is why the same file runs from cron, CI, a Lambda, or another Deepnote notebook. | ||
|
|
||
| DEEPNOTE_TOKEN=... NA_NOTEBOOK_ID=... EU_NOTEBOOK_ID=... APAC_NOTEBOOK_ID=... \ | ||
| python3 examples/pipelines/python/pipeline.py | ||
| """ | ||
|
|
||
| from __future__ import annotations | ||
|
|
||
| import asyncio | ||
| import os | ||
| import sys | ||
| import time | ||
|
|
||
| from deepnote import Deepnote, DeepnoteRunError, RunResult, outputs | ||
|
|
||
| QUALITY_THRESHOLD = 0.95 | ||
|
|
||
| REGIONS = [ | ||
| (name, os.environ.get(env)) | ||
| for name, env in ( | ||
| ("North America", "NA_NOTEBOOK_ID"), | ||
| ("Europe", "EU_NOTEBOOK_ID"), | ||
| ("Asia Pacific", "APAC_NOTEBOOK_ID"), | ||
| ) | ||
| ] | ||
|
|
||
|
|
||
| async def main() -> int: | ||
| regions = [(name, notebook_id) for name, notebook_id in REGIONS if notebook_id] | ||
| if not regions: | ||
| print("Set NA_NOTEBOOK_ID / EU_NOTEBOOK_ID / APAC_NOTEBOOK_ID to the notebooks to run.") | ||
| return 1 | ||
|
|
||
| started = time.monotonic() | ||
|
|
||
| async with Deepnote.from_env() as deepnote: | ||
| # A notebook plus the names of the values it publishes. `last_json` survives Deepnote | ||
| # reassigning block ids when it creates a notebook from a file. | ||
| def analysis(notebook_id: str): | ||
| return deepnote.notebooks.define(notebook_id, outputs={"reading": outputs.last_json()}) | ||
|
|
||
| # Fan out. Independent work is concurrent because `gather` is, not because a framework said | ||
| # so. `return_exceptions=True` keeps one failed region from cancelling the others: every | ||
| # notebook runs to its own conclusion and the failures are handled together below. | ||
| outcomes = await asyncio.gather( | ||
| *( | ||
| analysis(notebook_id).run_and_wait( | ||
| inputs={"region": name, "trailing_months": 6}, | ||
| on_status=lambda status, name=name: print(f" {name}: {status}"), | ||
| ) | ||
| for name, notebook_id in regions | ||
| ), | ||
| return_exceptions=True, | ||
| ) | ||
|
|
||
| results: dict[str, RunResult] = {} | ||
| failed = False | ||
| for (name, _), outcome in zip(regions, outcomes, strict=True): | ||
| if isinstance(outcome, DeepnoteRunError): | ||
| # A failed notebook is a real outcome, not an SDK error: the result carries the snapshot, | ||
| # which is where the failing block's own message is. | ||
| print(f" failed: run {outcome.run_id} — {outcome.result.error}") | ||
| failed = True | ||
| elif isinstance(outcome, BaseException): | ||
| raise outcome | ||
| else: | ||
| results[name] = outcome | ||
| if failed: | ||
| return 1 | ||
|
|
||
| # Gate. An ordinary comprehension over values that are already named. | ||
| below = [name for name, result in results.items() if result.values["reading"]["qualityScore"] < QUALITY_THRESHOLD] | ||
|
|
||
| print(f"\n {len(results)} regions in {time.monotonic() - started:.1f}s") | ||
| print(f" below threshold: {', '.join(below) or 'none'}") | ||
| print(f" runs: {', '.join(result.run_id for result in results.values())}\n") | ||
| return 0 | ||
|
|
||
|
|
||
| if __name__ == "__main__": | ||
| sys.exit(asyncio.run(main())) |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -78,6 +78,18 @@ Output "euShare" could not be read from totals.eu of block "stats-block" of run | |
| Deepnote reassigns block ids on creation. If Deepnote later grows a server-side notion of named | ||
| outputs, this surface does not change — only the resolver behind it does. | ||
|
|
||
| ### The same thing in Python | ||
|
|
||
| `packages/pipelines/python` is the Python SDK — the same two layers, the same boundary, Python idiom: | ||
|
|
||
| ```python | ||
| async with Deepnote.from_env() as deepnote: | ||
| result = await deepnote.notebooks["nb_extract"].run_and_wait(inputs={"region": "eu"}) | ||
|
Comment on lines
+86
to
+87
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🎯 Functional Correctness | 🟡 Minor | ⚡ Quick win Make the Python snippet executable. Line 81 uses Add 🤖 Prompt for AI Agents |
||
| ``` | ||
|
|
||
| It is published to PyPI as `deepnote-sdk` (`pip install deepnote-sdk`; the import name is `deepnote`). | ||
| See its [README](./python/README.md). | ||
|
|
||
| ## Run several notebooks as one pipeline | ||
|
|
||
| Fan out, gate on the results, decide: | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
🎯 Functional Correctness | 🟡 Minor | ⚡ Quick win
Run the example from the repository root.
Line 8 resolves
examples/pipelines/python/pipeline.pyfrompackages/pipelines/python, so Python cannot find the file. Install withpython -m pip install -e packages/pipelines/pythonfrom the repository root, then run the example path from that directory.Based on learnings: run commands from the repository root.
🤖 Prompt for AI Agents
Source: Learnings