嵌入式意味着按设计就没有调度器
DuckDB 运行在你的进程内部。没有守护进程,没有端口,也没有任何东西在凌晨三点等着执行查询。这是设计意图而非疏漏,意味着调度必须来自外部。
- DuckDB 管道通常就是一次 CLI 调用,这让 cron 成为默认答案,而大多数团队也止步于此
- cron 只提供执行:S3 超时时不会重试,没有历史记录,昨晚静默地什么也没做时也没有任何信号
- 社区的 cron 扩展在进程内调度,进程结束调度也随之消失,且不留下运行历史
DuckDB 运行在你的进程内部。没有守护进程,没有端口,也没有任何东西在凌晨三点等着执行查询。这是设计意图而非疏漏,意味着调度必须来自外部。
DuckDB 的吸引力在于没有集群、没有服务器、没有数仓账单。在它前面放一个重量级编排器会把这些全都还回去。Airflow 需要调度器、元数据数据库和 Python DAG 框架。Kestra 需要 JDBC 数据库、对象存储和四个组件。Dagu 是一个以文件保存状态的二进制文件。
DuckDB 生产环境的标准建议是用操作系统级的锁防止写入重叠,因为数据库只接受一个写入者。当这道防线是 cron 加锁文件时,它离出错只有一个脚本之遥。把它声明在工作流上,可以消除这一类缺陷。
这约束的是该工作流的并发运行,并不会让 DuckDB 支持多写入者:其他进程写入同一文件仍需你自行防范。
# duckdb-nightly-rollup.yaml
schedule: "0 3 * * *"
max_active_runs: 1
resources:
limits:
memory: "8Gi"
steps:
- id: rollup
action: duckdb@v1
with:
database: /data/analytics.duckdb
query: |
INSTALL httpfs; LOAD httpfs;
CREATE OR REPLACE TABLE daily_sales AS
SELECT order_date, region, sum(amount) AS amount
FROM read_parquet('s3://warehouse/raw/orders/*.parquet')
GROUP BY 1, 2;
retry_policy:
limit: 2
interval_sec: 300
- id: export
action: duckdb@v1
with:
database: /data/analytics.duckdb
query: |
COPY daily_sales TO '/data/export/daily_sales.parquet' (FORMAT parquet);
depends: rollup
- id: count_rows
action: duckdb@v1
with:
database: /data/analytics.duckdb
readonly: true
query: SELECT count(*) AS row_count FROM daily_sales;
depends: export
- id: verify
env:
- COUNT_JSON: ${steps.count_rows.outputs.result}
run: test "$(printf '%s\n' "$COUNT_JSON" | jq -r '.[0].row_count')" -gt 0
depends: count_rows
handler_on:
failure:
run: /opt/analytics/notify-failure.sh
mail_on:
failure: true
每次运行都全量重载,会把一个快速廉价的本地查询变成缓慢昂贵的查询。增量加载需要记住上次成功运行的位置,而这份记忆必须在没有数据库存放它的情况下跨越重启。
# duckdb-incremental-load.yaml
schedule: "*/15 * * * *"
max_active_runs: 1
steps:
- id: load_cursor
action: state.get
output: CURSOR
with:
key: cursors/events-loaded-through
default:
loaded_through: "2026-01-01T00:00:00Z"
- id: window
run: |
printf 'since=%s\n' "$(printf '%s\n' "$CURSOR" | jq -r .value.loaded_through)" >> "$DAGU_OUTPUT_FILE"
printf 'until=%s\n' "$(date -u +%Y-%m-%dT%H:%M:%SZ)" >> "$DAGU_OUTPUT_FILE"
outputs:
- name: since
- name: until
depends: load_cursor
- id: append_new_events
action: duckdb@v1
with:
database: /data/analytics.duckdb
query: |
INSTALL httpfs; LOAD httpfs;
INSERT INTO events
SELECT * FROM read_parquet('s3://warehouse/events/*.parquet')
WHERE ingested_at > TIMESTAMP '${steps.window.outputs.since}'
AND ingested_at <= TIMESTAMP '${steps.window.outputs.until}';
depends: window
retry_policy:
limit: 3
interval_sec: 60
- id: save_cursor
action: state.set
with:
key: cursors/events-loaded-through
value:
loaded_through: "${steps.window.outputs.until}"
depends: append_new_events
mail_on:
failure: true
Dagu 负责调度 DuckDB,但不会改变 DuckDB 本身,有些工作负载并不适合这个组合。
FAQ
没有,这是有意为之。DuckDB 提供了单二进制的 CLI,本身就能完成工作,所以 Dagu 把它当作普通命令步骤运行。专用执行器只会增加一层需要跟随 DuckDB 版本更新的东西,却提供不了比 CLI 更多的能力。
在工作流上把 max_active_runs 设为 1。仍在进行的运行会阻塞下一次计划运行,而不是打开第二个写入者,这就是 DuckDB 文档让你自己搭建的锁文件的声明式版本。它只管辖该工作流的运行,因此仍需让其他进程远离同一个文件。
在工作流上设置 resources.limits.memory,让运行在拖垮主机之前被限制;并为短暂性而非结构性的失败配置重试。一个需要超过主机内存的查询应当重写,而不是重试。
它在 DuckDB 进程内调度,因此调度只在该进程存活期间存在。没有运行历史、没有重试、任务失败时没有告警,第二天早上也没有可查看的东西。它适合长期运行的嵌入式应用,而不是服务器上的批处理。
可以,通过 httpfs 扩展,这也是大多数计划内 DuckDB 作业无需单独抽取步骤即可读取 Parquet 的方式。由于对象存储是网络依赖,这一步正是值得配置重试策略的地方。