Dagster 分区回填策略:逐分区、分批与单次运行

本文介绍三种回填分区资产的策略。为了初始化或重新处理数据而需要物化多个分区时,可以选择 Dagster 默认的每分区一次运行、分批运行,或者通过 BackfillPolicy 单次运行处理所有分区。每种策略在运行开销、故障隔离和资源利用方面都有不同取舍。

因素每分区一次分批单次运行
运行开销高,N 次运行中等,N/批量大小次低,1 次运行
故障隔离最好中等无
重试成本1 个分区1 个批次全部分区
可观测性逐分区逐批次仅整体

问题:回填 100 天历史数据

假设需要回填 100 天的历史事件,每一天都要处理并存储。不优化的话,可能需要启动 100 次独立运行,每次都有启动开销。但把所有数据放在一次运行里,又意味着一次失败就要重新处理全部 100 天。

应如何将分区分批,才能平衡开销、故障隔离与性能?

方案最适合的场景
每分区一次运行,默认不稳定数据源、API 限流、精细重试、逐分区观测
分批运行减少开销但保留故障隔离、单分区处理时间短、初始回填
单次运行Spark/Snowflake/Databricks、范围查询、尽量减少 Dagster+ 额度消耗

方案一:每分区一次运行,默认方式

Dagster 默认每个分区启动一次运行,提供最强的可观测性和故障隔离。一个分区失败,其他分区仍独立继续,也可以单独重试失败分区。100 个分区会产生 100 次运行,每次都有启动开销。

如果数据源不稳定,例如存在 API 限流或临时失败,或者需要逐分区精细重试、逐分区观测非常重要,这种方式最适合。

文件:src/project_mini/defs/partition_backfill_strategies/multi_run_backfill.py

import dagster as dg

daily_partitions = dg.DailyPartitionsDefinition(start_date="2024-01-01")


@dg.asset(partitions_def=daily_partitions)
def daily_events(context: dg.AssetExecutionContext):
    """Process events for a single day. Each partition runs separately."""
    partition_date = context.partition_key
    context.log.info(f"Processing events for {partition_date}")

    # Process data for this single partition
    events = fetch_events_for_date(partition_date)
    processed = transform_events(events)

    context.log.info(f"Processed {len(processed)} events for {partition_date}")
    return processed


def fetch_events_for_date(date: str) -> list:
    # Simulate fetching events for a specific date
    return [{"date": date, "event_id": i} for i in range(100)]


def transform_events(events: list) -> list:
    # Simulate transformation
    return [{"processed": True, **e} for e in events]

方案二:分批运行

BackfillPolicy.multi_run 让 Dagster 将分区分组。例如,设置 max_partitions_per_run=10 后,100 个分区变为 10 次运行,每次处理 10 个分区。运行次数降低 90%,同时保留中等程度的故障隔离:一个分区失败时,只需重试所在的 10 分区批次。

希望降低开销但保留一定故障隔离、单分区处理仅需几秒到几分钟,或者首次回填大量分区时,这种方式很合适。

文件:src/project_mini/defs/partition_backfill_strategies/batched_backfill.py

import dagster as dg

daily_partitions = dg.DailyPartitionsDefinition(start_date="2024-01-01")


@dg.asset(
    partitions_def=daily_partitions,
    backfill_policy=dg.BackfillPolicy.multi_run(max_partitions_per_run=10),
)
def daily_events(context: dg.AssetExecutionContext):
    """Process events for a batch of days in each run."""
    partition_keys = context.partition_keys
    context.log.info(f"Processing {len(partition_keys)} partitions in this run")

    # Process data for all partitions in this batch
    all_events = []
    for partition_key in partition_keys:
        events = fetch_events_for_date(partition_key)
        all_events.extend(events)

    processed = transform_events(all_events)
    context.log.info(f"Processed {len(processed)} total events across {len(partition_keys)} days")
    return processed


def fetch_events_for_date(date: str) -> list:
    # Simulate fetching events for a specific date
    return [{"date": date, "event_id": i} for i in range(100)]


def transform_events(events: list) -> list:
    # Simulate transformation
    return [{"processed": True, **e} for e in events]

使用 BackfillPolicy.multi_run 时,需要考虑:

  • 开销降低:批量大小为 10 时,运行次数减少 90%。
  • 失败影响范围:一个分区失败,整个批次重试。
  • 内存使用:每次运行包含更多分区,可能需要更多内存。
  • 处理模型:串行处理时,批次越大,耗时越长。

建议从以下批量大小开始:

单分区处理时间建议批量大小
不到 1 分钟20–50
1–5 分钟10–20
5–15 分钟5–10
超过 15 分钟1–5,或单次运行

根据实际失败率与基础设施限制调整。批次内部并行处理的方法见下文。

方案三:单次运行

BackfillPolicy.single_run 在一次运行中处理所有选中的分区,消除了逐次启动带来的开销。100 个分区只产生 1 次运行。不过,失败时必须一起重试全部分区。

使用 Spark、Snowflake、Databricks 等并行处理引擎,查询天然按日期范围运行,或者希望尽量减少 Dagster+ 额度消耗时,这种方式很理想。

文件:src/project_mini/defs/partition_backfill_strategies/single_run_backfill.py

import dagster as dg

daily_partitions = dg.DailyPartitionsDefinition(start_date="2024-01-01")


