DuckDB

调度 DuckDB 管道,而不必在数据库前面再放一个数据库。

DuckDB 在进程内运行且没有服务端,这正是它没有调度器、没有重试、也没有运行历史的原因。Dagu 用单个二进制文件补上这层运维能力,整个技术栈仍然只是两个可执行文件和一个数据目录。

一个二进制编排另一个二进制,没有需要运维的数据库
max_active_runs 保证 DuckDB 的单写入者约束
增量加载的游标可跨进程重启保留
在数据所在之处运行,包括封闭网络内部
01

嵌入式意味着按设计就没有调度器

DuckDB 运行在你的进程内部。没有守护进程,没有端口,也没有任何东西在凌晨三点等着执行查询。这是设计意图而非疏漏,意味着调度必须来自外部。

  • DuckDB 管道通常就是一次 CLI 调用,这让 cron 成为默认答案,而大多数团队也止步于此
  • cron 只提供执行:S3 超时时不会重试,没有历史记录,昨晚静默地什么也没做时也没有任何信号
  • 社区的 cron 扩展在进程内调度,进程结束调度也随之消失,且不留下运行历史
02

不要抵消你选择 DuckDB 的理由

DuckDB 的吸引力在于没有集群、没有服务器、没有数仓账单。在它前面放一个重量级编排器会把这些全都还回去。Airflow 需要调度器、元数据数据库和 Python DAG 框架。Kestra 需要 JDBC 数据库、对象存储和四个组件。Dagu 是一个以文件保存状态的二进制文件。

  • 为了调度一个无服务器的分析引擎而引入 Postgres,是一笔奇怪的交易
  • 工作流定义留在 git 中、紧挨着它运行的 SQL,而不是另一个平台的元数据里
  • 整套栈可以放在一台主机上,而那通常就是数据已经所在的主机
03

单写入者是一个调度问题

DuckDB 生产环境的标准建议是用操作系统级的锁防止写入重叠,因为数据库只接受一个写入者。当这道防线是 cron 加锁文件时,它离出错只有一个脚本之遥。把它声明在工作流上,可以消除这一类缺陷。

  • max_active_runs 设为 1 时,慢的一次运行会延后下一次,而不是损坏数据库
  • resources.limits.memory 在大型 join 拖垮主机之前限制运行
  • retry_policy 覆盖真正短暂的失败,例如扫描中途对象存储超时

这约束的是该工作流的并发运行,并不会让 DuckDB 支持多写入者:其他进程写入同一文件仍需你自行防范。

对对象存储上的 Parquet 做夜间汇总
# 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
04

增量加载需要比进程更长寿的游标

每次运行都全量重载,会把一个快速廉价的本地查询变成缓慢昂贵的查询。增量加载需要记住上次成功运行的位置,而这份记忆必须在没有数据库存放它的情况下跨越重启。

  • Dagu 跨运行保存一个小的 JSON 游标,因此每次运行只读取尚未加载的窗口
  • 游标在加载步骤成功之后才保存,失败的运行会原样保留它,下次运行重试同一窗口
  • 在两端界定窗口,可以避免与运行期间到达的新行发生竞争
带持久化水位线的增量追加
# 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
05

这个组合较弱的场景

Dagu 负责调度 DuckDB,但不会改变 DuckDB 本身,有些工作负载并不适合这个组合。

  • DuckDB 是单节点、单写入者。如果多个服务需要并发写入,任何编排器都解决不了
  • 它不适合高频小写入,那仍然是 PostgreSQL 的职责
  • 如果你已经在运维 Airflow 或自带调度器的数仓,为一条管道再加一个编排器很少划算

FAQ

Practical questions before adopting

有 DuckDB 专用的执行器或插件吗?

没有,这是有意为之。DuckDB 提供了单二进制的 CLI,本身就能完成工作,所以 Dagu 把它当作普通命令步骤运行。专用执行器只会增加一层需要跟随 DuckDB 版本更新的东西,却提供不了比 CLI 更多的能力。

如何避免两次运行损坏数据库文件?

在工作流上把 max_active_runs 设为 1。仍在进行的运行会阻塞下一次计划运行,而不是打开第二个写入者,这就是 DuckDB 文档让你自己搭建的锁文件的声明式版本。它只管辖该工作流的运行,因此仍需让其他进程远离同一个文件。

大查询耗尽内存时会怎样?

在工作流上设置 resources.limits.memory,让运行在拖垮主机之前被限制;并为短暂性而非结构性的失败配置重试。一个需要超过主机内存的查询应当重写,而不是重试。

DuckDB 有社区 cron 扩展,为什么不用它?

它在 DuckDB 进程内调度,因此调度只在该进程存活期间存在。没有运行历史、没有重试、任务失败时没有告警,第二天早上也没有可查看的东西。它适合长期运行的嵌入式应用,而不是服务器上的批处理。

计划运行中 DuckDB 能读取 S3 吗?

可以,通过 httpfs 扩展,这也是大多数计划内 DuckDB 作业无需单独抽取步骤即可读取 Parquet 的方式。由于对象存储是网络依赖,这一步正是值得配置重试策略的地方。

Next step

Start with one workflow.

Install Dagu, move one script that runs on cron today into YAML, and decide from a real run history.