Skip to content
Open
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
25 changes: 22 additions & 3 deletions python_modules/dagster/dagster/_core/events/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -410,6 +410,26 @@ def log_resource_event(log_manager: DagsterLogManager, event: "DagsterEvent") ->
log_manager.log_dagster_event(level=log_level, msg=event.message or "", dagster_event=event)


def _resource_init_start_message(
resource_keys: AbstractSet[str],
resource_key_to_op_names: Mapping[str, AbstractSet[str]] | None,
) -> str:
if not resource_key_to_op_names:
return "Starting initialization of resources [{}].".format(", ".join(sorted(resource_keys)))
Comment on lines +417 to +418

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 The guard if not resource_key_to_op_names falls through to the old message format for both None (no job_def) and an empty dict {}. An empty dict is returned by get_resource_origins_mapping when every resource in the job is required only through the I/O layer or type system (e.g., a job whose ops have no explicit required_resource_keys). In that case job_def was provided and the enriched format should still be used — each resource would correctly display as "required by the I/O layer" — but the old flat-list fallback is emitted instead. The check should use an identity comparison with None.

Suggested change
if not resource_key_to_op_names:
return "Starting initialization of resources [{}].".format(", ".join(sorted(resource_keys)))
if resource_key_to_op_names is None:
return "Starting initialization of resources [{}].".format(", ".join(sorted(resource_keys)))


# Iterate resource_keys (scoped to this step) rather than the job-wide mapping, so a step
# only lists its own ops. Keys not in the mapping are required by an io/input manager or
# type loader rather than by an op directly.
lines = ["Starting initialization of resources:"]
for resource_key in sorted(resource_keys):
op_names = sorted(resource_key_to_op_names.get(resource_key, set()))
if op_names:
lines.append(f"{resource_key} - required by op instances [{', '.join(op_names)}]")
else:
lines.append(f"{resource_key} - required by the I/O layer")
return "\n".join(lines)


