跳到主要内容

7 篇博文 含有标签「Airflow」

查看所有标签

容器日志吃满服务器磁盘?docker system df 的 reclaimable 是误报

· 阅读需 7 分钟

告警邮件:生产服务器根分区用到 81%(30G/40G,超过 80% 红线)。第一反应是查 docker system df——它显示 images 一栏 RECLAIMABLE 100% (7),看起来删掉镜像就能回血;但这 7 个镜像全在运行,真按提示 docker image prune -a 就是生产事故。

在开发 AI运营 时遇到此问题——基于大语言模型的智能分析平台,自动洞察市场趋势、用户行为与销售数据;出问题的是跑它 Airflow 数据管道的 Docker 服务器。

TL;DR

磁盘超红线时别信 docker system df 的 RECLAIMABLE——它是「无容器引用空间」的估算口径,active 镜像照样标 100%。正确姿势:sudo du -xh -d1 / 逐层找大头。这次的三大隐形消耗都不在业务数据里:容器可写层里无限累积的 task 日志(Dockerfile 没给日志配卷)、包管理器缓存(npm + pnpm 共 ~4G)、只增不删的备份脚本(每次全量克隆从不清理)。对应清理:缓存直删、可写层 force-recreate 回收、备份脚本加 TTL 剪枝——81% 回落到 64%。

问题现象

$ df -h /
Filesystem Size Used Use% Mounted on
/dev/vda1 40G 30G 81% /

$ docker system df
TYPE TOTAL ACTIVE SIZE RECLAIMABLE
Images 7 7 4.2GB 100% (7) ← 全是运行中镜像
Containers 5 5 810MB 0%

docker system df 的提示去做文章(清理镜像)无从下手——RECLAIMABLE 100% 但 ACTIVE 也是 7/7。真正的空间去哪了,docker system df 完全没体现:

sudo du -xh -d1 / | sort -rh | head
# 部分输出:
# 2.7G /root/.npm ← npm 缓存
# 1.1G /root/.local/share/pnpm ← pnpm store
# 857M /root/backups-git ← 备份工作区,只增不删
# (overlay2 内藏:airflow scheduler 可写层 468M + dag-processor 257M)

根因

三类消耗在「业务数据」视角里全部隐形。

容器可写层吃日志。 Airflow 的 task 日志写在容器内部、没挂外置卷——scheduler 可写层 468M、dag-processor 257M,约两周累积 780M,随时间单调增长。镜像层不变,变的都是可写层;docker system df 把它归在 Containers 的 SIZE 里,混在 810M 的总数里毫无存在感。

包管理器缓存只进不出。 频繁部署/构建的服务器上,~/.npm(2.7G)和 pnpm store(1.1G)持续累积,没有人会主动去清。

备份脚本只增不删。 备份脚本每次全量克隆仓库工作区推 GitHub,旧的克隆目录从不剪枝——857M 里大部分是历史重复。更隐蔽的是孤本:某些备份目录对应的分支从未推送成功,直接删有风险,需要先核对。

docker system df 的 RECLAIMABLE 是按「没有运行容器引用」估算的统计口径,active 镜像也可能显示 100% reclaimable(本次 7/7 全如此)——它是误报来源,不是清理依据。

解决方案

步骤 1:du 逐层定位,先看再删

sudo du -xh -d1 / | sort -rh | head        # 根分区逐层
sudo du -xh -d1 /var/lib/docker | sort -rh | head # Docker 目录下钻
docker ps -as --format "table {{.Names}}\t{{.Size}}" # 各容器可写层

du -x 不跨文件系统,避开 proc/sys 噪音和 overlay 混淆;docker ps -as 的 SIZE 列直接暴露每个容器的可写层大小——这是定位「日志写进容器」类问题的关键一条。

步骤 2:缓存直删

npm cache clean --force        # 或直接 rm -rf ~/.npm/_cacache
pnpm store prune # 只清未引用的包

缓存类共回收 ~5.1G,无风险、可随时重下。

步骤 3:可写层 force-recreate 回收

docker compose up -d --force-recreate   # 可写层随旧容器删除而回收

本次回收 ~780M。两个前提:选低峰期(服务会重启);容器内还有用要先捞出来——比如 docker cp 把 task 日志拷出,否则随容器一起没了。

步骤 4:备份脚本加 TTL,防复发

# 备份目录保留 7 天,超过自动剪
find /root/backups-git -maxdepth 1 -type d -mtime +7 -exec rm -rf {} +

把剪枝加进备份脚本尾部(本例 github-backup-push.sh 已部署);删除前核对未推送的孤本分支(git ls-remote 对比),确认无独有提交再删——本次 325M 孤本经确认后清理。

清理完成后:81% → 64%,缓存/剪枝/可写层三处合计回收 ~6.7G。

注意事项

  • docker system df 的 RECLAIMABLE ≠ 可删空间:active 镜像也可能标 100%(本次 7 个全在用)。清理决策以 du + docker ps -as 实测为准。
  • force-recreate 会重启服务且回收可写层——需要留存的日志先 docker cp 出来;根治是给日志配外置卷 + 保留期轮转,可写层回收只是这一次的止血。
  • 删备份目录前必核对:有没有从未推送成功的孤本分支,删之前 git ls-remote 确认。
  • 磁盘红线告警要配「定位命令」一起给,收到告警的人第一件事是 du -xh -d1 /,而不是猜。
  • 同一台服务器更早的「磁盘 93% + CPU 160%」排查复盘见 2 核 7G 服务器 Docker 资源黑洞三步排查

