Skip to content

Commit bf76c19

Browse files
committed
[pac] Add partitions_subset to AssetCheckEvaluationPlanned, partition to AssetCheckEvaluation
1 parent 36b1215 commit bf76c19

10 files changed

Lines changed: 500 additions & 84 deletions

File tree

python_modules/dagster/dagster/_core/asset_graph_view/serializable_entity_subset.py

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -50,6 +50,12 @@ class SerializableEntitySubset(Generic[T_EntityKey]):
5050
key: T_EntityKey
5151
value: EntitySubsetValue
5252

53+
@classmethod
54+
def empty(
55+
cls, key: T_EntityKey, partitions_def: Optional[PartitionsDefinition]
56+
) -> "SerializableEntitySubset[T_EntityKey]":
57+
return cls(key=key, value=partitions_def.empty_subset() if partitions_def else False)
58+
5359
@classmethod
5460
def from_coercible_value(
5561
cls,

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

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@
99
)
1010
from dagster._core.definitions.events import AssetKey, MetadataValue, RawMetadataValue
1111
from dagster._core.definitions.metadata import normalize_metadata
12+
from dagster._core.definitions.partitions.subset import PartitionsSubset
1213
from dagster._serdes import whitelist_for_serdes
1314

1415

@@ -19,6 +20,7 @@ class AssetCheckEvaluationPlanned:
1920

2021
asset_key: AssetKey
2122
check_name: str
23+
partitions_subset: Optional[PartitionsSubset] = None
2224

