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
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
- **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)

### Fixed
- **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)
- **Mixed snapshot and non-snapshot pipelines**: a layer-level `next_snapshot_and_version` callback now suppresses input-view creation only for snapshot specs, allowing CloudFiles and other flows to share the pipeline group. [Issue #443](https://github.com/databrickslabs/sdp-meta/issues/443)
- **Append-flow source metadata**: nested Spark rows are normalized before serialization, so CloudFiles append flows can select file metadata columns correctly. [Issue #444](https://github.com/databrickslabs/sdp-meta/issues/444)
- **Append-flow custom transformations**: bronze and silver custom transformation functions now run for append-flow inputs as well as the primary input. [Issue #445](https://github.com/databrickslabs/sdp-meta/issues/445)
Expand Down
139 changes: 127 additions & 12 deletions demo/launch_dab_template_demo.py
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,7 @@
import sys
import tempfile
from dataclasses import dataclass
from io import StringIO
from pathlib import Path
from typing import Dict, List, Optional

Expand Down Expand Up @@ -1143,6 +1144,74 @@ def stage_add_independent_pipeline(
)


def stage_assert_duplicate_pipeline_rejected(
scenario: Scenario,
bundle_dir: Path,
*,
uc_schema: Optional[str],
) -> None:
"""Prove Issue #458 rejects duplicate ownership without changing files."""
if scenario.name != "multi_pipeline_cloudfiles":
return

_banner(
"STAGE 3C",
"bundle-add-pipeline duplicate ownership rejection (Issue #458)",
)

def snapshot() -> Dict[Path, bytes]:
return {
path.relative_to(bundle_dir): path.read_bytes()
for path in bundle_dir.rglob("*")
if path.is_file()
and ".git" not in path.relative_to(bundle_dir).parts
}

before = snapshot()
schemas = _resolve_demo_schemas(scenario, uc_schema)
output = StringIO()
rc = bundle_add_pipeline(
BundleAddPipelineCommand(
bundle_dir=str(bundle_dir),
pipeline=PipelineSpec(
name="rejected_duplicate_owner",
layer="bronze",
dataflow_group="dab_demo_cf_group",
bronze_target_schema=(
f"{schemas['bronze_target_schema']}_must_not_apply"
),
),
),
output=output,
)
message = output.getvalue()
after = snapshot()
changed = sorted(
str(path)
for path in set(before).union(after)
if before.get(path) != after.get(path)
)
if rc == 0:
raise SystemExit(
"Issue #458 regression: duplicate group/layer ownership was "
"accepted"
)
if "Duplicate bronze ownership" not in message:
raise SystemExit(
"Issue #458 regression: duplicate request failed for an "
f"unexpected reason:\n{message}"
)
if changed:
raise SystemExit(
"Issue #458 regression: rejected duplicate request modified "
f"bundle file(s): {changed}"
)
print(
"[STAGE 3C] Duplicate bronze ownership rejected and bundle files "
"remained byte-for-byte unchanged."
)


def stage_recipe(scenario: Scenario, bundle_dir: Path, *, apply_recipe: bool,
uc_catalog_name: str, profile: Optional[str],
uc_source_catalog: Optional[str] = None,
Expand Down Expand Up @@ -1587,6 +1656,50 @@ def _resolve_onboarding_file_name(bundle_dir: Path) -> str:
return name


def _set_onboarding_file_path(
onboarding_job_path: Path,
onboarding_path: str,
) -> None:
"""Update the onboarding task or fail clearly on hand-edited YAML."""
onboarding_job = yaml.safe_load(onboarding_job_path.read_text()) or {}
tasks = (
(((onboarding_job.get("resources") or {}).get("jobs") or {})
.get("onboarding") or {}).get("tasks")
)
if not isinstance(tasks, list):
raise ValueError(
f"{onboarding_job_path}: job `onboarding` must define a task list"
)
onboarding_task = next(
(
task for task in tasks
if isinstance(task, dict)
and task.get("task_key") == "onboard_dataflowspecs"
),
None,
)
if onboarding_task is None:
raise ValueError(
f"{onboarding_job_path}: missing task `onboard_dataflowspecs`"
)
wheel_task = onboarding_task.get("python_wheel_task")
if not isinstance(wheel_task, dict):
raise ValueError(
f"{onboarding_job_path}: task `onboard_dataflowspecs` is missing "
"`python_wheel_task`"
)
named_parameters = wheel_task.get("named_parameters")
if not isinstance(named_parameters, dict):
raise ValueError(
f"{onboarding_job_path}: task `onboard_dataflowspecs` is missing "
"`python_wheel_task.named_parameters`"
)
named_parameters["onboarding_file_path"] = onboarding_path
onboarding_job_path.write_text(
yaml.safe_dump(onboarding_job, sort_keys=False)
)


def stage_deploy_and_run(bundle_dir: Path, profile: Optional[str], *,
uc_catalog: Optional[str] = None,
uc_schema: Optional[str] = None,
Expand Down Expand Up @@ -1624,18 +1737,15 @@ def stage_deploy_and_run(bundle_dir: Path, profile: Optional[str], *,
onboarding_job_path = (
bundle_dir / "resources" / "sdp_meta_onboarding_job.yml"
)
onboarding_job = yaml.safe_load(onboarding_job_path.read_text())
tasks = onboarding_job["resources"]["jobs"]["onboarding"]["tasks"]
onboarding_task = next(
task for task in tasks
if task.get("task_key") == "onboard_dataflowspecs"
)
onboarding_task["python_wheel_task"]["named_parameters"][
"onboarding_file_path"
] = onboarding_path
onboarding_job_path.write_text(
yaml.safe_dump(onboarding_job, sort_keys=False)
)
try:
_set_onboarding_file_path(
onboarding_job_path,
onboarding_path,
)
except (OSError, ValueError, yaml.YAMLError) as exc:
raise SystemExit(
f"Cannot update onboarding file path: {exc}"
) from exc
print(
f"[STAGE 6] Onboarding task file path -> {onboarding_path}"
)
Expand Down Expand Up @@ -1933,6 +2043,11 @@ def main() -> int:
uc_schema=args.uc_schema,
demo_data_volume_path=demo_data_volume_path,
)
stage_assert_duplicate_pipeline_rejected(
scenario,
bundle_dir,
uc_schema=args.uc_schema,
)
stage_recipe(
scenario, bundle_dir,
apply_recipe=args.apply_recipe,
Expand Down
7 changes: 7 additions & 0 deletions docs/docs/reference/cli-commands.md
Original file line number Diff line number Diff line change
Expand Up @@ -101,6 +101,13 @@ the matching bronze task. It widens the onboarding job when another layer is
introduced; later `bundle-add-flow` calls for that group inherit the pipeline's
layer and target schemas.

Each `data_flow_group` and layer may have exactly one pipeline owner. A
bronze owner and silver owner may intentionally share a group as a split
topology, but duplicate bronze, duplicate silver, and combined/split overlap
are rejected before files are written. The preflight checks the complete
merged topology—including every bundle target override—so a pre-existing
ownership conflict must be repaired before another pipeline can be added.

## `bundle-validate`

Validates every configured pipeline independently, checks its layer-specific
Expand Down
Loading
Loading