资产检查是一种测试,用于验证数据资产的特定属性,从而对数据执行质量检查。例如,可以创建检查来确保某一列不包含空值,验证表格资产是否符合指定的架构,或者判断资产中的数据是否需要刷新。
每个资产检查应只测试一个资产属性,这样测试才容易理解、能够复用,也便于长期追踪。
Dagster+ 额度消耗:资产检查不计入 Dagster+ 的额度使用量。
开始使用
使用资产检查通常包括以下步骤:
- 定义资产检查:通常使用
@asset_check或@multi_asset_check装饰器定义检查。检查既可以在资产内部执行,也可以与资产分开执行。 - 将检查传给 Definitions 对象:必须将资产检查添加到
Definitions,Dagster 才能识别它们。 - 选择执行方式:默认情况下,所有以某个资产为目标的作业也会运行关联的检查;也可以通过 Dagster 界面运行资产检查。
- 在界面中查看结果:检查结果会显示在界面中,可以通过元数据与严重性级别自定义显示内容。
- 为失败结果设置告警:使用 Dagster+ 时,可以选择在资产检查失败时触发告警。
定义单个资产检查
提示:Dagster 的 dbt 集成可以将现有的 dbt 测试建模为资产检查。更多信息请参阅 dagster-dbt 文档。
使用 @asset_check 装饰器定义一个资产检查。
下面的示例针对一个资产定义检查:如果资产的 order_id 列包含空值,检查就会失败。检查会在资产物化之后运行。
文件:src/<project_name>/defs/assets.py
import pandas as pd
import dagster as dg
@dg.asset
def orders():
orders_df = pd.DataFrame({"order_id": [1, 2], "item_id": [432, 878]})
orders_df.to_csv("orders.csv")
@dg.asset_check(asset=orders)
def orders_id_has_no_nulls():
orders_df = pd.read_csv("orders.csv")
num_null_order_ids = orders_df["order_id"].isna().sum()
# Return the result of the check
return dg.AssetCheckResult(
# Define passing criteria
passed=bool(num_null_order_ids == 0),
)
定义多个资产检查
在大多数情况下,检查资产的数据质量需要执行不止一项检查。
下面的示例使用 @multi_asset_check 装饰器定义两个检查:第一个在 order_id 列包含空值时失败;第二个在 item_id 列包含空值时失败。在这个示例中,两个检查会在资产物化后通过同一个操作运行。
文件:src/<project_name>/defs/assets.py
from collections.abc import Iterable
import pandas as pd
import dagster as dg
@dg.asset
def orders():
orders_df = pd.DataFrame({"order_id": [1, 2], "item_id": [432, 878]})
orders_df.to_csv("orders.csv")
@dg.multi_asset_check(
# Map checks to targeted assets
specs=[
dg.AssetCheckSpec(name="orders_id_has_no_nulls", asset="orders"),
dg.AssetCheckSpec(name="items_id_has_no_nulls", asset="orders"),
]
)
def orders_check() -> Iterable[dg.AssetCheckResult]:
orders_df = pd.read_csv("orders.csv")
# Check for null order_id column values
num_null_order_ids = orders_df["order_id"].isna().sum()
yield dg.AssetCheckResult(
check_name="orders_id_has_no_nulls",
passed=bool(num_null_order_ids == 0),
asset_key="orders",
)
# Check for null item_id column values
num_null_item_ids = orders_df["item_id"].isna().sum()
yield dg.AssetCheckResult(
check_name="items_id_has_no_nulls",
passed=bool(num_null_item_ids == 0),
asset_key="orders",
)
以编程方式生成资产检查
也可以使用工厂模式定义多个检查。下面定义的两个检查与前一个示例相同,不过这次使用了工厂模式和 @multi_asset_check 装饰器。
文件:src/<project_name>/defs/assets.py
from collections.abc import Iterable, Mapping, Sequence
import pandas as pd
import dagster as dg
@dg.asset
def orders():
orders_df = pd.DataFrame({"order_id": [1, 2], "item_id": [432, 878]})
orders_df.to_csv("orders.csv")
def make_orders_checks(
check_blobs: Sequence[Mapping[str, str]],
) -> dg.AssetChecksDefinition:
@dg.multi_asset_check(
specs=[
dg.AssetCheckSpec(name=check_blob["name"], asset=check_blob["asset"])
for check_blob in check_blobs
]
)
def orders_check() -> Iterable[dg.AssetCheckResult]:
orders_df = pd.read_csv("orders.csv")
for check_blob in check_blobs:
num_null_order_ids = orders_df[check_blob["column"]].isna().sum()
yield dg.AssetCheckResult(
check_name=check_blob["name"],
passed=bool(num_null_order_ids == 0),
asset_key=check_blob["asset"],
)
return orders_check
check_blobs = [
{
"name": "orders_id_has_no_nulls",
"asset": "orders",
"column": "order_id",
},
{
"name": "items_id_has_no_nulls",
"asset": "orders",
"column": "item_id",
},
]
@dg.definitions
def asset_checks():
return dg.Definitions(
asset_checks=[make_orders_checks(check_blobs)],
)
阻止下游资产物化
默认情况下,即使父资产的检查在一次运行中失败,运行仍会继续,下游资产也会物化。要阻止这种行为,请将 @asset_check 装饰器的 blocking 参数设为 True。
在下面的示例中,如果 orders_id_has_no_nulls 检查失败,下游的 augmented_orders 资产就不会物化。
文件:src/<project_name>/defs/assets.py
import pandas as pd
import dagster as dg
@dg.asset
def orders():
orders_df = pd.DataFrame({"order_id": [1, 2], "item_id": [432, 878]})
orders_df.to_csv("orders.csv")
# Check that targets `orders`; block materialization of `augmented_orders` on failure
@dg.asset_check(asset=orders, blocking=True)
def orders_id_has_no_nulls():
orders_df = pd.read_csv("orders.csv")
num_null_order_ids = orders_df["order_id"].isna().sum()
return dg.AssetCheckResult(
passed=bool(num_null_order_ids == 0),
)
# Asset downstream of `orders`
@dg.asset(deps=[orders])
def augmented_orders():
orders_df = pd.read_csv("orders.csv")
augmented_orders_df = orders_df.assign(description=["item_432", "item_878"])
augmented_orders_df.to_csv("augmented_orders.csv")
分区资产检查
预览功能:这项功能目前处于预览阶段,仍在积极开发,尚未被视为适合生产使用。你可能会遇到功能缺口,API 也可能发生变化。详情请参阅 API 生命周期阶段文档。
资产检查可以使用与被检查资产一致的分区。向 @asset_check 装饰器或 AssetCheckSpec 提供 partitions_def 后,每次执行检查都会针对资产的一个分区。这样就能逐个分区验证数据质量,并在界面中查看每个分区的检查状态。
检查的 partitions_def 必须与关联资产的 partitions_def 相同。
文件:src/<project_name>/defs/assets.py
import pandas as pd
import dagster as dg
partitions_def = dg.DailyPartitionsDefinition(start_date="2024-01-01")
@dg.asset(partitions_def=partitions_def)
def orders(context: dg.AssetExecutionContext):
orders_df = pd.DataFrame({"order_id": [1, 2], "item_id": [432, 878]})
orders_df.to_csv(f"orders_{context.partition_key}.csv")
@dg.asset_check(asset=orders, partitions_def=partitions_def)
def orders_id_has_no_nulls(context: dg.AssetCheckExecutionContext):
orders_df = pd.read_csv(f"orders_{context.partition_key}.csv")
num_null_order_ids = orders_df["order_id"].isna().sum()
return dg.AssetCheckResult(passed=bool(num_null_order_ids == 0))
# Or, define the check inline using check_specs on the asset:
@dg.asset(
partitions_def=partitions_def,
check_specs=[
dg.AssetCheckSpec(
name="orders_id_has_no_nulls",
asset="inline_orders",
partitions_def=partitions_def,
)
],
)
def inline_orders(context: dg.AssetExecutionContext):
orders_df = pd.DataFrame({"order_id": [1, 2], "item_id": [432, 878]})
orders_df.to_csv(f"orders_{context.partition_key}.csv")
yield dg.Output(value=None)
num_null_order_ids = orders_df["order_id"].isna().sum()
yield dg.AssetCheckResult(passed=bool(num_null_order_ids == 0))
调度和监控资产检查
在某些情况下,将资产检查与负责物化资产的作业分开运行会很有帮助。例如,每天运行一次全部数据质量检查,并在失败时发送告警。可以通过调度和传感器实现这一点。
下面的示例定义了两个作业:一个用于资产,另一个用于资产检查。两个调度分别独立地物化资产、执行资产检查。同时定义了一个传感器,在资产检查作业失败时发送电子邮件告警。
文件:src/<project_name>/defs/assets.py
import os
import pandas as pd
import dagster as dg
@dg.asset
def orders():
orders_df = pd.DataFrame({"order_id": [1, 2], "item_id": [432, 878]})
orders_df.to_csv("orders.csv")
@dg.asset_check(asset=orders)
def orders_id_has_no_nulls():
orders_df = pd.read_csv("orders.csv")
num_null_order_ids = orders_df["order_id"].isna().sum()
return dg.AssetCheckResult(
passed=bool(num_null_order_ids == 0),
)
# Only include the `orders` asset
asset_job = dg.define_asset_job(
"asset_job",
selection=dg.AssetSelection.assets(orders).without_checks(),
)
# Only include the `orders_id_has_no_nulls` check
check_job = dg.define_asset_job(
"check_job", selection=dg.AssetSelection.checks_for_assets(orders)
)
# Job schedules
asset_schedule = dg.ScheduleDefinition(job=asset_job, cron_schedule="0 0 * * *")
check_schedule = dg.ScheduleDefinition(job=check_job, cron_schedule="0 6 * * *")
# Send email on failure
check_sensor = dg.make_email_on_run_failure_sensor(
email_from="no-reply@example.com",
email_password=os.getenv("ALERT_EMAIL_PASSWORD"), # ty: ignore[invalid-argument-type]
email_to=["xxx@example.com"],
monitored_jobs=[check_job],
)
原文:Testing assets with asset checks。作者/维护者:Dagster 文档团队。本文为原文的中文译文;代码保留原文内容。











暂无评论内容