本文介绍三种回填分区资产的策略。为了初始化或重新处理数据而需要物化多个分区时,可以选择 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 文档团队。本文为原文的中文译文;代码保留原文内容。











暂无评论内容