使用资产检查测试 Dagster 数据资产

资产检查是一种测试,用于验证数据资产的特定属性,从而对数据执行质量检查。例如,可以创建检查来确保某一列不包含空值,验证表格资产是否符合指定的架构,或者判断资产中的数据是否需要刷新。

每个资产检查应只测试一个资产属性,这样测试才容易理解、能够复用,也便于长期追踪。

Dagster+ 额度消耗:资产检查不计入 Dagster+ 的额度使用量。

开始使用

使用资产检查通常包括以下步骤:

  1. 定义资产检查:通常使用 @asset_check 或 @multi_asset_check 装饰器定义检查。检查既可以在资产内部执行,也可以与资产分开执行。
  2. 将检查传给 Definitions 对象:必须将资产检查添加到 Definitions,Dagster 才能识别它们。
  3. 选择执行方式:默认情况下,所有以某个资产为目标的作业也会运行关联的检查;也可以通过 Dagster 界面运行资产检查。
  4. 在界面中查看结果:检查结果会显示在界面中,可以通过元数据与严重性级别自定义显示内容。
  5. 为失败结果设置告警:使用 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 文档团队。本文为原文的中文译文;代码保留原文内容。

© 版权声明
THE END
喜欢就支持一下吧
点赞0 分享
评论 抢沙发

请登录后发表评论

    暂无评论内容