Skip to content
Merged
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
3 changes: 3 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -3,9 +3,12 @@
## [v0.1.1]
### Added
- **Multiple pipeline topologies per DAB bundle**: `bundle-add-pipeline` adds independently configured bronze, silver, split, or combined pipelines with their own data-flow groups and target schemas. `bundle-validate` validates each pipeline and its job wiring independently. [Issue #446](https://github.com/databrickslabs/sdp-meta/issues/446)
- **Auto Loader schema-evolution demos**: focused and at-scale demos exercise additive schema evolution, strict-mode failure behavior, and multi-table onboarding from Unity Catalog Volumes. [PR #490](https://github.com/databrickslabs/sdp-meta/pull/490)

### Fixed
- **Databricks SDK compatibility floor**: the minimum supported SDK is now `databricks-sdk>=0.138.0,<1` across wheel, runtime, and Databricks App dependencies. CI installs exactly 0.138.0 and exercises the serverless `jobs.JobEnvironment` and `compute.Environment` models used by onboarding jobs. [Issue #457](https://github.com/databrickslabs/sdp-meta/issues/457)
- **[Databricks App] Make demo launches asynchronous and feature the 100-table Auto Loader demo**. [Issue #491](https://github.com/databrickslabs/sdp-meta/issues/491)
- **Data-quality expectations without quarantine targets**: bronze and silver flows can apply DQ expectations when no quarantine table is configured, without referencing a missing quarantine target. [PR #487](https://github.com/databrickslabs/sdp-meta/pull/487)
- **Databricks App deployment packaging**: macOS/Linux and Windows deployment staging now includes `examples/`, allowing the App container to build the SDP-META wheel after packaged MCP examples became a required wheel input.
- **DAB pipeline group ownership**: `bundle-add-pipeline` now rejects duplicate ownership of the same data-flow group and layer before writing files, while preserving supported bronze/silver split ownership and wiring its job dependency. Its fail-safe preflight checks the complete merged topology and every target override. `bundle-validate` also detects ownership conflicts introduced through manual YAML edits, including conflicts that exist only under a target override. [Issue #458](https://github.com/databrickslabs/sdp-meta/issues/458)
- **Customized DAB validation**: `bundle-validate` now honors default and environment-selected targets, accepts transitive bronze-to-silver task dependencies, and reports malformed targets or dependency cycles cleanly. [Issue #459](https://github.com/databrickslabs/sdp-meta/issues/459)
Expand Down
3 changes: 3 additions & 0 deletions databricks_app/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -243,6 +243,9 @@ export FLASK_DEBUG=true
# Optional: pick a non-default profile (matches the Apps "--profile" semantic)
# export DATABRICKS_CONFIG_PROFILE=<name>

# Optional: concurrent demo launch limit (default: 8)
# export SDP_META_MAX_ACTIVE_DEMOS=8

flask --app databricks_app/app.py run --host 127.0.0.1 --port 8000
```

Expand Down
140 changes: 127 additions & 13 deletions databricks_app/_jobs.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,27 +19,141 @@

from __future__ import annotations

import os
import time
import threading
import uuid


_jobs: dict = {}
_jobs_lock = threading.RLock()
_MAX_LOG_LINES = 5000
_MAX_COMPLETED_JOBS = 50
_COMPLETED_JOB_TTL_SECONDS = 60 * 60


def _new_job_token() -> str:
def _positive_int_env(name: str, default: int) -> int:
"""Read a positive integer environment setting with a safe fallback."""
try:
value = int(os.environ.get(name, str(default)))
except (TypeError, ValueError):
return default
return value if value > 0 else default


_MAX_ACTIVE_DEMO_JOBS = _positive_int_env(
'SDP_META_MAX_ACTIVE_DEMOS',
8,
)


class JobCapacityError(RuntimeError):
"""Raised when a bounded job category has no free launch slots."""


def _prune_jobs(now: float | None = None) -> None:
"""Evict expired and excess completed jobs from the process store."""
with _jobs_lock:
now = time.time() if now is None else now
stale = [
token
for token, job in _jobs.items()
if job.get('done')
and now - job.get('finished_at', now) > _COMPLETED_JOB_TTL_SECONDS
]
for token in stale:
_jobs.pop(token, None)

completed = sorted(
(
(job.get('finished_at', 0), token)
for token, job in _jobs.items()
if job.get('done')
),
reverse=True,
)
for _, token in completed[_MAX_COMPLETED_JOBS:]:
_jobs.pop(token, None)


def _new_job_token(
*,
kind: str = 'generic',
max_active_for_kind: int | None = None,
) -> str:
"""Allocate a fresh job-token entry in ``_jobs`` and return the token."""
token = uuid.uuid4().hex
_jobs[token] = {
'logs': [],
'done': False,
'returncode': None,
'stdout': '',
'stderr': '',
'modal_content': None,
'error': None,
}
return token
with _jobs_lock:
_prune_jobs()
if max_active_for_kind is not None:
active_count = sum(
1
for job in _jobs.values()
if not job.get('done') and job.get('kind') == kind
)
if active_count >= max_active_for_kind:
raise JobCapacityError(
f"Too many active {kind} jobs "
f"({active_count}/{max_active_for_kind})"
)
token = uuid.uuid4().hex
_jobs[token] = {
'kind': kind,
'logs': [],
'log_base_offset': 0,
'done': False,
'returncode': None,
'stdout': '',
'stderr': '',
'modal_content': None,
'error': None,
'created_at': time.time(),
'finished_at': None,
}
return token


def _append_job_log(job: dict, entry: dict) -> None:
"""Append one line while bounding retained per-job output."""
with _jobs_lock:
if len(job['logs']) >= _MAX_LOG_LINES:
job['logs'].pop(0)
job['log_base_offset'] += 1
job['logs'].append(entry)


def _update_job(token: str, **values) -> None:
"""Atomically update fields on an existing job."""
with _jobs_lock:
job = _jobs.get(token)
if job is not None:
job.update(values)


def _mark_job_done(token: str) -> None:
"""Finalize a job and enforce completed-job retention limits."""
with _jobs_lock:
job = _jobs.get(token)
if job is None:
return
job['done'] = True
job['finished_at'] = time.time()
_prune_jobs()


def _get_job(token: str):
"""Look up a job entry by token, returning ``None`` if absent."""
return _jobs.get(token)
with _jobs_lock:
_prune_jobs()
return _jobs.get(token)


def _get_job_snapshot(token: str):
"""Return a stable copy suitable for request handlers."""
with _jobs_lock:
_prune_jobs()
job = _jobs.get(token)
if job is None:
return None
snapshot = dict(job)
snapshot['logs'] = [dict(entry) for entry in job['logs']]
return snapshot
Loading
Loading