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
Original file line number Diff line number Diff line change
Expand Up @@ -172,6 +172,15 @@ def to_asset_check_evaluation(
else:
target_materialization_data = None

if step_context.has_partition_key:
check_spec = assets_def_for_check.get_spec_for_check_key(check_key)
if check_spec.partitions_def is not None:
partition = step_context.partition_key
else:
partition = None
else:
partition = None

return AssetCheckEvaluation(
check_name=check_key.name,
asset_key=check_key.asset_key,
Expand All @@ -181,6 +190,7 @@ def to_asset_check_evaluation(
severity=self.severity,
description=self.description,
blocking=assets_def_for_check.get_spec_for_check_key(check_key).blocking,
partition=partition,
)

def with_metadata(self, metadata: Mapping[str, RawMetadataValue]) -> "AssetCheckResult": # pyright: ignore[reportIncompatibleMethodOverride]
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -10,9 +10,11 @@
replace,
)
from dagster_shared.serdes import whitelist_for_serdes
from dagster_shared.utils.warnings import preview_warning

from dagster._annotations import PublicAttr, public
from dagster._core.definitions.asset_key import AssetCheckKey, AssetKey, CoercibleToAssetKey
from dagster._core.definitions.partitions.definition import PartitionsDefinition

if TYPE_CHECKING:
from dagster._core.definitions.assets.definition.asset_dep import AssetDep, CoercibleToAssetDep
Expand Down Expand Up @@ -58,6 +60,7 @@ class AssetCheckSpec(IHaveNew, LegacyNamedTupleMixin):
blocking: PublicAttr[bool]
metadata: PublicAttr[Mapping[str, Any]]
automation_condition: PublicAttr[Optional[LazyAutomationCondition]]
partitions_def: PublicAttr[Optional[PartitionsDefinition]]

"""Defines information about an asset check, except how to execute it.

Expand All @@ -80,6 +83,9 @@ class AssetCheckSpec(IHaveNew, LegacyNamedTupleMixin):
that multi-asset is responsible for enforcing that downstream assets within the
same step do not execute after a blocking asset check fails.
metadata (Optional[Mapping[str, Any]]): A dict of static metadata for this asset check.
automation_condition (Optional[AutomationCondition[AssetCheckKey]]): The AutomationCondition for this asset check.
partitions_def (Optional[PartitionsDefinition]): The PartitionsDefinition for this asset check. Must be either None
or the same as the PartitionsDefinition of the asset specified by `asset`.
"""

def __new__(
Expand All @@ -92,11 +98,15 @@ def __new__(
blocking: bool = False,
metadata: Optional[Mapping[str, Any]] = None,
automation_condition: Optional["AutomationCondition[AssetCheckKey]"] = None,
partitions_def: Optional[PartitionsDefinition] = None,
):
from dagster._core.definitions.assets.definition.asset_dep import (
coerce_to_deps_and_check_duplicates,
)

if partitions_def is not None:
preview_warning("Specifying a partitions_def on an AssetCheckSpec")

asset_key = AssetKey.from_coercible_or_definition(asset)

additional_asset_deps = coerce_to_deps_and_check_duplicates(
Expand All @@ -119,6 +129,7 @@ def __new__(
blocking=blocking,
metadata=metadata or {},
automation_condition=automation_condition,
partitions_def=partitions_def,
)

def get_python_identifier(self) -> str:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -189,6 +189,7 @@ def __init__(
v.get_spec_for_check_key(k).description,
v.get_spec_for_check_key(k).automation_condition,
v.get_spec_for_check_key(k).metadata,
v.get_spec_for_check_key(k).partitions_def,
)
for k, v in assets_defs_by_check_key.items()
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -217,13 +217,15 @@ def __init__(
description: Optional[str],
automation_condition: Optional["AutomationCondition[AssetCheckKey]"],
metadata: ArbitraryMetadataMapping,
partitions_def: Optional[PartitionsDefinition],
):
self.key = key
self.blocking = blocking
self._automation_condition = automation_condition
self._additional_deps = additional_deps
self._description = description
self._metadata = metadata
self._partitions_def = partitions_def

@property
def parent_entity_keys(self) -> AbstractSet[AssetKey]:
Expand All @@ -235,8 +237,7 @@ def child_entity_keys(self) -> AbstractSet[EntityKey]:

@property
def partitions_def(self) -> Optional[PartitionsDefinition]:
# all checks are unpartitioned
return None
return self._partitions_def

@property
def partition_mappings(self) -> Mapping[EntityKey, PartitionMapping]:
Expand Down Expand Up @@ -266,6 +267,10 @@ class BaseAssetGraph(ABC, Generic[T_AssetNode]):
def asset_nodes(self) -> Iterable[T_AssetNode]:
return self._asset_nodes_by_key.values()

@property
def asset_check_nodes(self) -> Iterable[AssetCheckNode]:
return self._asset_check_nodes_by_key.values()

@property
def nodes(self) -> Iterable[BaseEntityNode]:
return [
Expand Down Expand Up @@ -668,6 +673,28 @@ def validate_partitions(self):
f"Invalid partition mapping from {node.key.to_user_string()} to {parent.key.to_user_string()}"
) from e

# Validate that asset checks have compatible partitions_def with their target asset
for node in self.asset_check_nodes:
if node.partitions_def is None:
continue

target_asset_key = node.key.asset_key
if not self.has(target_asset_key):
raise DagsterInvalidDefinitionError(
f"Partitioned asset check '{node.key.to_user_string()}' targets "
f"asset '{target_asset_key.to_user_string()}' "
"but the asset does not exist in the graph."
)
# If the check is partitioned, it must have the same partitions_def as the asset
if node.partitions_def != self.get(target_asset_key).partitions_def:
raise DagsterInvalidDefinitionError(
f"Asset check '{node.key.to_user_string()}' targets asset '{target_asset_key.to_user_string()}' "
"but has a different partitions definition. "
f"Asset check partitions_def: {node.partitions_def}, "
f"Asset partitions_def: {self.get(target_asset_key).partitions_def}. "
"Partitioned asset checks must have the same partitions definition as their target asset."
)

def upstream_key_iterator(self, asset_key: AssetKey) -> Iterator[AssetKey]:
"""Iterates through all asset keys which are upstream of the given key."""
visited: set[AssetKey] = set()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -480,6 +480,9 @@ def _get_asset_check_node_from_remote_asset_check_node(
remote_node.asset_check.description,
remote_node.asset_check.automation_condition,
{}, # metadata not yet on AssetCheckNodeSnap
remote_node.asset_check.partitions_def_snapshot.get_partitions_definition()
if remote_node.asset_check.partitions_def_snapshot
else None,
)

##### COMMON ASSET GRAPH INTERFACE
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,9 @@
from dagster._core.definitions.decorators.op_decorator import _Op
from dagster._core.definitions.events import AssetKey, CoercibleToAssetKey
from dagster._core.definitions.output import Out
from dagster._core.definitions.partitions.definition.partitions_definition import (
PartitionsDefinition,
)
from dagster._core.definitions.policy import RetryPolicy
from dagster._core.definitions.source_asset import SourceAsset
from dagster._core.definitions.utils import DEFAULT_OUTPUT
Expand Down Expand Up @@ -113,6 +116,7 @@ def asset_check(
metadata: Optional[Mapping[str, Any]] = None,
automation_condition: Optional[AutomationCondition[AssetCheckKey]] = None,
pool: Optional[str] = None,
partitions_def: Optional[PartitionsDefinition] = None,
) -> Callable[[AssetCheckFunction], AssetChecksDefinition]:
"""Create a definition for how to execute an asset check.

Expand Down Expand Up @@ -151,6 +155,7 @@ def asset_check(
automation_condition (Optional[AutomationCondition]): An AutomationCondition which determines
when this check should be executed.
pool (Optional[str]): A string that identifies the concurrency pool that governs this asset check's execution.
partitions_def (Optional[PartitionsDefinition]): The PartitionsDefinition for this asset check.

Produces an :py:class:`AssetChecksDefinition` object.

Expand Down Expand Up @@ -218,6 +223,7 @@ def inner(fn: AssetCheckFunction) -> AssetChecksDefinition:
blocking=blocking,
metadata=metadata,
automation_condition=automation_condition,
partitions_def=partitions_def,
)

resource_defs_for_execution = wrap_resources_for_execution(resource_defs)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,8 @@
from dagster._core.definitions.asset_checks.asset_check_spec import AssetCheckKey, AssetCheckSpec
from dagster._core.definitions.job_definition import JobDefinition
from dagster._core.definitions.op_definition import OpDefinition
from dagster._core.definitions.partitions.partition_key_range import PartitionKeyRange
from dagster._core.definitions.partitions.utils.time_window import TimeWindow
from dagster._core.definitions.repository_definition.repository_definition import (
RepositoryDefinition,
)
Expand Down Expand Up @@ -135,6 +137,43 @@ def step_launcher(self) -> Optional[StepLauncher]:
def get_step_execution_context(self) -> StepExecutionContext:
return self.op_execution_context.get_step_execution_context()

#### partition related
@public
@property
@_copy_docs_from_op_execution_context
def has_partition_key(self) -> bool:
return self.op_execution_context.has_partition_key

@public
@property
@_copy_docs_from_op_execution_context
def partition_key(self) -> str:
return self.op_execution_context.partition_key

@public
@property
@_copy_docs_from_op_execution_context
def partition_keys(self) -> Sequence[str]:
return self.op_execution_context.partition_keys

@public
@property
@_copy_docs_from_op_execution_context
def has_partition_key_range(self) -> bool:
return self.op_execution_context.has_partition_key_range

@public
@property
@_copy_docs_from_op_execution_context
def partition_key_range(self) -> PartitionKeyRange:
return self.op_execution_context.partition_key_range

@public
@property
@_copy_docs_from_op_execution_context
def partition_time_window(self) -> TimeWindow:
return self.op_execution_context.partition_time_window
Comment thread
OwenKephart marked this conversation as resolved.

# misc

@public
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -899,6 +899,7 @@ class AssetCheckNodeSnap(IHaveNew):
additional_asset_keys: Sequence[AssetKey]
automation_condition: Optional[AutomationCondition]
automation_condition_snapshot: Optional[AutomationConditionSnapshot]
partitions_def_snapshot: Optional[PartitionsSnap]

def __new__(
cls,
Expand All @@ -911,6 +912,7 @@ def __new__(
additional_asset_keys: Optional[Sequence[AssetKey]] = None,
automation_condition: Optional[AutomationCondition] = None,
automation_condition_snapshot: Optional[AutomationConditionSnapshot] = None,
partitions_def_snapshot: Optional[PartitionsSnap] = None,
):
return super().__new__(
cls,
Expand All @@ -923,6 +925,7 @@ def __new__(
additional_asset_keys=additional_asset_keys or [],
automation_condition=automation_condition,
automation_condition_snapshot=automation_condition_snapshot,
partitions_def_snapshot=partitions_def_snapshot,
)

@property
Expand Down Expand Up @@ -1211,6 +1214,9 @@ def asset_check_node_snaps_from_repo(repo: RepositoryDefinition) -> Sequence[Ass
additional_asset_keys=[dep.asset_key for dep in spec.additional_deps],
automation_condition=automation_condition,
automation_condition_snapshot=automation_condition_snapshot,
partitions_def_snapshot=PartitionsSnap.from_def(spec.partitions_def)
if spec.partitions_def
else None,
)
)

Expand Down
Loading