class DagsterEventSerializer(NamedTupleSerializer["DagsterEvent"]):
def before_unpack(self, context, unpacked_dict: Any) -> dict[str, Any]:
event_type_value, event_specific_data = _handle_back_compat(
Expand Down Expand Up @@ -1337,15 +1357,14 @@ def resource_init_start(
execution_plan: "ExecutionPlan",
log_manager: DagsterLogManager,
resource_keys: AbstractSet[str],
resource_key_to_op_names: Mapping[str, AbstractSet[str]] | None = None,
) -> "DagsterEvent":
return DagsterEvent.from_resource(
DagsterEventType.RESOURCE_INIT_STARTED,
job_name=job_name,
execution_plan=execution_plan,
log_manager=log_manager,
message="Starting initialization of resources [{}].".format(
", ".join(sorted(resource_keys))
),
message=_resource_init_start_message(resource_keys, resource_key_to_op_names),
event_specific_data=EngineEventData(metadata={}, marker_start="resources"),
)

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -238,6 +238,7 @@ def execution_context_event_generator(
instance=instance,
emit_persistent_events=True,
event_loop=event_loop,
job_def=job_def,
)
yield from resources_manager.generate_setup_events()
scoped_resources_builder = check.inst(
Expand Down
43 changes: 42 additions & 1 deletion python_modules/dagster/dagster/_core/execution/resources_init.py
Original file line number Diff line number Diff line change
@@ -1,13 +1,14 @@
import inspect
from asyncio import AbstractEventLoop
from collections import deque
from collections import defaultdict, deque
from collections.abc import Generator, Mapping
from contextlib import ContextDecorator
from typing import AbstractSet, Any, Callable, cast # noqa: UP035

from dagster_shared.utils.timing import format_duration

import dagster._check as check
from dagster._core.definitions.dependency import OpNode
from dagster._core.definitions.job_definition import JobDefinition
from dagster._core.definitions.resource_definition import (
ResourceDefinition,
Expand Down Expand Up @@ -49,6 +50,7 @@ def resource_initialization_manager(
instance: DagsterInstance | None,
emit_persistent_events: bool | None,
event_loop: AbstractEventLoop | None,
job_def: JobDefinition | None = None,
):
generator = resource_initialization_event_generator(
resource_defs=resource_defs,
Expand All @@ -60,6 +62,7 @@ def resource_initialization_manager(
instance=instance,
emit_persistent_events=emit_persistent_events,
event_loop=event_loop,
job_def=job_def,
)
return EventGenerationManager(generator, ScopedResourcesBuilder)

Expand Down Expand Up @@ -116,6 +119,39 @@ def _get_deps_helper(resource_key):
return reqd_resources


def get_resource_origins_mapping(job_def: JobDefinition) -> Mapping[str, AbstractSet[str]]:
"""Maps each resource key to the op instances that require it, keyed by full node handle so
that aliased and nested ops are unambiguous.

A resource required by another resource (e.g. spark, pulled in by a pyspark resource) is
attributed to the ops that required the outer resource. Resources no op requires directly,
such as the default io_manager, are absent from the mapping.
"""
origins: dict[str, set[str]] = defaultdict(set)

for handle in job_def.graph.iterate_node_handles():
node = job_def.get_node(handle)
if not isinstance(node, OpNode):
continue

op_name = str(handle)

for resource_key in node.definition.required_resource_keys:
origins[resource_key].add(op_name)

for hook_def in job_def.get_all_hooks_for_handle(handle):
for resource_key in hook_def.required_resource_keys:
origins[resource_key].add(op_name)

resource_dependencies = resolve_resource_dependencies(job_def.resource_defs)
for resource_key, op_names in list(origins.items()):
for dep_key in get_dependencies(resource_key, resource_dependencies):
if dep_key != resource_key:
origins[dep_key].update(op_names)
Comment on lines +146 to +150

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 KeyError risk on unregistered resource keys

get_dependencies indexes into resource_dependencies with every key currently in origins. If an op (or hook) declares a required_resource_keys entry that is not present in job_def.resource_defs, resource_deps[resource_key] will raise a KeyError before the enriched log message is produced. In a valid, fully-validated job this can't happen, but if get_resource_origins_mapping is ever called before full job validation (e.g., during inspection or tooling), an uncaught KeyError would surface here as an obscure internal error rather than a clear validation message. A resource_key in resource_dependencies guard, or a resource_dependencies.get(resource_key, set()) call inside get_dependencies, would make the failure mode explicit.


return dict(origins)


def _core_resource_initialization_event_generator(
resource_defs: Mapping[str, ResourceDefinition],
resource_configs: Mapping[str, ResourceConfig],
Expand All @@ -127,6 +163,7 @@ def _core_resource_initialization_event_generator(
instance: DagsterInstance | None,
emit_persistent_events: bool | None,
event_loop,
job_def: JobDefinition | None = None,
):
job_name = "" # Must be initialized to a string to satisfy typechecker
contains_generator = False
Expand All @@ -147,6 +184,7 @@ def _core_resource_initialization_event_generator(
cast("ExecutionPlan", execution_plan),
resource_log_manager,
resource_keys_to_init,
get_resource_origins_mapping(job_def) if job_def else None,
)

resource_dependencies = resolve_resource_dependencies(resource_defs)
Expand Down Expand Up @@ -225,6 +263,7 @@ def resource_initialization_event_generator(
instance: DagsterInstance | None,
emit_persistent_events: bool | None,
event_loop: AbstractEventLoop | None,
job_def: JobDefinition | None = None,
):
check.inst_param(log_manager, "log_manager", DagsterLogManager)
resource_keys_to_init = check.opt_set_param(
Expand All @@ -233,6 +272,7 @@ def resource_initialization_event_generator(
check.opt_inst_param(execution_plan, "execution_plan", ExecutionPlan)
check.opt_inst_param(dagster_run, "dagster_run", DagsterRun)
check.opt_inst_param(instance, "instance", DagsterInstance)
check.opt_inst_param(job_def, "job_def", JobDefinition)

if execution_plan and execution_plan.step_handle_for_single_step_plans():
step = execution_plan.get_step(
Expand Down Expand Up @@ -260,6 +300,7 @@ def resource_initialization_event_generator(
instance=instance,
emit_persistent_events=emit_persistent_events,
event_loop=event_loop,
job_def=job_def,
)
except GeneratorExit:
# Shouldn't happen, but avoid runtime-exception in case this generator gets GC-ed
Expand Down
168 changes: 168 additions & 0 deletions python_modules/dagster/dagster_tests/core_tests/test_job_init.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,9 +2,13 @@
import pytest
from dagster import DagsterInstance
from dagster._core.definitions.job_base import InMemoryJob
from dagster._core.events import DagsterEvent, DagsterEventType
from dagster._core.execution import resources_init
from dagster._core.execution.api import create_execution_plan
from dagster._core.execution.context_creation_job import PlanExecutionContextManager
from dagster._core.execution.resources_init import (
get_required_resource_keys_to_init,
get_resource_origins_mapping,
resource_initialization_event_generator,
resource_initialization_manager,
single_resource_event_generator,
Expand Down Expand Up @@ -117,3 +121,167 @@ def test_clean_event_generator_exit():
@dg.op
def fake_op(_):
pass


def test_get_resource_origins_mapping():
@dg.resource
def spark(_):
return "spark"

@dg.resource(required_resource_keys={"spark"})
def pyspark(_):
return "pyspark"

@dg.resource
def required_by_nobody(_):
return "nobody"

@dg.op(required_resource_keys={"pyspark"})
def uses_pyspark(_):
pass

@dg.op(required_resource_keys={"spark"})
def uses_spark_directly(_):
pass

@dg.graph
def nested():
uses_spark_directly()

@dg.job(
resource_defs={
"spark": spark,
"pyspark": pyspark,
"required_by_nobody": required_by_nobody,
}
)
def origins_job():
uses_pyspark()
nested()

mapping = get_resource_origins_mapping(origins_job)

assert mapping["pyspark"] == {"uses_pyspark"}

# spark is required transitively (via pyspark) and directly by a nested op, reported by handle
assert mapping["spark"] == {"uses_pyspark", "nested.uses_spark_directly"}

# resources no op requires are absent from the mapping
assert "required_by_nobody" not in mapping
assert "io_manager" not in mapping


def gen_two_op_resource_job():
"""A job where each op requires a different resource."""

@dg.resource
def resource_a(_):
return "a"

@dg.resource
def resource_b(_):
return "b"

@dg.op(required_resource_keys={"resource_a"})
def op_one(_):
pass

@dg.op(required_resource_keys={"resource_b"})
def op_two(_):
pass

@dg.job(resource_defs={"resource_a": resource_a, "resource_b": resource_b})
def two_op_resource_job():
op_one()
op_two()

return two_op_resource_job


def _resource_init_started_message(job_def, *, pass_job_def: bool) -> str:
"""Drives the real resource-init generator and returns the RESOURCE_INIT_STARTED message."""
instance = DagsterInstance.ephemeral()
execution_plan = create_execution_plan(job_def)
run = instance.create_run_for_job(job_def=job_def, execution_plan=execution_plan)
log_manager = DagsterLogManager.create(loggers=[], dagster_run=run)
resolved_run_config = ResolvedRunConfig.build(job_def)

events = list(
resource_initialization_event_generator(
resource_defs=job_def.resource_defs,
resource_configs=resolved_run_config.resources,
log_manager=log_manager,
execution_plan=execution_plan,
dagster_run=run,
resource_keys_to_init=get_required_resource_keys_to_init(execution_plan, job_def),
instance=instance,
emit_persistent_events=True,
event_loop=None,
job_def=job_def if pass_job_def else None,
)
)

started = [
event
for event in events
if isinstance(event, DagsterEvent)
and event.event_type_value == DagsterEventType.RESOURCE_INIT_STARTED.value
]
assert len(started) == 1, f"expected exactly one RESOURCE_INIT_STARTED, got {len(started)}"
return started[0].message


def test_resource_init_message_reports_op_origins():
"""With job_def available, the message says which ops require each resource (issue #2307)."""
message = _resource_init_started_message(gen_two_op_resource_job(), pass_job_def=True)

assert message == (
"Starting initialization of resources:\n"
"io_manager - required by the I/O layer\n"
"resource_a - required by op instances [op_one]\n"
"resource_b - required by op instances [op_two]"
)


def test_resource_init_message_falls_back_without_job_def():
"""Without job_def (e.g. build_resources), the message is byte-for-byte the original."""
message = _resource_init_started_message(gen_two_op_resource_job(), pass_job_def=False)

assert message == "Starting initialization of resources [io_manager, resource_a, resource_b]."


def test_job_def_reaches_innermost_resource_generator(monkeypatch):
"""job_def must be forwarded from resource_initialization_manager down to the innermost
generator. A param that is accepted but not forwarded still typechecks, so assert it arrives.
"""
job_def = gen_two_op_resource_job()
instance = DagsterInstance.ephemeral()
execution_plan = create_execution_plan(job_def)
run = instance.create_run_for_job(job_def=job_def, execution_plan=execution_plan)
log_manager = DagsterLogManager.create(loggers=[], dagster_run=run)
resolved_run_config = ResolvedRunConfig.build(job_def)

captured = {}
real_core_generator = resources_init._core_resource_initialization_event_generator # noqa: SLF001

def spy(**kwargs):
captured["job_def"] = kwargs.get("job_def", "<<not forwarded>>")
return real_core_generator(**kwargs)

monkeypatch.setattr(resources_init, "_core_resource_initialization_event_generator", spy)

manager = resource_initialization_manager(
resource_defs=job_def.resource_defs,
resource_configs=resolved_run_config.resources,
log_manager=log_manager,
execution_plan=execution_plan,
dagster_run=run,
resource_keys_to_init=get_required_resource_keys_to_init(execution_plan, job_def),
instance=instance,
emit_persistent_events=True,
event_loop=None,
job_def=job_def,
)
list(manager.generate_setup_events())

assert captured["job_def"] is job_def
2 changes: 2 additions & 0 deletions python_modules/libraries/dagstermill/dagstermill/manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -86,6 +86,7 @@ def _setup_resources(
instance: DagsterInstance | None,
emit_persistent_events: bool | None,
event_loop: AbstractEventLoop | None,
job_def: JobDefinition | None = None,
):
"""Drop-in replacement for
`dagster._core.execution.resources_init.resource_initialization_manager`. It uses a
Expand All @@ -101,6 +102,7 @@ def _setup_resources(
instance=instance,
emit_persistent_events=emit_persistent_events,
event_loop=event_loop,
job_def=job_def,
)
self.resource_manager = DagstermillResourceEventGenerationManager(
generator, ScopedResourcesBuilder
Expand Down