常见问题

docker system df 的 RECLAIMABLE 显示 100% 能直接删吗?

不能。RECLAIMABLE 是「当前无容器引用」空间的估算口径,运行中的 active 镜像也可能被标 100% reclaimable——实测 7 个镜像全在用仍显示如此。它回答的是「理论上有多少未被引用」,不回答「能不能删」。动手前用 du -xh -d1docker ps -as 实测定位。

docker ps 的 SIZE 和 docker system df 的 SIZE 有什么区别?

docker ps -as 的 SIZE 是单个容器可写层的大小;docker system df 是镜像/容器/卷/缓存的分类汇总。容器内日志无限累积只体现在可写层(docker ps -as 能看到),不会体现在任何 image size 上——这就是「镜像看起来不大、磁盘却在涨」的原因。

Docker 服务器磁盘满了怎么清理?

分三类:包管理缓存直删(npm cache clean --forcepnpm store prune),无风险;容器可写层用 docker compose up -d --force-recreate 回收(会重启服务,日志先备份);备份、日志、克隆目录这类加 TTL 定期剪枝防复发。任何删除动作前,先用 du -xh -d1 / 确认大头位置。

CCLEE

独立开发者,24年电商行业实战经验,专注将AI能力落地于真实商业场景。

合作咨询

数据管道快照槽位错位?0 行段丢弃导致位置漂移

· 阅读需 6 分钟

核对一次生产任务的决策快照时,发现第 5 个阶段的特征数据落在了数组槽 3,而不是设计文档里写的槽 4——下游和抽屉组件按「段号−1」取值,取到的是上一阶段的输出。

在开发 AI运营 时遇到此问题——基于大语言模型的智能分析平台,自动洞察市场趋势、用户行为与销售数据;快照是管道留给前端展示和事后审计的决策依据。

TL;DR

管道把各阶段(段)输出按执行顺序 push 进快照数组,而 if rows.empty: skip 会把 0 行段整个丢弃,后续所有段前移一位——「段号−1 = 槽位」的静态映射随时被打破,且哪种段为空取决于运行态,槽位每次都可能不同。解法两条路:消费端按行内键名实时认槽(推荐),或写入端给空段保留占位、维持槽位恒定。

问题现象

装配逻辑长这样:

snapshot = {"features": [], "rule_output": []}
for seg in segments: # 段①…段⑤ 顺序执行
df = execute_sql(seg.sql)
if not df.empty: # 0 行段在这里被丢弃
snapshot["features"].append(df.to_dict("records"))

设计假设是「段⑤ → 槽 4」。但生产快照核验发现段⑤的特征落在槽 3:

全段有数据:      段①→0  段②→1  段③→2  段④→3  段⑤→4   ✓ 符合假设
段④ 空表被弃: 段①→0 段②→1 段③→2 段⑤→3 ✗ 前移
段③④ 都空: 段①→0 段②→1 段⑤→2 ✗ 再前移

同一个代码版本,不同店铺/不同权限下快照槽位完全不同——保护词白名单是空表时段④被弃,权限关闭时段③被弃,槽位跟着运行态漂移。

根因

位置寻址撞上了稀疏装配。 快照数组是运行时把「有输出的段」压缩拼接的产物,本质是个稀疏集合;而下游按「段号−1」硬编码取值,等价于假设「每个段必然产出至少一行」。这个假设在三种常见情形下都会碎:白名单空表、功能开关关闭、业务数据天然为空——0 行是常态而不是异常。

更深一层,if not df.empty 这个判空本身没写错,错的是契约的隐含前提:设计文档写了「槽位 = 段号−1」,却没人把它声明成显式契约。所有按位置消费的下游都在继承一个未被承认、也无人维护的假设。

解决方案

方案 A(推荐):消费端按行键名认槽

让每行数据自带段标识键,消费方在读取时实时解析位置,不做任何静态映射:

def locate_segment(features: list, seg_key: str) -> dict:
for row in features:
if seg_key in row: # 行内自带段标识,按内容寻址
return row
raise KeyError(f"segment '{seg_key}' missing in snapshot")

槽位漂移从此无关紧要——找的是「键名长这样的段」,不是「第 N 个元素」。唯一要求是所有消费方统一走这个解析入口(写进 processor docstring 和消费方契约,明确禁止硬编码槽位)。

方案 B:写入端保留空段占位

如果下游暂时改不动,可以让装配端维持「槽位 = 段号」恒定:

snapshot["features"].append(
df.to_dict("records") if not df.empty else {"__empty__": True}
)

代价是快照里出现占位对象,所有消费方都得处理它;作为过渡方案可用,长期仍建议收敛到方案 A。

步骤 3:用多种运行态做契约测试

把「全段有数据 / 单段空 / 多段空」三种运行态做成快照 fixture,断言消费方在三种形态下解析结果一致。只测全满场景,等于没测。

注意事项

  • 规格文档里任何「位置对应关系」都必须显式声明寻址方式(按键名/按 ID),并注明「禁止按下标硬编码」;隐含假设一定会被某个运行态打破。
  • 判空跳过(if empty: skip)是最常见的压缩来源——同类静默丢数据还有 Airflow PostgresHook 多语句 SQL 只返回第一段结果,同样是「不报错、悄悄少东西」。
  • 改造消费方时,先用三种运行态 fixture 回归,再上生产;只验证「全段有数据」的场景会漏掉全部错位路径。

