Skip to content

Commit 0cfb191

Browse files
committed
[pac] Add partitions_subset to AssetCheckEvaluationPlanned, partition to AssetCheckEvaluation
1 parent 9e9b67f commit 0cfb191

12 files changed

Lines changed: 640 additions & 109 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
@@ -10,6 +10,7 @@
1010
)
1111
from dagster._core.definitions.events import AssetKey, MetadataValue, RawMetadataValue
1212
from dagster._core.definitions.metadata import normalize_metadata
13+
from dagster._core.definitions.partitions.subset import PartitionsSubset
1314
from dagster._serdes import whitelist_for_serdes
1415

1516

@@ -20,6 +21,7 @@ class AssetCheckEvaluationPlanned:
2021

2122
asset_key: AssetKey
2223
check_name: str
24+
partitions_subset: Optional[PartitionsSubset] = None
2325

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

6569
asset_key: AssetKey
@@ -70,6 +74,7 @@ class AssetCheckEvaluation(IHaveNew):
7074
severity: AssetCheckSeverity
7175
description: Optional[str]
7276
blocking: Optional[bool]
77+
partition: Optional[str]
7378

7479
def __new__(
7580
cls,
@@ -81,6 +86,7 @@ def __new__(
8186
severity: AssetCheckSeverity = AssetCheckSeverity.ERROR,
8287
description: Optional[str] = None,
8388
blocking: Optional[bool] = None,
89+
partition: Optional[str] = None,
8490
):
8591
return super().__new__(
8692
cls,
@@ -94,6 +100,7 @@ def __new__(
94100
severity=severity,
95101
description=description,
96102
blocking=blocking,
103+
partition=partition,
97104
)
98105

99106
@property

python_modules/dagster/dagster/_core/definitions/assets/graph/remote_asset_graph.py

Lines changed: 6 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -744,12 +744,13 @@ def get_repository_handle(self, key: EntityKey) -> RepositoryHandle:
744744
def get_repo_scoped_node(
745745
self, key: EntityKey, repository_selector: "RepositorySelector"
746746
) -> Optional[Union[RemoteRepositoryAssetNode, RemoteAssetCheckNode]]:
747-
if isinstance(key, AssetKey):
748-
if not self.has(key):
749-
return None
750-
return self.get(key).resolve_to_repo_scoped_node(repository_selector)
747+
if not self.has(key):
748+
return None
749+
node = self.get(key)
750+
if isinstance(node, RemoteWorkspaceAssetNode):
751+
return node.resolve_to_repo_scoped_node(repository_selector)
751752
else:
752-
raise Exception("Key must be an asset key for get_repo_scoped_node")
753+
return node # type: ignore
753754

754755
def split_entity_keys_by_repository(
755756
self, keys: AbstractSet[EntityKey]

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: 72 additions & 89 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@
22
import os
33
import warnings
44
from collections.abc import Mapping, Sequence, Set
5-
from typing import TYPE_CHECKING, Any, Optional, cast
5+
from typing import TYPE_CHECKING, Any, Optional, cast, overload
66

77
import dagster._check as check
88
from dagster._core.definitions.asset_checks.asset_check_evaluation import (
@@ -50,10 +50,13 @@
5050
if TYPE_CHECKING:
5151
from dagster._core.definitions.asset_checks.asset_check_spec import AssetCheckKey
5252
from dagster._core.definitions.assets.graph.base_asset_graph import (
53+
AssetCheckNode,
5354
BaseAssetGraph,
5455
BaseAssetNode,
56+
BaseEntityNode,
5557
)
5658
from dagster._core.definitions.job_definition import JobDefinition
59+
from dagster._core.definitions.partitions.definition import PartitionsDefinition
5760
from dagster._core.definitions.repository_definition.repository_definition import (
5861
RepositoryLoadData,
5962
)
@@ -66,10 +69,7 @@
6669
from dagster._core.remote_representation.code_location import CodeLocation
6770
from dagster._core.remote_representation.external import RemoteJob
6871
from dagster._core.snap import ExecutionPlanSnapshot, JobSnap
69-
from dagster._core.snap.execution_plan_snapshot import (
70-
ExecutionStepOutputSnap,
71-
ExecutionStepSnap,
72-
)
72+
from dagster._core.snap.execution_plan_snapshot import ExecutionStepSnap
7373
from dagster._core.workspace.context import BaseWorkspaceRequestContext
7474

7575

@@ -108,6 +108,7 @@ def create_run(
108108
"""Create a run with the given parameters."""
109109
from dagster._core.definitions.asset_key import AssetCheckKey
110110
from dagster._core.definitions.assets.graph.remote_asset_graph import RemoteAssetGraph
111+
from dagster._core.definitions.partitions.context import partition_loading_context
111112
from dagster._core.remote_origin import RemoteJobOrigin
112113
from dagster._core.snap import ExecutionPlanSnapshot, JobSnap
113114
from dagster._utils.tags import normalize_tags
@@ -256,7 +257,8 @@ def create_run(
256257
dagster_run = self._instance.run_storage.add_run(dagster_run)
257258

258259
if execution_plan_snapshot and not assets_are_externally_managed(dagster_run):
259-
self._log_asset_planned_events(dagster_run, execution_plan_snapshot, asset_graph)
260+
with partition_loading_context(dynamic_partitions_store=self._instance):
261+
self._log_asset_planned_events(dagster_run, execution_plan_snapshot, asset_graph)
260262

261263
return dagster_run
262264

@@ -315,8 +317,8 @@ def construct_run_with_snapshots(
315317
adjusted_output = output
316318

317319
if asset_key:
318-
asset_node = self._get_repo_scoped_asset_node(
319-
asset_graph, asset_key, remote_job_origin
320+
asset_node = self._get_repo_scoped_entity_node(
321+
asset_key, asset_graph, remote_job_origin
320322
)
321323
if asset_node:
322324
partitions_definition = asset_node.partitions_def
@@ -767,12 +769,28 @@ def get_keys_to_reexecute(
767769
{key for key in to_reexecute if isinstance(key, AssetCheckKey)},
768770
)
769771

770-
def _get_repo_scoped_asset_node(
772+
@overload
773+
def _get_repo_scoped_entity_node(
771774
self,
775+
key: AssetKey,
776+
asset_graph: "BaseAssetGraph",
777+
remote_job_origin: Optional["RemoteJobOrigin"] = None,
778+
) -> Optional["BaseAssetNode"]: ...
779+
780+
@overload
781+
def _get_repo_scoped_entity_node(
782+
self,
783+
key: "AssetCheckKey",
772784
asset_graph: "BaseAssetGraph",
773-
asset_key: AssetKey,
774785
remote_job_origin: Optional["RemoteJobOrigin"] = None,
775-
) -> Optional["BaseAssetNode"]:
786+
) -> Optional["AssetCheckNode"]: ...
787+
788+
def _get_repo_scoped_entity_node(
789+
self,
790+
key: "EntityKey",
791+
asset_graph: "BaseAssetGraph",
792+
remote_job_origin: Optional["RemoteJobOrigin"] = None,
793+
) -> Optional["BaseEntityNode"]:
776794
from dagster._core.definitions.assets.graph.remote_asset_graph import (
777795
RemoteWorkspaceAssetGraph,
778796
)
@@ -783,16 +801,29 @@ def _get_repo_scoped_asset_node(
783801
# in all cases, return the BaseAssetNode for the supplied asset key if it exists.
784802
if isinstance(asset_graph, RemoteWorkspaceAssetGraph):
785803
return cast(
786-
"Optional[BaseAssetNode]",
804+
"Optional[BaseEntityNode]",
787805
asset_graph.get_repo_scoped_node(
788-
asset_key, check.not_none(remote_job_origin).repository_origin.get_selector()
806+
key, check.not_none(remote_job_origin).repository_origin.get_selector()
789807
),
790808
)
791809

792-
if not asset_graph.has(asset_key):
810+
if not asset_graph.has(key):
793811
return None
794812

795-
return asset_graph.get(asset_key)
813+
return asset_graph.get(key)
814+
815+
def _get_partitions_def(
816+
self,
817+
key: "EntityKey",
818+
asset_graph: "BaseAssetGraph",
819+
remote_job_origin: Optional["RemoteJobOrigin"],
820+
run: "DagsterRun",
821+
) -> Optional["PartitionsDefinition"]:
822+
# don't fetch the partitions def if the run is not partitioned
823+
if not run.is_partitioned:
824+
return None
825+
entity_node = self._get_repo_scoped_entity_node(key, asset_graph, remote_job_origin)
826+
return entity_node.partitions_def if entity_node else None
796827

797828
def _log_asset_planned_events(
798829
self,
@@ -819,7 +850,7 @@ def _log_asset_planned_events(
819850
if asset_key:
820851
events.extend(
821852
self.get_materialization_planned_events_for_asset(
822-
dagster_run, asset_key, job_name, step, output, asset_graph
853+
dagster_run, asset_key, job_name, step, asset_graph
823854
)
824855
)
825856

@@ -830,6 +861,13 @@ def _log_asset_planned_events(
830861
target_asset_key = asset_check_key.asset_key
831862
check_name = asset_check_key.name
832863

864+
partitions_def = self._get_partitions_def(
865+
asset_check_key, asset_graph, dagster_run.remote_job_origin, dagster_run
866+
)
867+
partitions_subset = dagster_run.get_resolved_partitions_subset(
868+
partitions_def
869+
)
870+
833871
event = DagsterEvent(
834872
event_type_value=DagsterEventType.ASSET_CHECK_EVALUATION_PLANNED.value,
835873
job_name=job_name,
@@ -840,6 +878,7 @@ def _log_asset_planned_events(
840878
event_specific_data=AssetCheckEvaluationPlanned(
841879
asset_key=target_asset_key,
842880
check_name=check_name,
881+
partitions_subset=partitions_subset,
843882
),
844883
step_key=step.key,
845884
)
@@ -878,94 +917,26 @@ def get_materialization_planned_events_for_asset(
878917
asset_key: AssetKey,
879918
job_name: str,
880919
step: "ExecutionStepSnap",
881-
output: "ExecutionStepOutputSnap",
882920
asset_graph: "BaseAssetGraph[BaseAssetNode]",
883921
) -> Sequence["DagsterEvent"]:
884922
"""Moved from DagsterInstance._log_materialization_planned_event_for_asset."""
885-
from dagster._core.definitions.partitions.context import partition_loading_context
886-
from dagster._core.definitions.partitions.definition import DynamicPartitionsDefinition
887923
from dagster._core.events import AssetMaterializationPlannedData, DagsterEvent
888924

889925
events = []
890926

891-
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-
)
896-
897-
if partition_tag and (partition_range_start or partition_range_end):
898-
raise DagsterInvariantViolationError(
899-
f"Cannot have {ASSET_PARTITION_RANGE_START_TAG} or"
900-
f" {ASSET_PARTITION_RANGE_END_TAG} set along with"
901-
f" {PARTITION_NAME_TAG}"
902-
)
903-
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-
)
912-
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-
)
918-
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"
926-
)
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"
927+
partitions_def = self._get_partitions_def(
928+
asset_key, asset_graph, dagster_run.remote_job_origin, dagster_run
946929
)
947930

948-
if not individual_partitions and not partitions_subset:
931+
partitions_subset = dagster_run.get_resolved_partitions_subset(partitions_def)
932+
if partitions_subset is None:
949933
materialization_planned = DagsterEvent.build_asset_materialization_planned_event(
950934
job_name,
951935
step.key,
952936
AssetMaterializationPlannedData(asset_key, partition=None, partitions_subset=None),
953937
)
954938
events.append(materialization_planned)
955-
elif individual_partitions:
956-
for individual_partition in individual_partitions:
957-
materialization_planned = DagsterEvent.build_asset_materialization_planned_event(
958-
job_name,
959-
step.key,
960-
AssetMaterializationPlannedData(
961-
asset_key,
962-
partition=individual_partition,
963-
partitions_subset=partitions_subset,
964-
),
965-
)
966-
events.append(materialization_planned)
967-
968-
else:
939+
elif self._instance.event_log_storage.supports_partition_subset_in_asset_materialization_planned_events:
969940
materialization_planned = DagsterEvent.build_asset_materialization_planned_event(
970941
job_name,
971942
step.key,
@@ -974,6 +945,18 @@ def get_materialization_planned_events_for_asset(
974945
),
975946
)
976947
events.append(materialization_planned)
948+
else:
949+
for partition_key in partitions_subset.get_partition_keys():
950+
materialization_planned = DagsterEvent.build_asset_materialization_planned_event(
951+
job_name,
952+
step.key,
953+
AssetMaterializationPlannedData(
954+
asset_key,
955+
partition=partition_key,
956+
partitions_subset=None,
957+
),
958+
)
959+
events.append(materialization_planned)
977960

978961
return events
979962

0 commit comments

Comments
 (0)