2325
@property
2426
def asset_check_key(self) -> AssetCheckKey:
@@ -59,6 +61,8 @@ class AssetCheckEvaluation(IHaveNew):
5961
A text description of the result of the check evaluation.
6062
blocking (Optional[bool]):
6163
Whether the check is blocking.
64+
partition (Optional[str]):
65+
The partition that the check was evaluated on, if applicable.
6266
"""
6367

6468
asset_key: AssetKey
@@ -69,6 +73,7 @@ class AssetCheckEvaluation(IHaveNew):
6973
severity: AssetCheckSeverity
7074
description: Optional[str]
7175
blocking: Optional[bool]
76+
partition: Optional[str]
7277

7378
def __new__(
7479
cls,
@@ -80,6 +85,7 @@ def __new__(
8085
severity: AssetCheckSeverity = AssetCheckSeverity.ERROR,
8186
description: Optional[str] = None,
8287
blocking: Optional[bool] = None,
88+
partition: Optional[str] = None,
8389
):
8490
return super().__new__(
8591
cls,
@@ -91,6 +97,7 @@ def __new__(
9197
severity=severity,
9298
description=description,
9399
blocking=blocking,
100+
partition=partition,
94101
)
95102

96103
@property

python_modules/dagster/dagster/_core/definitions/declarative_automation/legacy/valid_asset_subset.py

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -98,14 +98,14 @@ def all(
9898
key=asset_key, value=AllPartitionsSubset(partitions_def, ctx)
9999
)
100100

101-
@staticmethod
101+
@classmethod
102102
def empty(
103-
asset_key: AssetKey, partitions_def: Optional[PartitionsDefinition]
103+
cls, key: AssetKey, partitions_def: Optional[PartitionsDefinition]
104104
) -> "ValidAssetSubset":
105105
if partitions_def is None:
106-
return ValidAssetSubset(key=asset_key, value=False)
106+
return cls(key=key, value=False)
107107
else:
108-
return ValidAssetSubset(key=asset_key, value=partitions_def.empty_subset())
108+
return cls(key=key, value=partitions_def.empty_subset())
109109

110110
@staticmethod
111111
def from_asset_partitions_set(

python_modules/dagster/dagster/_core/instance/runs/run_domain.py

Lines changed: 43 additions & 69 deletions
Original file line numberDiff line numberDiff line change
@@ -830,6 +830,15 @@ def _log_asset_planned_events(
830830
target_asset_key = asset_check_key.asset_key
831831
check_name = asset_check_key.name
832832

833+
partitions_def = (
834+
asset_graph.get(asset_check_key).partitions_def
835+
if asset_graph.get(asset_check_key)
836+
else None
837+
)
838+
partitions_subset = dagster_run.get_resolved_partitions_subset(
839+
partitions_def
840+
)
841+
833842
event = DagsterEvent(
834843
event_type_value=DagsterEventType.ASSET_CHECK_EVALUATION_PLANNED.value,
835844
job_name=job_name,
@@ -838,8 +847,9 @@ def _log_asset_planned_events(
838847
f" asset {target_asset_key.to_string()}"
839848
),
840849
event_specific_data=AssetCheckEvaluationPlanned(
841-
target_asset_key,
850+
asset_key=target_asset_key,
842851
check_name=check_name,
852+
partitions_subset=partitions_subset,
843853
),
844854
step_key=step.key,
845855
)
@@ -883,16 +893,13 @@ def get_materialization_planned_events_for_asset(
883893
) -> Sequence["DagsterEvent"]:
884894
"""Moved from DagsterInstance._log_materialization_planned_event_for_asset."""
885895
from dagster._core.definitions.partitions.context import partition_loading_context
886-
from dagster._core.definitions.partitions.definition import DynamicPartitionsDefinition
887896
from dagster._core.events import AssetMaterializationPlannedData, DagsterEvent
888897

889898
events = []
890899

891900
partition_tag = dagster_run.tags.get(PARTITION_NAME_TAG)
892-
partition_range_start, partition_range_end = (
893-
dagster_run.tags.get(ASSET_PARTITION_RANGE_START_TAG),
894-
dagster_run.tags.get(ASSET_PARTITION_RANGE_END_TAG),
895-
)
901+
partition_range_start = dagster_run.tags.get(ASSET_PARTITION_RANGE_START_TAG)
902+
partition_range_end = dagster_run.tags.get(ASSET_PARTITION_RANGE_END_TAG)
896903

897904
if partition_tag and (partition_range_start or partition_range_end):
898905
raise DagsterInvariantViolationError(
@@ -901,79 +908,46 @@ def get_materialization_planned_events_for_asset(
901908
f" {PARTITION_NAME_TAG}"
902909
)
903910

904-
partitions_subset = None
905-
individual_partitions = None
906-
if partition_range_start or partition_range_end:
907-
if not partition_range_start or not partition_range_end:
908-
raise DagsterInvariantViolationError(
909-
f"Cannot have {ASSET_PARTITION_RANGE_START_TAG} or"
910-
f" {ASSET_PARTITION_RANGE_END_TAG} set without the other"
911-
)
911+
asset_node = self._get_repo_scoped_asset_node(
912+
asset_graph, asset_key, dagster_run.remote_job_origin
913+
)
914+
partitions_def = asset_node.partitions_def if asset_node else None
912915

913-
asset_node = check.not_none(
914-
self._get_repo_scoped_asset_node(
915-
asset_graph, asset_key, dagster_run.remote_job_origin
916-
)
917-
)
916+
partitions_subset = dagster_run.get_resolved_partitions_subset(partitions_def)
918917

919-
partitions_def = asset_node.partitions_def
920-
if (
921-
isinstance(partitions_def, DynamicPartitionsDefinition)
922-
and partitions_def.name is None
923-
):
924-
raise DagsterInvariantViolationError(
925-
"Creating a run targeting a partition range is not supported for assets partitioned with function-based dynamic partitions"
918+
with partition_loading_context(dynamic_partitions_store=self._instance):
919+
if partitions_subset is None:
920+
materialization_planned = DagsterEvent.build_asset_materialization_planned_event(
921+
job_name,
922+
step.key,
923+
AssetMaterializationPlannedData(
924+
asset_key, partition=None, partitions_subset=None
925+
),
926926
)
927-
928-
if partitions_def is not None:
929-
with partition_loading_context(dynamic_partitions_store=self._instance):
930-
if self._instance.event_log_storage.supports_partition_subset_in_asset_materialization_planned_events:
931-
partitions_subset = partitions_def.subset_with_partition_keys(
932-
partitions_def.get_partition_keys_in_range(
933-
PartitionKeyRange(partition_range_start, partition_range_end),
934-
)
935-
).to_serializable_subset()
936-
individual_partitions = []
937-
else:
938-
individual_partitions = partitions_def.get_partition_keys_in_range(
939-
PartitionKeyRange(partition_range_start, partition_range_end),
940-
)
941-
elif check.not_none(output.properties).is_asset_partitioned and partition_tag:
942-
individual_partitions = [partition_tag]
943-
944-
assert not (individual_partitions and partitions_subset), (
945-
"Should set either individual_partitions or partitions_subset, but not both"
946-
)
947-
948-
if not individual_partitions and not partitions_subset:
949-
materialization_planned = DagsterEvent.build_asset_materialization_planned_event(
950-
job_name,
951-
step.key,
952-
AssetMaterializationPlannedData(asset_key, partition=None, partitions_subset=None),
953-
)
954-
events.append(materialization_planned)
955-
elif individual_partitions:
956-
for individual_partition in individual_partitions:
927+
events.append(materialization_planned)
928+
elif self._instance.event_log_storage.supports_partition_subset_in_asset_materialization_planned_events:
957929
materialization_planned = DagsterEvent.build_asset_materialization_planned_event(
958930
job_name,
959931
step.key,
960932
AssetMaterializationPlannedData(
961-
asset_key,
962-
partition=individual_partition,
963-
partitions_subset=partitions_subset,
933+
asset_key, partition=None, partitions_subset=partitions_subset
964934
),
965935
)
966936
events.append(materialization_planned)
967-
968-
else:
969-
materialization_planned = DagsterEvent.build_asset_materialization_planned_event(
970-
job_name,
971-
step.key,
972-
AssetMaterializationPlannedData(
973-
asset_key, partition=None, partitions_subset=partitions_subset
974-
),
975-
)
976-
events.append(materialization_planned)
937+
else:
938+
for partition_key in partitions_subset.get_partition_keys():
939+
materialization_planned = (
940+
DagsterEvent.build_asset_materialization_planned_event(
941+
job_name,
942+
step.key,
943+
AssetMaterializationPlannedData(
944+
asset_key,
945+
partition=partition_key,
946+
partitions_subset=None,
947+
),
948+
)
949+
)
950+
events.append(materialization_planned)
977951

978952
return events
979953

python_modules/dagster/dagster/_core/storage/asset_check_execution_record.py

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -60,6 +60,7 @@ class AssetCheckExecutionRecord(
6060
# Old records won't have an event if the status is PLANNED.
6161
("event", Optional[EventLogEntry]),
6262
("create_timestamp", float),
63+
("partition", Optional[str]),
6364
],
6465
),
6566
LoadableBy[AssetCheckKey],
@@ -72,13 +73,15 @@ def __new__(
7273
status: AssetCheckExecutionRecordStatus,
7374
event: Optional[EventLogEntry],
7475
create_timestamp: float,
76+
partition: Optional[str],
7577
):
7678
check.inst_param(key, "key", AssetCheckKey)
7779
check.int_param(id, "id")
7880
check.str_param(run_id, "run_id")
7981
check.inst_param(status, "status", AssetCheckExecutionRecordStatus)
8082
check.opt_inst_param(event, "event", EventLogEntry)
8183
check.float_param(create_timestamp, "create_timestamp")
84+
check.opt_str_param(partition, "partition")
8285

8386
event_type = event.dagster_event_type if event else None
8487
if status == AssetCheckExecutionRecordStatus.PLANNED:
@@ -105,6 +108,7 @@ def __new__(
105108
status=status,
106109
event=event,
107110
create_timestamp=create_timestamp,
111+
partition=partition,
108112
)
109113

110114
@property
@@ -129,6 +133,7 @@ def from_db_row(cls, row, key: AssetCheckKey) -> "AssetCheckExecutionRecord":
129133
else None
130134
),
131135
create_timestamp=utc_datetime_from_naive(row["create_timestamp"]).timestamp(),
136+
partition=row["partition"],
132137
)
133138

134139
@classmethod

python_modules/dagster/dagster/_core/storage/dagster_run.py

Lines changed: 54 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -41,6 +41,9 @@
4141

4242
if TYPE_CHECKING:
4343
from dagster._core.definitions.assets.graph.base_asset_graph import EntityKey
44+
from dagster._core.definitions.partitions.definition.partitions_definition import (
45+
PartitionsDefinition,
46+
)
4447
from dagster._core.definitions.schedule_definition import ScheduleDefinition
4548
from dagster._core.definitions.sensor_definition import SensorDefinition
4649
from dagster._core.remote_representation.external import RemoteSchedule, RemoteSensor
@@ -372,6 +375,57 @@ def get_root_run_id(self) -> Optional[str]:
372375
def get_parent_run_id(self) -> Optional[str]:
373376
return self.tags.get(PARENT_RUN_ID_TAG)
374377

378+
def get_resolved_partitions_subset(
379+
self, partitions_def: Optional["PartitionsDefinition"]
380+
) -> Optional[PartitionsSubset]:
381+
"""Get the partitions subset targeted by a run based on its partition tags."""
382+
from dagster._core.definitions.partitions.definition import DynamicPartitionsDefinition
383+
from dagster._core.definitions.partitions.partition_key_range import PartitionKeyRange
384+
from dagster._core.errors import DagsterInvariantViolationError
385+
from dagster._core.storage.tags import (
386+
ASSET_PARTITION_RANGE_END_TAG,
387+
ASSET_PARTITION_RANGE_START_TAG,
388+
PARTITION_NAME_TAG,
389+
)
390+
391+
# some runs store this information directly
392+
if self.partitions_subset is not None:
393+
return self.partitions_subset
394+
395+
# otherwise, fetch information from the tags
396+
partition_tag = self.tags.get(PARTITION_NAME_TAG)
397+
partition_range_start = self.tags.get(ASSET_PARTITION_RANGE_START_TAG)
398+
partition_range_end = self.tags.get(ASSET_PARTITION_RANGE_END_TAG)
399+
400+
if partition_range_start or partition_range_end:
401+
if not partition_range_start or not partition_range_end:
402+
raise DagsterInvariantViolationError(
403+
f"Cannot have {ASSET_PARTITION_RANGE_START_TAG} or"
404+
f" {ASSET_PARTITION_RANGE_END_TAG} set without the other"
405+
)
406+
407+
if (
408+
isinstance(partitions_def, DynamicPartitionsDefinition)
409+
and partitions_def.name is None
410+
):
411+
raise DagsterInvariantViolationError(
412+
"Creating a run targeting a partition range is not supported for assets "
413+
"partitioned with function-based dynamic partitions"
414+
)
415+
416+
if partitions_def is not None:
417+
return partitions_def.subset_with_partition_keys(
418+
partitions_def.get_partition_keys_in_range(
419+
PartitionKeyRange(partition_range_start, partition_range_end),
420+
)
421+
).to_serializable_subset()
422+
elif partition_tag and partitions_def is not None:
423+
return partitions_def.subset_with_partition_keys(
424+
[partition_tag]
425+
).to_serializable_subset()
426+
427+
return None
428+
375429
@cached_property
376430
def dagster_execution_info(self) -> Mapping[str, str]:
377431
"""Key-value pairs encoding metadata about the current Dagster run, typically attached to external execution resources.

python_modules/dagster/dagster/_core/storage/event_log/base.py

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
import os
22
from abc import ABC, abstractmethod
33
from collections.abc import Iterable, Mapping, Sequence, Set
4+
from types import EllipsisType
45
from typing import TYPE_CHECKING, Annotated, NamedTuple, Optional, Union
56

67
from dagster_shared.record import ImportFrom, record
@@ -633,13 +634,16 @@ def get_asset_check_execution_history(
633634
limit: int,
634635
cursor: Optional[int] = None,
635636
status: Optional[Set[AssetCheckExecutionRecordStatus]] = None,
637+
partition: Optional[str] = None,
636638
) -> Sequence[AssetCheckExecutionRecord]:
637639
"""Get executions for one asset check, sorted by recency."""
638640
pass
639641

640642
@abstractmethod
641643
def get_latest_asset_check_execution_by_key(
642-
self, check_keys: Sequence[AssetCheckKey]
644+
self,
645+
check_keys: Sequence[AssetCheckKey],
646+
partition: Union[str, None, EllipsisType] = ...,
643647
) -> Mapping[AssetCheckKey, AssetCheckExecutionRecord]:
644648
"""Get the latest executions for a list of asset checks."""
645649
pass

0 commit comments

Comments
 (0)