常见问题

数据工程里怎么处理 schema drift?

把位置契约换成键名契约:快照、消息、接口按字段名或段标识寻址,而不是数组下标。上游发生未经约定的结构变化(空段被跳过、字段增删)时,按名寻址的下游最多报「找不到」,不会静默拿到错误数据。

schema drift 和 schema evolution 有什么区别?

Schema evolution 是显式管理的版本演进(加字段、发版本、迁移消费方);schema drift 是被动发生的漂移——上游一改、下游不知不觉错位。本例的槽位前移就是典型 drift:没人改契约,是数据形态变了。

怎么检测数据管道里这类槽位错位?

两层:契约测试覆盖多种运行态(全段有数据/单段空/多段空),断言消费方解析一致;生产侧定期抽检快照,核对槽位内容自带的段标识与预期段是否对应。发现「内容与位置对不上」即是 drift。

CCLEE

独立开发者,24年电商行业实战经验,专注将AI能力落地于真实商业场景。

合作咨询

Airflow 触发 dagRun 静默失败?logical_date 唯一约束在作怪

· 阅读需 5 分钟

在用 Airflow REST API 重复触发同一个 DAG 做灰度验证时,请求返回了 4xx 但响应体里没有 dag_run_id,DAG 实际根本没有运行——而脚本却把它当成了成功。

在开发 AI 运营 时遇到此问题——基于大语言模型的智能分析,自动洞察市场趋势、用户行为、销售数据,提供精准运营策略。广告决策链路的灰度切换需要在 Airflow 上反复触发同一次分析做对照,结果部分触发悄无声息地失败了。

TL;DR

Airflow 对每个 DAG 的 logical_date 有唯一约束dag_run_id 同样必须唯一)。用相同的 logical_date 重复 POST /dags/{dag_id}/dagRuns,Airflow 会拒绝并返回 4xx,响应体里没有 dag_run_id。如果你只检查 HTTP 状态码、不检查返回的 dag_run_id,就会误以为触发成功。解法:每次触发用不同的 logical_date(和不同的 dag_run_id)。

问题现象

为了对照测试,用固定日期 2026-01-01 连续触发同一个 DAG:

$ curl -s -X POST "$AIRFLOW/api/v2/dags/my_dag/dagRuns" \
-H "Authorization: Bearer $TOKEN" \
-H "Content-Type: application/json" \
-d '{
"dag_run_id": "manual-run-1",
"logical_date": "2026-01-01T00:00:00Z"
}'
# 第一次:返回正常的 dag_run 对象,包含 dag_run_id ✅

$ curl -s -X POST "$AIRFLOW/api/v2/dags/my_dag/dagRuns" \
-H "Authorization: Bearer $TOKEN" \
-H "Content-Type: application/json" \
-d '{
"dag_run_id": "manual-run-2",
"logical_date": "2026-01-01T00:00:00Z" # ⚠️ 同一个 logical_date
}'
# 第二次:返回错误对象,没有 dag_run_id ❌
{
"detail": "...",
"status": 400,
"title": "Bad Request",
"type": "https://airflow.apache.org/docs/apache-airflow/2/stable-rest-api-ref.html#/default/Error"
}

如果调用方只判断「HTTP 是否 2xx」就停止解析,或者直接读 JSON 不校验 dag_run_id 字段,第二次失败就会被静默吞掉——日志里看不到异常,Airflow UI 里也找不到这次 run。

根因

Airflow 用 dag_run_id 作为每次运行的主键,同时在 metadata 数据库的 dag_run 表上对 (dag_id, logical_date) 维护唯一性。logical_date 是调度的「逻辑时间」——调度器按它判断某个调度槽位是否已经跑过。一旦同一个 DAG 下已存在某 logical_date 的 run,再用相同值触发,Airflow 就会拒绝,避免重复执行。

问题在于这个失败是 HTTP 4xx + 错误 JSON,不是连接错误或 5xx。很多脚本只做 response.status_code == 200 的粗判断,或者拿到 JSON 后直接取字段而不校验是否存在 dag_run_id,于是把「拒绝创建」当成了「创建成功」。

解决方案

核心:每次触发用不同的 logical_date(以及不同的 dag_run_id)。 灰度/回放场景下,给每次触发拼一个不重复的日期即可:

# 每次循环用不同的 logical_date(2026-01-01 / 02 / 03 …)
for i in 1 2 3; do
curl -s -X POST "$AIRFLOW/api/v2/dags/my_dag/dagRuns" \
-H "Authorization: Bearer $TOKEN" \
-H "Content-Type: application/json" \
-d "{
\"dag_run_id\": \"manual-run-$i\",
\"logical_date\": \"2026-01-0${i}T00:00:00Z\"
}"
done

更稳妥的是用递增时间戳,保证 logical_datedag_run_id 永不重复。更重要的是:必须校验响应体里的 dag_run_id 字段,把它当成「触发真正成功」的唯一证据:

import requests

def trigger_dag(dag_id: str, logical_date: str, conf: dict | None = None) -> str:
resp = requests.post(
f"{AIRFLOW}/api/v2/dags/{dag_id}/dagRuns",
headers={"Authorization": f"Bearer {TOKEN}", "Content-Type": "application/json"},
json={"dag_run_id": f"manual-{logical_date}", "logical_date": logical_date, "conf": conf or {}},
)
# ❌ 不够:只看状态码,4xx 会被当异常但容易漏判
# resp.raise_for_status()
data = resp.json()
# ✅ 正确:dag_run_id 存在才算真正创建成功
if "dag_run_id" not in data:
raise RuntimeError(f"Trigger failed: {resp.status_code} {data}")
return data["dag_run_id"]

