adjust deals job

This commit is contained in:
2025-09-23 19:22:49 +02:00
parent 7a600f6264
commit 883ecf86be
2 changed files with 3 additions and 5 deletions

View File

@@ -176,9 +176,7 @@ def works(context: dg.AssetExecutionContext) -> Iterator[dg.Output[pl.DataFrame]
},
output_required=False,
dagster_type=patito_model_to_dagster_type(Deal),
automation_condition=dg.AutomationCondition.on_missing().ignore(
dg.AssetSelection.assets(cleaned_deals.key)
),
automation_condition=dg.AutomationCondition.eager(),
)
def new_deals(
context: dg.AssetExecutionContext, partitions: dict[str, pl.LazyFrame | None]

View File

@@ -12,8 +12,8 @@ define_asset_job = partial(dg.define_asset_job, **kwargs)
deals_job = dg.define_asset_job(
"deals_job",
selection=[assets.deals.key],
partitions_def=assets.deals.partitions_def,
selection=dg.AssetSelection.assets(assets.new_deals.key).upstream(),
partitions_def=assets.new_deals.partitions_def,
)