@dg.asset(
    partitions_def=daily_partitions,
    backfill_policy=dg.BackfillPolicy.single_run(),
)
def daily_events(context: dg.AssetExecutionContext):
    """Process events for multiple days in a single run."""
    start_datetime, end_datetime = context.partition_time_window
    context.log.info(f"Processing events from {start_datetime} to {end_datetime}")

    # Process data for the entire partition range at once
    events = fetch_events_for_range(start_datetime, end_datetime)
    processed = transform_events(events)

    context.log.info(f"Processed {len(processed)} events in single run")
    return processed


def fetch_events_for_range(start, end) -> list:
    # Simulate fetching events for a date range (e.g., SQL WHERE clause)
    return [{"start": str(start), "end": str(end), "event_id": i} for i in range(1000)]


def transform_events(events: list) -> list:
    # Simulate transformation
    return [{"processed": True, **e} for e in events]
场景推荐策略
API 限流或临时失败每分区一次
处理时间短、数据源可靠分批,每次 10–50 个
Spark/Snowflake 范围查询单次运行
优化 Dagster+ 成本单次运行或大批次
首次回填 1000 个以上分区分批,每次 50–100 个

在批次内部并行处理

使用 BackfillPolicy.multi_run 时,一次运行会包含多个分区。下面介绍在同一次运行内并行处理的不同方法。

策略适用场景最大并发开销复杂度
批量查询SQL 数据库不适用,单次查询很低很低
线程池I/O 密集任务10–100 个线程低低
进程池CPU 密集任务CPU 核心数量中等低

策略一:批量查询,对数据库最快

在单次数据库查询中处理全部分区。

文件:src/project_mini/defs/partition_backfill_strategies/parallel_batch_query.py

import dagster as dg

customer_partitions = dg.StaticPartitionsDefinition(
    ["customer_a", "customer_b", "customer_c", "customer_d", "customer_e"]
)


@dg.asset(
    partitions_def=customer_partitions,
    backfill_policy=dg.BackfillPolicy.multi_run(max_partitions_per_run=10),
)
def process_customers_batch(context: dg.AssetExecutionContext) -> None:
    """Process all partitions in a single database query."""
    customer_ids = context.partition_keys

    # Single query for all customers - most efficient for databases
    query = "SELECT * FROM customers WHERE customer_id IN %s"
    results = execute_query(query, customer_ids)

    context.log.info(f"Processed {len(results)} customers in single query")


def execute_query(query: str, params: list) -> list:
    # Simulated database query
    return [{"customer_id": p, "data": "..."} for p in params]

最适合 SQL 数据库,以及提供批量端点的 REST API。

策略二:线程池,适合 I/O 密集操作

使用线程并行执行 I/O。

文件:src/project_mini/defs/partition_backfill_strategies/parallel_threadpool.py

from concurrent.futures import ThreadPoolExecutor

import dagster as dg

customer_partitions = dg.StaticPartitionsDefinition(
    ["customer_a", "customer_b", "customer_c", "customer_d", "customer_e"]
)


@dg.asset(
    partitions_def=customer_partitions,
    backfill_policy=dg.BackfillPolicy.multi_run(max_partitions_per_run=10),
)
def fetch_customer_data(context: dg.AssetExecutionContext) -> None:
    """Use threads for parallel I/O operations."""
    customer_ids = context.partition_keys

    def fetch_one(customer_id):
        # Simulated API call
        return {"customer_id": customer_id, "data": f"data for {customer_id}"}

    # Process up to 5 customers concurrently
    with ThreadPoolExecutor(max_workers=5) as executor:
        results = list(executor.map(fetch_one, customer_ids))

    context.log.info(f"Fetched {len(results)} customers with thread pool")

最适合 HTTP 请求、文件 I/O 和数据库查询。并发度由 max_workers 限制,本例最多同时处理 5 个。

策略三:进程池,适合 CPU 密集操作

使用进程并行执行 CPU 密集工作。

文件:src/project_mini/defs/partition_backfill_strategies/parallel_processpool.py

from concurrent.futures import ProcessPoolExecutor

import dagster as dg

customer_partitions = dg.StaticPartitionsDefinition(
    ["customer_a", "customer_b", "customer_c", "customer_d", "customer_e"]
)


def analyze_one(customer_id: str) -> dict:
    """CPU-intensive analysis - must be defined at module level for ProcessPool."""
    # Simulated CPU-intensive work
    result = sum(i * i for i in range(100000))
    return {"customer_id": customer_id, "analysis": result}


@dg.asset(
    partitions_def=customer_partitions,
    backfill_policy=dg.BackfillPolicy.multi_run(max_partitions_per_run=10),
)
def analyze_customer_data(context: dg.AssetExecutionContext) -> None:
    """Use processes for parallel CPU-intensive work."""
    customer_ids = context.partition_keys

    # Use multiple CPU cores
    with ProcessPoolExecutor(max_workers=4) as executor:
        results = list(executor.map(analyze_one, customer_ids))

    context.log.info(f"Analyzed {len(results)} customers with process pool")

最适合 CPU 密集型计算和数据转换。并发能力受 CPU 核心数量限制。


原文:Partition backfill strategies。作者/维护者:Dagster 文档团队。本文为原文的中文译文;代码保留原文内容。

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

请登录后发表评论

    暂无评论内容