# 每次用不同 logical_date,重复触发安全
for i in range(1, 4):
trigger_dag("my_dag", f"2026-01-0{i}T00:00:00Z")

dag_run_id 同样要保持唯一——它是主键,重复会被直接拒绝。用「前缀 + logical_date」组合是常见做法,既唯一又能在 UI 里一眼识别。

常见问题

airflow trigger_dagrun 怎么通过 REST API 触发 DAG?

POST /api/v2/dags/{dag_id}/dagRuns,请求体至少包含 dag_run_idlogical_date 两个字段(可选 conf 传参数)。这两个字段在同一个 DAG 下都必须唯一,否则 Airflow 返回 4xx。代码里推荐用 TriggerDagRunOperator,它内部也会生成唯一的 run id。

airflow 重复触发同一个 DAG 为什么失败?

因为 Airflow 在 metadata 数据库的 dag_run 表上对 (dag_id, logical_date) 维护唯一约束,dag_run_id 本身也是主键。重复的 logical_datedag_run_id 都会被拒绝并返回 4xx。要做回放或灰度对照,每次触发换一个新的 logical_date(或递增时间戳)即可。

注意事项

  • 校验 dag_run_id,不要只看状态码:4xx + 错误 JSON 是 Airflow 表达「拒绝创建」的正常方式,只判断 status_code 容易把失败误判为成功。
  • logical_date 用过去时间:填未来时间会被当成定时调度,不会立即执行;要立即跑用过去日期。
  • API 版本差异:Airflow 2.x 是 /api/v2/dags/{dag_id}/dagRuns,3.x 路径和字段有调整,迁移时务必核对对应版本的 REST API 文档。
  • 回放优先用 CLI 的 --logical-dateairflow dags trigger 支持指定日期,但对 logical_date 的唯一约束同样生效,重复日期一样会失败。

CCLEE

独立开发者,24年电商行业实战经验,专注将AI能力落地于真实商业场景。

合作咨询

Airflow 删除 DAG 后它还在列表里?元数据没清干净 + 正确清理顺序

· 阅读需 7 分钟

在 Airflow 里删掉某个 DAG 的 .py 文件想下线它,Web UI 的 DAG 列表和数据库里却还挂着这个 dag_id;更诡异的是——如果先清元数据再删文件,刚刚清空的行会「复活」重新出现。

在开发 AI运营 时遇到此问题——基于大语言模型的智能分析流水线,下线旧版报表 DAG 时要把它的元数据一起清干净,否则 UI 列表和定时扫描会被残留行干扰。

TL;DR

Airflow 的 dag-processor 会定期扫描 DAG 目录重新注册,而 airflow dags reserialize 也不会清理「文件已删除」的孤儿 dag 行——所以光删 .py 文件,UI 和数据库里的 dag_id 不会自动消失;反过来先清元数据再删文件,processor 扫到文件还在,会把清空的行重新注册(「复活」)。正确顺序:①先删文件让 processor 不再注册 → ②按外键顺序 SQL DELETE 清元数据 → ③跑 airflow dags reserialize 验证

问题现象

下线 shop_report_aggregation 这个 DAG,删了它的 .py 文件后:

$ ls /opt/airflow/project/airflow_dags/shop_report_aggregation.py
ls: cannot access '.../shop_report_aggregation.py': No such file or directory

$ # 但数据库里还在
$ docker exec cclhub-db psql -U airflow -d airflow -c \
"SELECT dag_id, is_paused, is_active FROM dag WHERE dag_id='shop_report_aggregation';"
dag_id | is_paused | is_active
--------------------------+-----------+-----------
shop_report_aggregation | f | t ← 仍残留

不只 dag 表,serialized_dagdag_codedag_version 表里对应行也全在,于是 Web UI 的 DAG 列表继续显示这个已「删除」的 DAG。

更坑的是反向操作——先清元数据、后删文件:

T0  DELETE FROM dag WHERE dag_id='shop_report_aggregation';   ← 清空
T1 (此时还没删 .py 文件)
T2 dag-processor 扫描周期到达,发现文件存在、dag 表无对应行 → 重新注册
T3 SELECT ... FROM dag WHERE dag_id='shop_report_aggregation'; ← 又回来了(复活)

根因

两个机制叠加:

1. dag-processor 定期扫描并重新注册。 Airflow 的 dag-processor(Scheduler 的一部分)按 processor_poll_interval(默认约 5 分钟)周期性扫描 dags_folder 目录,解析每个 .py 文件并 upsert 进元数据表(dagserialized_dagdag_version)。只要文件还在,下一个扫描周期就会重新写入对应行。 这是「复活」的直接来源——你清了行,文件还在,processor 把它当新 DAG 重新登记。

2. reserialize 不管「文件已消失」的旧行。 airflow dags reserialize 的职责是把现有 DAG 文件重新序列化、刷新 serialized_dag;它不会去删除「文件已经不存在」的孤儿 dag 行。而 airflow dags cleanup 默认只清理过期的 dag_run 运行历史,也不动 dag / serialized_dag / dag_code / dag_version 这几张元数据表。所以删了文件后,元数据行成了无人清理的孤儿。

