跳到主要内容

2 篇博文 含有标签「pandas」

查看所有标签

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能力落地于真实商业场景。

合作咨询