Skip to content

Commit 33ba92c

Browse files
prhaclaude
andcommitted
Add core definitions and API for partitioned asset checks
This PR adds the core API layer for partitioned asset checks: - Add partitions_def field to AssetCheckSpec - Update asset_check decorator to support partitions_def parameter - Add partition_key property to AssetCheckExecutionContext - Update AssetCheckEvaluation and AssetCheckResult to include partition info - Add event definitions for partitioned check evaluations - Update external data representation for remote execution 🤖 Generated with [Claude Code](https://claude.com/claude-code) Co-Authored-By: Claude <noreply@anthropic.com>
1 parent 6db61e3 commit 33ba92c

7 files changed

Lines changed: 180 additions & 71 deletions

File tree

python_modules/dagster/dagster/_core/definitions/asset_checks/asset_check_result.py

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -172,6 +172,21 @@ def to_asset_check_evaluation(
172172
else:
173173
target_materialization_data = None
174174

175+
if step_context.has_partition_key:
176+
assets_def = step_context.job_def.asset_layer.get_assets_def_for_node(
177+
step_context.node_handle
178+
)
179+
assert assets_def
180+
spec = assets_def_for_check.get_spec_for_check_key(check_key)
181+
if spec.partitions_def is not None and spec.partitions_def.has_partition_key(
182+
step_context.partition_key
183+
):
184+
partition = step_context.partition_key
185+
else:
186+
partition = None
187+
else:
188+
partition = None
189+
175190
return AssetCheckEvaluation(
176191
check_name=check_key.name,
177192
asset_key=check_key.asset_key,
@@ -181,6 +196,7 @@ def to_asset_check_evaluation(
181196
severity=self.severity,
182197
description=self.description,
183198
blocking=assets_def_for_check.get_spec_for_check_key(check_key).blocking,
199+
partition=partition,
184200
)
185201

186202
def with_metadata(self, metadata: Mapping[str, RawMetadataValue]) -> "AssetCheckResult": # pyright: ignore[reportIncompatibleMethodOverride]

python_modules/dagster/dagster/_core/definitions/asset_checks/asset_check_spec.py

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,7 @@
1313

1414
from dagster._annotations import PublicAttr, public
1515
from dagster._core.definitions.asset_key import AssetCheckKey, AssetKey, CoercibleToAssetKey
16+
from dagster._core.definitions.partitions.definition import PartitionsDefinition
1617

1718
if TYPE_CHECKING:
1819
from dagster._core.definitions.assets.definition.asset_dep import AssetDep, CoercibleToAssetDep
@@ -58,6 +59,7 @@ class AssetCheckSpec(IHaveNew, LegacyNamedTupleMixin):
5859
blocking: PublicAttr[bool]
5960
metadata: PublicAttr[Mapping[str, Any]]
6061
automation_condition: PublicAttr[Optional[LazyAutomationCondition]]
62+
partitions_def: PublicAttr[Optional[PartitionsDefinition]]
6163

6264
"""Defines information about an asset check, except how to execute it.
6365
@@ -80,6 +82,7 @@ class AssetCheckSpec(IHaveNew, LegacyNamedTupleMixin):
8082
that multi-asset is responsible for enforcing that downstream assets within the
8183
same step do not execute after a blocking asset check fails.
8284
metadata (Optional[Mapping[str, Any]]): A dict of static metadata for this asset check.
85+
partitions_def (Optional[PartitionsDefinition]): The partitions definition for this asset check.
8386
"""
8487

8588
def __new__(
@@ -92,6 +95,7 @@ def __new__(
9295
blocking: bool = False,
9396
metadata: Optional[Mapping[str, Any]] = None,
9497
automation_condition: Optional["AutomationCondition[AssetCheckKey]"] = None,
98+
partitions_def: Optional[PartitionsDefinition] = None,
9599
):
96100
from dagster._core.definitions.assets.definition.asset_dep import (
97101
coerce_to_deps_and_check_duplicates,
@@ -119,6 +123,7 @@ def __new__(
119123
blocking=blocking,
120124
metadata=metadata or {},
121125
automation_condition=automation_condition,
126+
partitions_def=partitions_def,
122127
)
123128

124129
def get_python_identifier(self) -> str:

python_modules/dagster/dagster/_core/definitions/decorators/asset_check_decorator.py

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@
2828
from dagster._core.definitions.decorators.op_decorator import _Op
2929
from dagster._core.definitions.events import AssetKey, CoercibleToAssetKey
3030
from dagster._core.definitions.output import Out
31+
from dagster._core.definitions.partitions.definition import PartitionsDefinition
3132
from dagster._core.definitions.policy import RetryPolicy
3233
from dagster._core.definitions.source_asset import SourceAsset
3334
from dagster._core.definitions.utils import DEFAULT_OUTPUT
@@ -113,6 +114,7 @@ def asset_check(
113114
metadata: Optional[Mapping[str, Any]] = None,
114115
automation_condition: Optional[AutomationCondition[AssetCheckKey]] = None,
115116
pool: Optional[str] = None,
117+
partitions_def: Optional[PartitionsDefinition] = None,
116118
) -> Callable[[AssetCheckFunction], AssetChecksDefinition]:
117119
"""Create a definition for how to execute an asset check.
118120
@@ -218,6 +220,7 @@ def inner(fn: AssetCheckFunction) -> AssetChecksDefinition:
218220
blocking=blocking,
219221
metadata=metadata,
220222
automation_condition=automation_condition,
223+
partitions_def=partitions_def,
221224
)
222225

223226
resource_defs_for_execution = wrap_resources_for_execution(resource_defs)

python_modules/dagster/dagster/_core/events/__init__.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -39,6 +39,7 @@
3939
AssetCheckEvaluation,
4040
AssetCheckEvaluationPlanned,
4141
)
42+
from dagster._core.definitions.asset_checks.asset_check_spec import AssetCheckKey
4243
from dagster._core.definitions.asset_health.asset_health import AssetHealthStatus
4344
from dagster._core.definitions.events import (
4445
AssetLineageInfo,

python_modules/dagster/dagster/_core/execution/context/asset_check_execution_context.py

Lines changed: 39 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,8 @@
66
from dagster._core.definitions.asset_checks.asset_check_spec import AssetCheckKey, AssetCheckSpec
77
from dagster._core.definitions.job_definition import JobDefinition
88
from dagster._core.definitions.op_definition import OpDefinition
9+
from dagster._core.definitions.partitions.partition_key_range import PartitionKeyRange
10+
from dagster._core.definitions.partitions.utils import TimeWindow
911
from dagster._core.definitions.repository_definition.repository_definition import (
1012
RepositoryDefinition,
1113
)
@@ -135,6 +137,43 @@ def step_launcher(self) -> Optional[StepLauncher]:
135137
def get_step_execution_context(self) -> StepExecutionContext:
136138
return self.op_execution_context.get_step_execution_context()
137139

140+
#### partition related
141+
@public
142+
@property
143+
@_copy_docs_from_op_execution_context
144+
def has_partition_key(self) -> bool:
145+
return self.op_execution_context.has_partition_key
146+
147+
@public
148+
@property
149+
@_copy_docs_from_op_execution_context
150+
def partition_key(self) -> str:
151+
return self.op_execution_context.partition_key
152+
153+
@public
154+
@property
155+
@_copy_docs_from_op_execution_context
156+
def partition_keys(self) -> Sequence[str]:
157+
return self.op_execution_context.partition_keys
158+
159+
@public
160+
@property
161+
@_copy_docs_from_op_execution_context
162+
def has_partition_key_range(self) -> bool:
163+
return self.op_execution_context.has_partition_key_range
164+
165+
@public
166+
@property
167+
@_copy_docs_from_op_execution_context
168+
def partition_key_range(self) -> PartitionKeyRange:
169+
return self.op_execution_context.partition_key_range
170+
171+
@public
172+
@property
173+
@_copy_docs_from_op_execution_context
174+
def partition_time_window(self) -> TimeWindow:
175+
return self.op_execution_context.partition_time_window
176+
138177
# misc
139178

140179
@public

0 commit comments

Comments
 (0)