┌─ dag-processor ──────────────────────────────┐
│ 扫描 dags_folder │
│ ├─ 文件在 → upsert dag / serialized_dag ... │ ← 复活来源
│ └─ 文件不在 → 跳过,不删旧行 │ ← 孤儿残留
└──────────────────────────────────────────────┘

结论:要让元数据真正消失,必须让 processor 没有文件可注册(先删文件),再手动清掉残留的元数据行。

解决方案

第 1 步:先删文件

.py 文件从 DAG 目录消失,dag-processor 就不会再注册它。

# 生产环境通常经 git pull 同步到 volume 挂载的 DAG 目录
# /opt/airflow/project/airflow_dags/
git pull # 让 shop_report_aggregation.py 从仓库移除并同步到目录

# 或直接删除(确认无其他依赖后)
rm /opt/airflow/project/airflow_dags/shop_report_aggregation.py

第 2 步:按外键顺序清元数据

按外键依赖顺序 DELETE,避免约束冲突。dag_run 删除会 CASCADE 到 task_instance

BEGIN;

-- 1. 运行历史(CASCADE 带 task_instance)
DELETE FROM dag_run WHERE dag_id = 'shop_report_aggregation';

-- 2. 序列化 DAG
DELETE FROM serialized_dag WHERE dag_id = 'shop_report_aggregation';

-- 3. 版本
DELETE FROM dag_version WHERE dag_id = 'shop_report_aggregation';

-- 4. dag 主表
DELETE FROM dag WHERE dag_id = 'shop_report_aggregation';

-- 5. dag_code 按源码 hash 存,多个 DAG 可能共享同一份代码;
-- 只删已经没有任何 serialized_dag 引用的 orphan hash
DELETE FROM dag_code
WHERE dag_hash NOT IN (SELECT dag_hash FROM serialized_dag);

COMMIT;

第 3 步:验证

airflow dags reserialize

# 确认 dag 表不再重建该行
docker exec cclhub-db psql -U airflow -d airflow -c \
"SELECT count(*) FROM dag WHERE dag_id='shop_report_aggregation';"
# count
# -------
# 0 ✅

reserializedag / serialized_dag / dag_code / dag_version 对该 dag_id 全部为 0,且下一个 processor 扫描周期过去也不再重建,说明清理稳定。

顺带一提,同一条流水线上 pandas NaN 进 XCom 导致任务无声崩溃是另一个值得收藏的坑。

注意事项

注意事项

  • dag_code 是按源码 hash 共享的:多个 DAG 可能引用同一份源码 hash,删除前务必用 orphan 判定(dag_hash NOT IN (SELECT dag_hash FROM serialized_dag)),不要按 dag_id 直接删——这张表压根没有 dag_id 列。
  • 别指望 airflow dags cleanup 清元数据:它默认只删过期的 dag_run(由 max_active_runs / retention 控制),不动 dag / serialized_dag / dag_code / dag_version。清元数据得手写 SQL。
  • 删文件后等一个扫描周期再清更稳:极端竞态下,删文件和清元数据之间若正好夹一个 processor 扫描,可能在文件已被你删但 processor 还没刷新的窗口里写入。实践中按「先删文件→再清元数据→reserialize 验证」顺序操作即可,必要时清完再跑一次 reserialize 确认。
  • 删 DAG 前确认无下游依赖:其他 DAG 可能用 ExternalTaskSensor 等待这个 DAG,或 TriggerDagRunOperator 触发它。下线前 grep 一遍 dag_id 引用。

常见问题

Airflow 删除 DAG 的 .py 文件后为什么还在列表里?

因为删除文件不会清理数据库元数据。dag / serialized_dag / dag_code / dag_version 这几张表的旧行仍然存在,Web UI 读这些表来渲染列表,所以已删的 DAG 还会显示。Airflow 没有内置命令自动清这些孤儿行,需要手动按外键顺序 SQL DELETE。

Airflow 怎么彻底删除一个 DAG 及其全部元数据?

三步:①先删 .py 文件,让 dag-processor 不再注册它;②按外键顺序 SQL DELETE 清理(dag_runserialized_dagdag_versiondag → orphan dag_code);③跑 airflow dags reserialize,然后查 dag 表确认该 dag_id 行数不再重建为 0。

Airflow 清理 DAG 元数据的正确顺序是什么?为什么不能先清元数据再删文件?

必须先删文件、后清元数据。反过来操作的话,.py 文件还在,dag-processor 下一个扫描周期会重新把清空的 dag 行注册回来,元数据「复活」。只有先让文件消失、processor 无文件可注册,再清残留的元数据行,才能彻底下线。


CCLEE

独立开发者,24年电商行业实战经验,专注将AI能力落地于真实商业场景。

合作咨询

Airflow XCom 报 'Out of range float values are not JSON compliant'?pandas NaN 惹的祸

· 阅读需 8 分钟

在 Airflow 任务用 ti.xcom_push() 把 pandas 处理后的结果推给下游任务时,任务直接崩溃——ValueError: Out of range float values are not JSON compliant: nan,而且应用自建的 logs 表里找不到任何错误记录。

在开发 AI运营 时遇到此问题——基于大语言模型的智能分析流水线,Airflow DAG 从 SQL 取数、经 pandas 处理后,通过 XCom 在任务间传递结果。

TL;DR

XCom 底层用 JSON 序列化,而 Airflow 调用 json.dumps(..., allow_nan=False) 严格遵循 JSON 标准——标准里压根没有 NaN / Infinity。pandas 从 SQL NULL 转来的浮点 NaN 一旦进了 xcom_push 的数据,序列化立即抛 ValueError。解法:在 push 前递归遍历数据,把 NaN / ±Inf 转成 None(JSON null)。

问题现象

某次季度报表 DAG 失败,但症状很迷惑——应用自建的 logs 表里,该 trace 的步骤 1–4 全部正常,连「分析完成」都打了两次(重试),然后就断了,步骤 5 缺失,且没有任何 error 行

trace=91c126c3
├─ step 1 SQL 取数 ✅
├─ step 2 pandas 处理 ✅
├─ step 3 规则判定 ✅
├─ step 4 LLM 分析完成 ✅ ← 之后重试了一次
└─ step 5 XCom 推送结果 ❌ ← 缺失,无 error 记录

去 Airflow 任务日志才看到真正的 traceback:

# 容器内 /opt/airflow/logs/dag_id=ai_analysis_v2/run_id=.../task_id=analyze_results/attempt=N.log
ValueError: Out of range float values are not JSON compliant: nan
File ".../ai_analysis_tasks.py", line 142, in analyze_results
ti.xcom_push(key='sql_metadata', value=result)

崩溃点精确落在 ti.xcom_push——任务把结果推给 XCom 的那一刻。

根因

三层叠加,缺一不可:

1. JSON 标准不含 NaN / Infinity RFC 8259 定义的 JSON 只允许数字字面量是有限数。虽然 Python 的 json.dumps 默认会把 NaN 写成裸 NaN、把 Infinity 写成 Infinity,但这是 Python 的私有扩展,不是合法 JSON——任何严格 parser(包括 Airflow 用的)读到都会拒绝。

2. Airflow XCom 序列化时显式 allow_nan=False XCom 默认走 JSON serializer,序列化时关闭了 NaN 容忍,遇到 NaN 直接抛 ValueError: Out of range float values are not JSON compliant,而不是偷偷写出非法 JSON。

3. pandas 把 SQL NULL 读成 NaN pandas.read_sql 对 SQL 的 NULL 列返回 float('nan')。一旦这列参与计算后被 to_dict('records') 带进结果对象,NaN 就顺着数据流进了 xcom_push

import pandas as pd

# SQL 某行某列是 NULL → pandas 读成 NaN
df = pd.DataFrame({"ad_roi": [1.2, None, 0.8]})
records = df.to_dict("records")
# [{'ad_roi': 1.2}, {'ad_roi': nan}, {'ad_roi': 0.8}] ← nan 混了进来

# 下游任务 push 时崩
ti.xcom_push(key="result", value=records)
# ValueError: Out of range float values are not JSON compliant: nan

这次之所以长期没触发,是因为平时跑的数据那些列都有值;直到某客户某个季度完全没有广告投放、ad_roi 整列 NULLNaN 才第一次大规模进入 XCom 路径。

为什么 logs 表没有 error? 因为崩溃发生在 XCom 序列化阶段,处于任务函数的 try/except 之外——异常直接冒泡给 Airflow 调度器,只写进 Airflow 自己的任务日志,应用层自建的 logs 表的 catch 根本没机会记录。这是这类故障最迷惑的地方:看起来「无声失败」。

解决方案

在数据进入 XCom 前,递归清洗掉所有 NaN / ±Inf

1. 写一个纯函数递归清洗

import math

def json_safe_value(obj):
"""
递归把 NaN / +Inf / -Inf 转成 None,使数据可被 JSON 严格序列化。
兼容 dict / list / tuple / scalar,遇到未知类型原样返回。
"""
if isinstance(obj, float):
if math.isnan(obj) or math.isinf(obj):
return None
return obj
if isinstance(obj, dict):
return {k: json_safe_value(v) for k, v in obj.items()}
if isinstance(obj, (list, tuple)):
return [json_safe_value(v) for v in obj]
return obj

为什么不能用 df.fillna(None)?因为 pandas 的 fillna(None) 在数值列上行为依版本和 dtype 不稳定,有时会把 NaN 强制转型而非置空;而且它只处理 DataFrame,管不到已经 to_dict 之后嵌在 dict/list 里的浮点。递归清洗在「数据已变成 Python 原生结构」这一层兜底,最稳。

2. 在 push 前统一兜底

最省心的做法是把清洗挂在所有 xcom_push 的必经之路上(比如一个归一化函数),而不是每个 push 点都记得调:

def push_safe(ti, key, value):
"""XCom push 前清洗 NaN/Inf,杜绝序列化崩溃。"""
ti.xcom_push(key=key, value=json_safe_value(value))

# 任务内
push_safe(ti, "sql_metadata", result)
push_safe(ti, "processor_output", processor_result)

3. 补上「无声失败」的可观测性

光修序列化还不够——异常发生在 catch 外、应用 logs 表不记录这个缺口要一起补。给任务挂一个失败装饰器,顶层异常先落库再 re-raise:

import functools
import logging

logger = logging.getLogger(__name__)

def log_task_failure(fn):
@functools.wraps(fn)
def wrapper(*args, **kwargs):
try:
return fn(*args, **kwargs)
except Exception:
logger.error("task %s failed", fn.__name__, exc_info=True)
# 这里把 traceback 写进应用自建 logs 表
raise
return wrapper

@log_task_failure
def analyze_results(**context):
...

这样即使以后再出现 catch 外的异常,应用 logs 表也能留下 error 行,不再「无声失败」。

修完后重跑同一份 conf:DAG 全绿、落库 success,原本 NaNad_roi 在库里落成 null,下游正常。

同一条 Airflow 分析流水线上,让数据悄悄出问题的坑不止这一个——PostgresHook 多语句 SQL 静默丢结果是另一个典型案例。

注意事项

注意事项

  • json.dumps 默认 allow_nan=True 会埋雷:它会偷偷写出裸 NaN / Infinity 这个非法 JSON,当下游用严格 parser(如 Airflow XCom、JS 的 JSON.parse)读取时才崩。永远在序列化跨进程边界的数据时显式 allow_nan=False 提前暴露问题。
  • ±Infinity 同样踩雷float('inf') / float('-inf')NaN 一样被 JSON 标准排除,json_safe_value 要一并处理。
  • XCom 不止 JSON 一种 serializer:Airflow 也支持二进制对象序列化,能存任意 Python 对象,但这种 XCom 不可读、不跨版本、且反序列化任意对象有安全风险,生产环境坚持用 JSON 并把数据清洗干净。
  • 排查心法:当 logs 表 trace 中断且无 error 行时,直接去 Airflow 任务日志(容器内 /opt/airflow/logs/dag_id=.../task_id=.../)找 traceback——「应用层无日志」不等于「没出错」。

常见问题

Airflow 报 Out of range float values are not JSON compliant 怎么解决?

这是 XCom 用 json.dumps(allow_nan=False) 序列化时遇到了 NaN / Infinity,而 JSON 标准不含这两种值。根因通常是 pandas 把 SQL NULL 读成了 float('nan'),跟着数据流进了 xcom_push。解法是在 push 前递归把 NaN / ±Inf 转成 None(JSON null),用一个 json_safe_value 纯函数统一兜底即可。

为什么 Airflow 任务失败但自建 logs 表没有错误记录?

如果异常发生在 XCom 序列化阶段、且位于任务函数的 try/except 之外,错误只会冒泡给 Airflow 调度器、写进 Airflow 任务日志(容器内 /opt/airflow/logs/),应用层自建的 logs 表的 catch 拿不到,于是表现为「无声失败」。排查这类情况要直接看 Airflow task log 的 traceback,别只盯应用日志。

Airflow XCom 能不能直接存 pandas 的 NaN?

不能。XCom 默认走 JSON 序列化,而 JSON 标准只有有限数字,没有 NaN / Infinity。正确做法是 push 前把 NaN 转成 None(对应 JSON null)。换成二进制对象序列化虽能绕过类型限制,但结果不可读、不跨进程/版本、还有反序列化安全风险,生产环境不推荐。


CCLEE

独立开发者,24年电商行业实战经验,专注将AI能力落地于真实商业场景。

合作咨询

Airflow PostgresHook 多语句 SQL 静默丢结果?按分号切分逐条执行

· 阅读需 7 分钟

在 Airflow DAG 把 .sql 模板文件整段读出后传给 PostgresHook.get_pandas_df() 时,前置 SELECT 的结果被静默丢弃——DAG 报「SQL 查询无结果」,但把同一段 SQL 复制到 psql 又能正常返回数据。

在开发 AI运营 时遇到此问题——基于大语言模型的智能分析流水线,Airflow DAG 从 SQL 模板文件读取多查询报表模板并执行。

TL;DR

PostgresHook.get_pandas_df(sql) 内部走 pandas.read_sql(sql, conn)psycopg2 cursor.execute(sql)。当 sql 是含多条 ; 分隔 SELECT 的单字符串时,DBAPI 只暴露最后一个结果集的游标,前置查询结果被静默丢弃且不报错。修复方法:按顶层分号切分成 list[str],逐条 get_pandas_df 收集结果,或直接传 list 让 DbApiHook 顺序执行。

问题现象

DAG 任务执行 shop_monthly_overview.sql 报「SQL 查询无结果」:

sql_count = 1   ← 模板里明明写了 4 条查询
result = "❌ SQL 查询无结果"

但同一份 SQL 复制到 psql 直连同库同参数,4 条 SELECT 都有数据。

复现实验

在 Airflow 容器内直接验证 get_pandas_df 对多语句的行为:

from airflow.providers.postgres.hooks.postgres import PostgresHook

hook = PostgresHook(postgres_conn_id="postgres_default")

# 三条 SELECT 串成单字符串
sql = "SELECT 1 AS a; SELECT 2 AS b; SELECT 99 AS c WHERE 1=0;"

df = hook.get_pandas_df(sql)
print(df.columns.tolist()) # ['c'] ← 只拿到末条的列
print(df) # Empty ← 末条本身 0 行

预期应得到三条结果,实际只拿到末条(SELECT 99 ... WHERE 1=0,0 行),前两条完全消失,没有任何报错或警告。

根因

调用链是 PostgresHook.get_pandas_dfDbApiHook.get_pandas_dfpandas.io.sql.read_sqlpsycopg2 cursor.execute(sql)

DBAPI 协议(PEP 249)允许 execute 接受含多条 ; 分隔语句的字符串,PostgreSQL 服务端会依次执行全部语句,但游标只暴露最后一个结果集——这是 PostgreSQL wire protocol 的固有行为,不是 Airflow 或 pandas 的 bug。

┌────────────────────────────────────────────────────┐
│ SELECT 1; ← 执行,结果集 1 立即被丢弃 │
│ SELECT 2; ← 执行,结果集 2 立即被丢弃 │
│ SELECT 99 WHERE 1=0; ← 执行,结果集 3 暴露给游标 │
└────────────────────────────────────────────────────┘

pandas.read_sql 只 fetch 到结果集 3

源头是 task_execute_sql.sql 文件整段读出后当成一条字符串传进 get_pandas_df

# ❌ 问题代码
sql_text = open(sql_path).read() # 含 4 条 SELECT 的整段
df = pg_hook.get_pandas_df(sql_text) # 只拿到末条结果

为什么 psql 能正常返回?因为 psql 前端会主动遍历所有结果集并依次打印,而 DBAPI 游标不会。

解决方案

方案 A(推荐):按顶层分号切分后逐条执行

适合 .sql 模板文件场景——文件含注释、引号、多查询,需要稳健的切分。

def split_sql_statements(sql: str) -> list:
"""
按顶层分号切分 SQL,正确处理:
- 单引号字符串内的分号('a;b' 不切)
- SQL 标准 '' 转义('it''s' 不切)
- -- 行注释内的分号(-- note; not split 不切)
"""
statements = []
buf = []
i, n = 0, len(sql)
in_quote = False

while i < n:
ch = sql[i]

# 在单引号字符串内
if in_quote:
buf.append(ch)
if ch == "'":
# '' = 字面量单引号,不结束字符串
if i + 1 < n and sql[i + 1] == "'":
buf.append(sql[i + 1])
i += 2
continue
in_quote = False
i += 1
continue

# 顶层
if ch == "'":
in_quote = True
buf.append(ch)
elif ch == '-' and i + 1 < n and sql[i + 1] == '-':
# 行注释,原样吞到行尾(注释里的 ; 不切分)
while i < n and sql[i] != '\n':
buf.append(sql[i])
i += 1
continue
elif ch == ';':
stmt = ''.join(buf).strip()
if stmt:
statements.append(stmt)
buf = []
i += 1
continue
else:
buf.append(ch)
i += 1

# 末尾无分号的残留块
stmt = ''.join(buf).strip()
if stmt:
statements.append(stmt)

return statements


# 调用方
sql_text = open(sql_path).read()
statements = split_sql_statements(sql_text)

# 逐条执行,收集所有结果
all_results = []
for idx, stmt in enumerate(statements, start=1):
df = pg_hook.get_pandas_df(stmt)
if not df.empty:
all_results.append({
"sql_index": idx,
"sql": stmt,
"data": df.to_dict("records"),
"columns": df.columns.tolist(),
"row_count": len(df),
})

方案 B:直接传 list 给 DbApiHook

Airflow DbApiHook.runget_records 接受 list[str] 参数会按顺序执行——但 get_pandas_df 在 list 模式下的返回行为各 provider 实现不一致,生产环境建议用方案 A 自己控制。

为什么不用 sqlparse.split

社区答案常推荐 sqlparse.split(sqlparse.format(sql, strip_comments=True)),但 strip_comments=True丢掉注释,如果你的下游 processor 依赖注释中的元信息(如 -- dimension: shop),就会丢失上下文。手写切分器保留注释原文,行为可控。

注意事项

注意事项

  • 不要用 sql.split(';') 简单切分——会误切 WHERE name = 'a;b' 这类引号内的分号,以及 -- 注释; 行注释里的分号
  • split_sql_statements 只处理单引号字符串和 -- 行注释;如果你的 SQL 用 /* 块注释 */ 或 dollar-quoted string($$...$$),需要扩展切分器
  • 修复后下游 processor 的 sql_index 语义会变(1-based 顺序索引),同步检查所有 df.iloc[sql_index] 类用法
  • 如果你的 SQL 是程序生成而非文件读取,更安全的做法是生成时就用 list,避免后续切分
  • 顺带提一个相邻的坑:如果你在 Drizzle ORM 里也遇到过 SQL 表达式被静默参数化的问题,可以看 Drizzle sql 模板混用参数化值与 SQL 表达式——同样是「框架替你做了你没预期到的转换」类陷阱

常见问题

Airflow PostgresHook 怎么执行多条 SQL 语句?

list[str] 而不是单条字符串。DbApiHook.get_pandas_dfrun 接受 sql 参数为 list 时按顺序逐条执行;单字符串含多条分号分隔语句时 psycopg2 只返回末条结果集。生产环境推荐自己切分后逐条调用,方便控制结果聚合和 sql_index 索引。

为什么 get_pandas_df 多语句 SQL 只返回最后一条结果?

pandas.io.sql.read_sqlpsycopg2 cursor.execute 执行整段字符串,DBAPI 协议对多语句只暴露最后一个结果集的游标,前置 SELECT 结果被服务端立即丢弃,不报错也不警告。psql 能正常返回是因为 psql 前端会主动遍历所有结果集,DBAPI 游标不会。

怎么安全地按分号切分含注释和引号的 SQL?

逐字符扫描,仅在「非单引号内、非 -- 行注释内」的顶层分号处切分。单引号字面量用 SQL 标准 '' 转义;不要用 str.split(';'),会误切注释和字符串里的分号。如果用 sqlparse.split,注意 strip_comments=True 会丢掉注释原文。


CCLEE

独立开发者,24年电商行业实战经验,专注将AI能力落地于真实商业场景。

合作咨询