跳到主要内容

37 篇博文 含有标签「Bug修复」

查看所有标签

SQL 聚合后指标暴涨几十倍?别直接 SUM 比率列

· 阅读需 6 分钟

在把按天返回的广告数据聚合成周报时,PPC、CPM、ROI 等比率指标暴涨几十倍——PPC 从日粒度实测的 4.70 变成了周报表里的 111.39。

在开发 AI运营 时遇到此问题——基于大语言模型的智能分析,自动洞察市场趋势、用户行为、销售数据,提供精准运营策略。广告数据接口按天返回明细行,入库前要聚合成周粒度供报表消费。聚合上线后首查周报:PPC 111.39(日实测 4.70)、CPM 1585.9(日实测 68.6)、ROI 169.98(日实测 1.82)——13 个比率/均值列全部失真。

TL;DR​

聚合查询里对 CTR、PPC、ROI 这类比率/均值列直接 SUM,得到的是「每天比率的总和」而不是均值,数值按聚合天数成倍放大。两条规则:聚合时只加总可加列(展现、点击、花费等总量),跳过所有比率列;聚合后用总量重算比率(花费÷点击、点击÷展现)。如果手里只有比率没有总量,用分母作权重做加权平均,绝不简单平均。

问题现象​

聚合代码对数值列「全部 SUM」,比率列被静默带进去:

-- 错误写法:所有列都 SUM
SELECT
campaign_id,
date_trunc('week', day) AS week,
SUM(clicks) AS clicks,
SUM(impressions) AS impressions,
SUM(ppc) AS ppc, -- 7 天日 PPC 之和!
SUM(ctr) AS ctr, -- 7 天日 CTR 之和!
SUM(roi) AS roi -- 7 天日 ROI 之和!
FROM daily_ad_report
GROUP BY campaign_id, date_trunc('week', day);

实测对比(某计划单周):

指标日粒度实测SUM 周值放大倍数
ppc4.70111.39~24×
cpm68.61585.9~23×
roi1.82169.98~93×

没有任何报错,数据照常入库,报表照常渲染——只有把周报和日明细并排对比,才能发现数值差了一到两个数量级。

根因​

比率与均值是不可加的派生量。 ppc = cost / clicks 的分母每天不同,把 7 天的日 PPC 直接相加,数学上得到的是「7 个相对值的和」,没有任何业务含义。CTR、ROI 同理。

「全部数值列求和」是静默陷阱。 聚合代码通常按列循环统一处理,比率列混在其中不报错、不告警,只是结果悄悄失真。列越多、比率列占比越高,越难肉眼发现。

ROI 放大 93 倍反而更具迷惑性。 它看起来像「投放效果极好」,如果下游直接消费周报做预算决策,错误的数字会一路传到运营动作里。这次事故里,下游还有只读消费方直接读这张周表——表值修对之前,所有消费方都在读错数据。

解决方案​

第一步:聚合只保留可加列​

CREATE VIEW weekly_ad_totals AS
SELECT
campaign_id,
date_trunc('week', day) AS week,
SUM(impressions) AS impressions,
SUM(clicks) AS clicks,
SUM(cost) AS cost,
SUM(gmv) AS gmv
FROM daily_ad_report
GROUP BY campaign_id, date_trunc('week', day);

可加列的特征:它们是「计数/总量」(展现、点击、花费、订单数),跨时间区间相加仍有意义。

第二步:聚合后用总量统一重算比率​

SELECT
campaign_id,
week,
impressions,
clicks,
cost,
CASE WHEN clicks > 0
THEN cost / NULLIF(clicks, 0)::numeric
ELSE 0 END AS ppc,
CASE WHEN impressions > 0
THEN clicks::numeric / NULLIF(impressions, 0)
ELSE 0 END AS ctr,
CASE WHEN cost > 0
THEN (gmv - cost)::numeric / NULLIF(cost, 0)
ELSE 0 END AS roi
FROM weekly_ad_totals;

两个细节:PostgreSQL 整数除法会截断,除法前先 ::numeric;分母为 0 统一返回 0,保持与明细层口径一致。

第三步:只有比率、拿不到总量时用加权平均​

-- 用展现量加权聚合日 CTR(展开式:SUM(ctr × impressions) / SUM(impressions))
SELECT
date_trunc('week', day) AS week,
SUM(clicks)::numeric / NULLIF(SUM(impressions), 0) AS ctr_weighted
FROM daily_ad_report
GROUP BY date_trunc('week', day);

加权平均的本质就是「还原分子分母再相除」——只要还拿得到权重列,就永远优先于简单平均。

改完后周报 13 个比率列全部与日明细实测一致,历史脏数据用同一套公式回填,下游只读消费方不改一行代码自动变对。

注意事项

「所有数值列求和」的通用聚合代码是这类事故的源头:维护一份可加列白名单,比率/均值列显式排除,新增指标列时先回答「它跨天相加还有意义吗」。

多粒度报表(周报、月报)从同一张日表派生时,把「用总量重算比率」收敛成一个函数/视图,别在每份报表 SQL 里复制公式——口径漂移往往从复制开始。

修完聚合逻辑记得回填历史数据:聚合错误通常已持续多个周期,只改代码不回填,报表会继续展示旧错值。

常见问题​

百分比可以直接 SUM 吗?​

不能。百分比/比率是相对值,各行分母不同,直接 SUM 得到的是 N 个相对值之和,按聚合天数放大且无业务含义。正确做法是聚合时跳过比率列,聚合后用总量重算(点击÷展现、花费÷点击);只有比率没有总量时,用分母作权重做加权平均。

百分比能加起来求平均吗?​

只有分母相同时才可以。分母不同的百分比直接平均等于不加权平均,结果偏向分母小的项(小流量日的极端比率会被放大)。正确做法是分别加总分子和分母再相除,数学上等价于以分母为权重的加权平均。

CCLEE

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

合作咨询

明细表有行、汇总查无此条?跨 SQL 判定集粒度分叉

· 阅读需 6 分钟

排查一个数据看板线索:某广告 offer 在「诊断明细」里被判疑似停投、并附了再投资建议,点进「建议投放」清单却查无此 offer——明细页指着一张清单,清单里没有这行。

在开发 AI运营 时遇到此问题——基于大语言模型的智能分析平台,自动洞察市场趋势、用户行为与销售数据;诊断明细与建议清单分别来自管道里两个相邻查询的产出。

TL;DR​

两个查询对同一个业务谓词(「疑似停投」)用了不同的聚合粒度:明细查询按 计划×offer 粒度判定——任一计划末周无消耗即判疑似;候选清单按 offer 整体粒度判定——挂着的全部计划末周均无消耗才入选。某 offer 在其中一个计划末周仍有消耗,于是「明细判疑似、汇总不入选」,呈现层指向的明细悬空。修法三件套:口径统一收紧到整体粒度、两侧判定 CTE 逐字同源、引用对方行集前校验包含关系。

问题现象​

两个查询各自定义「疑似停投」:

-- 查询④ 诊断明细:pair 粒度(计划 × offer)
WITH consumption_paused_q4 AS (
SELECT offer_id, campaign_id
FROM spend_weekly
GROUP BY offer_id, campaign_id -- 每个计划单独看
HAVING SUM(spend) FILTER (WHERE is_last_week) = 0
)

-- 查询⑤ 建议投放候选:offer 粒度(跨全部计划汇总)
WITH suspected_paused_q5 AS (
SELECT offer_id
FROM spend_weekly
GROUP BY offer_id -- 整个 offer 一起看
HAVING SUM(spend) FILTER (WHERE is_last_week) = 0
)

出问题的 offer 挂在多个计划下,其中一个计划末周仍有消耗(3.09):

查询④(pair 粒度):  计划A 末周消耗 0     → 判疑似 ✓,并给出建议
查询⑤(offer 粒度): 跨计划汇总末周消耗 3.09 → 不入选 ✗
呈现层:明细页说「见建议投放表」,建议投放表里没有它

任务不报错、行数对得上大数,悬空只有顺着某条明细点过去才会撞见。

根因​

「同源」契约有覆盖盲区。管道规范要求相邻查询的特征列 CTE 逐字同源——这条契约被严格遵守了;但它只约束特征列,判定集(WHERE 之前的那个业务谓词)不在契约范围内。「疑似停投」在两个查询里被独立实现了两次,粒度不同:pair 级判定对「单个计划内无消耗」敏感,offer 级判定对「所有计划都无消耗」敏感。同一个 offer 两种结论,数学上必然存在交叉带。

跨查询引用把分叉放大成了悬空:呈现层拿查询④的行当明细、查询⑤的表做入口,却没人校验过「④的行集 ⊆ ⑤的行集」。聚合口径不一致是数据仓库最经典的一致性陷阱之一——此前写过的 SUM 比率列导致聚合后指标暴涨是它的另一种形态:聚合发生在了错误的层级上。

解决方案​

步骤 1:先定口径,再写 SQL​

业务问题只有一个答案:「这个 offer 还投不投」是 offer 级决策,判定就该收紧到整体粒度——挂着的全部计划末周均无消耗、且无任何运营标注行,才判疑似。口径变更升 rule_version,可追溯。

步骤 2:判定集 CTE 逐字同源​

把判定 CTE 抽成同一段文本,两个查询直接引用;契约同步升级:逐字同源覆盖判定集(含粒度),不止特征列。从此同一谓词只有一处定义,改口径只改一处。

步骤 3:跨查询引用前,校验行集包含关系​

把「标记集 ⊆ 对方行集」做成固定校验(可入库为测试):

-- 悬空检测:④ 判了疑似、⑤ 却查无此行
SELECT q4.offer_id
FROM consumption_paused_q4 q4
LEFT JOIN suspected_paused_q5 q5 USING (offer_id)
WHERE q5.offer_id IS NULL;

这条查询返回 0 行,呈现层才有资格把两个产出拼在同一张页面上。

步骤 4:重放验证后上线​

新口径对历史快照重放:逐行核对翻转方向(疑似↔正常)全部正确、零误伤、零悬空,再重跑生产验证通过,才完成收口。

注意事项

  • 「同源」契约的覆盖面必须包含判定集粒度,特征列同源救不了谓词两处定义。
  • 同一业务谓词(判停投、判爆款、判流失…)全管道只允许一处定义;发现第二处实现即是事故预备役。
  • 新增跨查询/跨模块引用(A 的输出行指向 B 的产出表)前,先跑行集包含校验(标记集 ⊆ 对方行集),不要等用户点到悬空链接。
  • 口径变更必须升版本并对历史数据重放,只看「新数据跑通」会漏掉存量结论的翻转风险。

常见问题​

为什么两条 SQL 对同一份数据给出不同结果?​

常见三个差异:聚合粒度(GROUP BY 维度不同——本例的计划×offer 与 offer 整体)、过滤口径(WHERE/HAVING 条件不同)、取数时点(查询时间不同)。粒度分叉最隐蔽:两条 SQL 各自都对,结论却可以相反。

SQL 的 HAVING 和 WHERE 在聚合判定里怎么选?​

WHERE 在分组前过滤行,HAVING 在分组后过滤组。但选对关键字之前先选对粒度——「以什么为一组」决定了判定的灵敏度:粒度越细越容易命中(单组满足即判),越粗越保守(全部满足才判)。两条 SQL 粒度不同,判定结论就可能相反。

如何避免报表之间的数据口径不一致?​

同一业务谓词只在一处定义,判定 CTE 多查询逐字同源;跨查询引用对方行集前跑包含关系校验(A ⊆ B);口径变更升版本号并对历史数据重放验证。一致性不靠约定俗成,靠契约加校验。

CCLEE

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

合作咨询

Python json.dumps 序列化 set 后 in 判断静默失效?default=str 的隐藏陷阱

· 阅读需 7 分钟

在用 json.dumps(data, default=str) 把一个含 Python set 的字典持久化、再回读用 in 判断成员时,结果静默出错——没有任何报错,但 in 判断全乱。

在开发 AI运营 时遇到此问题——基于大语言模型的智能分析,自动洞察市场趋势、用户行为、销售数据,提供精准运营策略。某次分析模板的决策重放(replay)功能里,需要把「缺失月份集合」序列化进快照、回放时再读出来判断某月是否缺失。结果重放后,本应判定为「缺失」的月份被误判为「不缺失」,而整个链路没有任何异常抛出。

TL;DR​

default=str 不是万能兜底。它会把 set 交给 str(),在 JSON 里存成 "{1, 2}" 这样的字面量字符串而非数组;回读后类型已不可逆,对它做 in 判断会退化成子串匹配,静默返回错误结果。涉及 set 时,正确做法是序列化前转 list、读取时 set() 重建。

问题现象​

下面这段代码完整复现了静默出错的过程:

import json

# 一个含 set 的字典——比如"需要补数据的缺失月份"
data = {"missing_months": {"3", "5", "12"}}

# 用 default=str 兜底序列化(常见的"别让它报错"写法)
serialized = json.dumps(data, default=str)
print(serialized)
# {"missing_months": "{'3', '5', '12'}"} ← 变成了字符串,不是数组!

# 回读
back = json.loads(serialized)
value = back["missing_months"]
print(type(value)) # <class 'str'> ← 已经不是 set 了

# 静默 bug:本想判断某月份是否在"缺失集合"里
print("1" in value) # True ← 1 根本不在 {3,5,12},但 "1" 是 "12" 的子串!
print("3" in value) # True ← 碰巧对
print("9" in value) # False

"1" in value 返回 True,但原集合 {"3", "5", "12"} 根本不含 "1"。没有异常、没有警告,判断结果就这样悄悄错了。这种 bug 在依赖判断结果做分支(如「这个月缺数据吗?缺则补采」)的链路里尤其致命。

根因​

分三层看:

第一层:set 本就不可 JSON 序列化。 JSON 只有 array(对应 list)和 object,没有集合类型。直接 json.dumps({"x": {1, 2}}) 会抛 TypeError: Object of type set is not JSON serializable。

第二层:default=str 把报错变成了静默污染。 json.dumps 的 default 参数在遇到无法序列化的对象时被调用,期望返回一个可序列化的值。str 作为 default 时,会把对象交给 str()——set 就被转成了它的 Python 字面量表示 {'3', '5', '12'},作为字符串存进 JSON:

>>> json.dumps({"m": {"3", "5", "12"}}, default=str)
'{"m": "{\'3\', \'5\', \'12\'}"}'

报错消失了,代价是类型从 set 变成了 str,且这个过程不会给你任何提示。

第三层:in 对 str 和 set 语义不同。 这是静默 bug 的核心。对 set/list,x in s 是成员判断;对 str,x in s 退化成子串匹配。回读后的值是字符串 "{'3', '5', '12'}",于是 "1" in "{'3', '5', '12'}" 判断的是字符 "1" 是否作为子串出现——而 "12" 里恰好有 "1",所以返回 True。

这和 Airflow PostgresHook 多语句 SQL 静默丢结果 是同一类陷阱:最危险的 bug 不是抛异常,而是「静默地给错结果」,因为没有任何信号提醒你去查。

解决方案​

核心原则:JSON 里只存标准类型,集合语义在读取端重建。

方案一:序列化前显式转 list(推荐)​

最直接、最可控——明确知道哪里有 set,就地转成 list:

import json

# 序列化前:set → list(标准 JSON 数组)
data = {"missing_months": list({"3", "5", "12"})}
serialized = json.dumps(data)
print(serialized)
# {"missing_months": ["3", "5", "12"]} ← 正确的 JSON 数组

# 回读后重建 set
back = json.loads(serialized)
months = set(back["missing_months"])
print("1" in months) # False ✓
print("3" in months) # True ✓

序列化结果是一个干净的 JSON 数组,跨语言、可读、可还原。

方案二:自定义 default 函数(数据来源复杂时)​

如果数据结构较深、不确定哪里混入了 set,用一个专门处理集合类型的 default 函数,既不丢失语义,又能兜底其他非标准类型:

import json

def safe_default(obj):
# 集合类型 → list,保留为标准 JSON 数组
if isinstance(obj, (set, frozenset)):
return sorted(obj) # 排序让输出稳定可预测
# 其他无法序列化的类型再退回 str,但要清楚这会丢类型
return str(obj)

data = {"missing_months": {"3", "5", "12"}, "created_at": some_datetime}
serialized = json.dumps(data, default=safe_default)
# {"missing_months": ["3", "5", "12"], "created_at": "..."}

back = json.loads(serialized)
months = set(back["missing_months"])
print("1" in months) # False ✓

相比无脑 default=str,这个函数把「需要保真的类型」(集合)单独处理,只有真正无法表示的类型才退回 str,把静默风险控制到最小。

注意事项

  • default=str 是「静默」而非「安全」:它消除了报错,却把 set/tuple/datetime/自定义对象全部压扁成字符串,类型信息不可逆。回读后所有依赖原类型的运算(in 成员判断、算术、比较)都可能出错。
  • tuple 也有类似问题:str((1, 2)) 是 "(1, 2)",同样会让回读后的 in 退化成子串匹配。处理集合类容器的思路一致:序列化成 list。
  • 跨进程/跨语言是试金石:如果这份 JSON 会被 Node.js、Go 等读取,default=str 产出的 "{1, 2}" 在那边只是一个普通字符串,连 Python 字面量都不是,还原几乎不可能。坚持存标准 JSON 类型才能保证可移植。
  • 优先在源头转换:与其事后用 default 兜底,不如在构造数据结构时就用 list 存集合语义,从根上避免 set 进入序列化管线。

常见问题​

Python set 怎么转 json?​

set 不是 JSON 原生类型,直接 json.dumps 会抛 TypeError。正确做法是序列化前用 list(set) 转成列表,存成标准 JSON 数组;读取时再 set(back["key"]) 重建。这样既不报错,又能完整还原集合语义,跨语言也兼容。

json.dumps 报 Object of type set is not JSON serializable 怎么解决?​

根因是 set 不可 JSON 序列化。最稳妥的解法是序列化前把 set 转成 list;也可以传一个 default 函数,在里面对 isinstance(obj, (set, frozenset)) 返回 list(obj)。要避免用 default=str 兜底——它虽不报错,却把 set 存成了字符串,回读后类型无法还原。

为什么 default=str 序列化 set 后 in 判断结果错了?​

default=str 会把 set 交给 str(),变成字面量字符串 '{1, 2}' 存进 JSON。回读后值类型是 str 而非 set,x in s 就从「成员判断」退化成「子串匹配」——比如 "1" in "{'3','5','12'}" 因 "12" 含字符 "1" 而返回 True,但原集合并不含 "1"。解法是序列化 list、读取时 set() 重建。

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_dag、dag_code、dag_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 进元数据表(dag、serialized_dag、dag_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 ✅

reserialize 后 dag / 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_run → serialized_dag → dag_version → dag → 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 整列 NULL,NaN 才第一次大规模进入 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,原本 NaN 的 ad_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能力落地于真实商业场景。

合作咨询

用 Zod 校验 LLM 输出却静默失败?别用 .strict()

· 阅读需 7 分钟

用 Zod 给 LLM 的 tool_call / function call 输出做校验,模型偶尔多吐一个字段——比如你只定义了 amount / category,它顺手填了个 note——整条校验就挂了,动作被静默丢弃,用户只收到一句「没识别到」,实则是一条 tool_call 被 whole-reject。

在开发 Life 记账助手 时遇到此问题——自然语言记账健康助手,说人话就能记,AI 自动抽取金额、类目、账户;用户说「删掉昨天那杯咖啡」时,模型在 delete locator 里多塞了个 note: "咖啡" 想按备注定位。

TL;DR​

Zod 的 .strict() 等于「对象不许有任何未知键,多一个就报错」。这套约束适合校验你完全控制的客户端,但 LLM 的 function call 输出是模型生成的、本质不可控——它会填入自己「以为该有」的字段,尤其当多个 tool 共用相似 schema 时。一个无关字段就把整条 tool_call 杀掉,校验返回 null,动作静默丢失。解法:去掉 .strict(),用 Zod 默认的 strip(静默删除未知键)容错,配合 safeParse 兜底。

问题现象​

delete/update 的 locator schema 定义了几个已知字段,但用 .strict() 收紧:

import { z } from "zod";

// ❌ 危险:带 .strict()
const LocatorSchema = z.object({
date: z.string().optional(),
category: z.string().optional(),
noteContains: z.string().optional(),
}).strict(); // ← 未知键一律报错

// 解析 LLM 的 tool_call 参数
function parseToolCall(raw: unknown) {
const parsed = LocatorSchema.safeParse(raw);
if (!parsed.success) {
return null; // ← 整条 tool_call 被丢弃
}
return parsed.data;
}

用户说「删掉昨天那杯咖啡」,模型给出(合理但多了一个字段的)输出:

{
"date": "昨天",
"noteContains": "咖啡",
"note": "咖啡"
}

模型同时填了 noteContains(schema 内)和 note(schema 外,它以为该有)。.strict() 对 note 这个未知键直接判失败,parseToolCall 返回 null,这条 delete 动作被静默丢弃——用户收到「没识别到」,实际是被整条拒绝。

根因​

.strict() 改的是 Zod 对未知键的策略,而 LLM 输出天然会带未知键。

Zod z.object() 对未知键有三种策略:

写法未知键行为适合场景
默认(strip)静默删除LLM 输出、宽松外部输入
.strict()报错(unknown key)你完全控制的客户端 API
.passthrough()保留原样下游要用未知键时

.strict() 的设计意图是「契约严格性」——服务端定义了什么字段,客户端就该只给什么,多给即违约。这套逻辑对传统 API 成立,因为客户端是开发者写的、可以要求守约。

但 LLM function calling 颠覆了这个前提:

  1. 输出来自模型生成,不是开发者写的客户端。 模型基于 schema 的 description 和示例猜测该填什么,跨域复用的 schema(比如 locator 在 budget / mood / todo 多个域共用)更会让它混淆,填入「它以为该有」的字段。
  2. 字段填错是常态,不是异常。 模型偶尔多吐一个 note、少吐一个可选字段,是 LLM 应用的预期行为,不该用「整条失败」来惩罚。
  3. 失败被静默吞掉。 safeParse 失败后返回 null,上游拿到 null 只能笼统地说「没识别到」,真正的根因(一个 unknown key)藏在 parsed.error 里没人看。
LLM 输出 { date, noteContains, note }
│
▼
.strict() 遇到未知键 note
│
▼
safeParse → { success: false }
│
▼
parseToolCall 返回 null(动作丢弃)
│
▼
用户收到「没识别到」(实则 whole-reject)

解决方案​

1. 去掉 .strict(),用默认 strip 容错​

// ✅ 推荐:不带 .strict(),Zod 默认 strip 未知键(静默删除)
const LocatorSchema = z.object({
date: z.string().optional(),
category: z.string().optional(),
noteContains: z.string().optional(),
});
// 模型多吐的 note 会被静默删掉,已知字段照常解析

去掉 .strict() 后,「删掉昨天那杯咖啡」正常解析为 { date, noteContains },多余的 note 被 strip,delete 动作正确执行。

2. 如果未知键本身有用,用 .passthrough() 显式保留​

当模型多吐的字段其实承载了你想用的语义(比如它填 note 是想表达「按备注定位」),别丢,保留下来再决定怎么消费:

const LocatorSchema = z.object({
date: z.string().optional(),
category: z.string().optional(),
noteContains: z.string().optional(),
}).passthrough(); // 保留未知键,parsed.data.note 仍可读

更好的做法是把它收编成已知字段——发现模型反复填某个未知键,说明 schema 缺了这个能力位,补上(比如这里的 noteContains 就是收编「按备注定位」需求后新增的)。

3. 失败要可观测,别静默返 null​

无论哪种策略,safeParse 失败时都要把具体的 error 落日志,而不是吞成 null:

function parseToolCall(raw: unknown) {
const parsed = LocatorSchema.safeParse(raw);
if (!parsed.success) {
// 把 Zod 的具体报错(哪个键、什么问题)落日志,便于定位
logger.warn(
{ raw, issues: parsed.error.issues },
"locator parse failed"
);
return null;
}
return parsed.data;
}

这样真出问题时,日志里有完整的 issues(含 unknown key 的路径),而不是一句无从下手的「没识别到」。

修完后,「删掉昨天那杯咖啡」→ delete_record { locator: { noteContains: "咖啡" } } 正确解析,不再静默丢失。

注意事项​

注意事项

  • .strict() 适合校验「你控制的客户端」,不适合「LLM 生成的输出」。 判断标准:数据来源是你写的代码 → 可以 strict;数据来源是模型生成 → 用默认 strip 或 passthrough。
  • strip 会丢失未知字段。 如果那个字段承载了模型的意图(如例子里的 note),用 .passthrough() 保留,或直接收编成已知字段,别让意图被静默删掉。
  • 永远用 safeParse 而非 parse。 parse 校验失败会抛异常,在 tool 调度链里可能中断整个流程;safeParse 返回结果对象,失败可控。
  • LLM tool schema 设计要给容错空间。 字段尽量 .optional()、description 写清用途、提供 few-shot 示例;预期到模型会「多填/少填」,schema 层就该兜得住。

常见问题​

为什么 Zod .strict() 会让 LLM 输出校验失败?​

.strict() 要求对象不含任何未知键,多一个就抛 unknown key 错误。LLM 的 function call 输出是模型基于 schema description 猜测生成的,常会填入它「以为该有」的字段(尤其跨域复用的 schema),一旦命中未知键,.strict() 就让整条校验失败、整个 tool_call 被丢弃。

用 Zod 校验 LLM function calling 输出应该用 strict 吗?​

不建议。.strict() 适合校验你完全控制的客户端(开发者写的代码可以要求守约),但 LLM 输出不可控、多填少填是常态。去掉 .strict() 用 Zod 默认的 strip(静默删除未知键)容错性更好;如果未知键承载了想用的语义,用 .passthrough() 保留,或直接收编成已知字段。

Zod 默认对未知字段是 strip 还是报错?​

默认是 strip——静默删除未知键、不报错;.strict() 改为遇到未知键就报错;.passthrough() 改为保留未知键原样。校验 LLM 这类不可控输出时推荐默认 strip 或 passthrough,避免 .strict() 因一个无关字段杀掉整条数据。


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_df → DbApiHook.get_pandas_df → pandas.io.sql.read_sql → psycopg2 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.run 和 get_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_df 和 run 接受 sql 参数为 list 时按顺序逐条执行;单字符串含多条分号分隔语句时 psycopg2 只返回末条结果集。生产环境推荐自己切分后逐条调用,方便控制结果聚合和 sql_index 索引。

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

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

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

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


CCLEE

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

合作咨询

Docker Compose 服务重启后起不来?检查 restart 策略

· 阅读需 5 分钟

在 RAG 知识库项目中排查依赖 Milvus 的服务启动失败,以下是完整排查过程。

TL;DR​

宿主机重启(或容器崩溃)后,一组服务没有自动恢复,应用端口无监听、docker ps -a 里容器全是 Exited。根因是 docker-compose.yml 没配 restart 策略(默认 no),容器挂了就永远躺着。解法:给所有生产服务加 restart: always,让基础设施在崩溃或重启后自愈。

Python 任务全标 failed 却不报错?try/except 吞掉了异常

· 阅读需 5 分钟

在 RAG 知识库项目中排查文档同步任务全部标记 failed 的静默故障,以下是完整排查过程。

TL;DR​

重构一个公共方法改了参数签名,但漏改了一个调用方。调用方按旧契约传参抛 TypeError,而这个调用被包在 try/except 里,异常被悄悄吞进 failed 计数——服务不崩溃、日志没有 ERROR,只有计数字段悄悄上涨。这类「静默故障」是最难查的 bug。两个解法:重构签名后 grep 所有调用方同步;except 块必须记日志或重抛,绝不静默吞掉。

Node.js AsyncLocalStorage 在回调里读不到值?EventEmitter 越界丢失上下文

· 阅读需 5 分钟

请求日志中间件在 res.on('finish') 回调里读 AsyncLocalStorage 的 traceId,getStore() 返回 undefined,每条响应日志的 traceId 都是空的。

在为客户开发 电商数据采集工具 时遇到此问题——服务端用 ALS 把每条请求的 traceId 贯穿整条处理链路,但响应日志死活关联不上,排查发现是「晚回调」丢了上下文。

TL;DR​

res.on('finish') 这类 EventEmitter 回调,触发时已经脱离了注册它时的 async context,als.getStore() 自然拿不到请求的 store。最稳的解法是在同步段把值取到闭包变量,回调里直接用闭包值;需要完整 store 时则在回调内 als.run(store, fn) 重建上下文。

问题现象​

一个看起来毫无问题的请求日志中间件:

// middleware/requestLog.js
import { als } from '../utils/als.js';

app.use((req, res, next) => {
res.on('finish', () => {
const store = als.getStore();
logger.info({
traceId: store?.traceId, // 响应日志里这里永远是 undefined
statusCode: res.statusCode,
}, 'request');
});
next();
});

中间件顺序没问题,traceId 在请求处理链路里(路由、业务函数)都读得到,唯独 res.on('finish') 里读不到。更迷惑的是:把 als.getStore() 挪到 next() 之前的同步段,它就有值。

根因​

AsyncLocalStorage 靠 Node 的 async_hooks 把 store 绑定到当前激活的 async context 上,顺着异步调用链往下传。als.run(store, fn) 的语义是:在 fn 执行期间(及其派生的异步任务里),getStore() 都能拿到这个 store。

问题出在 EventEmitter。res.on('finish', cb) 做的事是把 cb 注册成监听器,等响应发送完毕后由 EventEmitter 的事件循环触发。触发 cb 的那个 async context,是 EventEmitter 派发事件时所在的上下文——不是注册它时的请求上下文。而且响应发送通常发生在请求处理链路之后,请求对应的 als.run 作用域可能已经退出。

所以 cb 里 als.getStore() 拿到的是「当前激活上下文」的 store,而那个上下文根本不属于这次请求,结果就是 undefined(或更糟,串到别的上下文)。

凡是「注册时一个上下文、触发时另一个上下文」的回调都有这个坑:res.on('finish')、once、某些 setTimeout/setInterval、chrome.alarms 监听器等等。

解决方案​

按场景给两个模式,按需选。

模式 A(推荐):同步段闭包捕获​

如果你的回调只需要 store 里的某几个值(最常见就是 traceId),最简单也最可靠——在同步段(store 一定存活的时刻)把值取出来存进闭包,回调里直接用闭包变量,彻底不依赖 ALS:

app.use((req, res, next) => {
// 同步段:此时一定在 als.run 作用域内,getStore() 必有值
const traceId = als.getStore()?.traceId;
const start = Date.now();

res.on('finish', () => {
// 回调里用闭包里的 traceId,不再碰 ALS
logger.info({
traceId, // 稳定拿到
statusCode: res.statusCode,
durationMs: Date.now() - start,
}, 'request');
});

next();
});

这一步把「异步上下文是否还活着」这个不确定性,换成了一个确定的闭包引用。回调何时触发都不影响——值已经在闭包里了。

模式 B:als.run 重建上下文​

当回调里要调用一坨内部都依赖 getStore() 的代码(比如 logger 的 mixin、Sentry 的 scope 注入),逐个改成闭包不现实,就在回调入口重建上下文:

res.on('finish', () => {
const traceId = capturedTraceId; // 同步段捕获的值
if (traceId) {
// 在回调内重新建立 ALS 上下文,后续 record() 内部 getStore() 能正常拿到
als.run({ traceId }, () => record(res, start));
} else {
record(res, start);
}
});

als.run(store, fn) 会为 fn 建立一个新的、独立的 async context 并把 store 绑上去,fn 内部及它派生的异步调用都能读到。这比 als.enterWith 更安全——后者改写的是「当前共享上下文」,在并发场景下会串值,那是另一个坑,见 AsyncLocalStorage 并发读到错误的值?enterWith 改用 run 隔离上下文。

注意事项

  • 判断某回调是否会丢上下文,看它是不是「注册和触发分离」。res.on('finish')、once、跨 tick 的 setTimeout 都要警惕;而 await、fetch().then() 这类顺着 async chain 走的则天然继承,不用处理。
  • 模式 A 优先。它把问题降维成一个普通闭包,可读性最好,也不会引入「重建上下文」的隐式行为;只有回调内部有大量依赖 getStore() 的既有代码时,才上模式 B。
  • 别用 als.enterWith 在回调里补救——它在并发下会改写共享父上下文导致串扰,是比「丢上下文」更难查的 bug。

常见问题​

为什么 res.on('finish') 回调里读不到 AsyncLocalStorage 的值?​

res.on('finish', cb) 把 cb 注册为 EventEmitter 监听器,响应发送完毕后才触发。触发时的 async context 是事件派发所在的上下文,不是注册它的请求上下文,请求的 als.run 作用域可能已退出,因此 getStore() 返回 undefined。

怎么让 EventEmitter 回调重新拿到 AsyncLocalStorage 上下文?​

最简单的是在同步段把需要的值取到闭包变量,回调里直接用闭包值;如果回调内部有大量依赖 getStore() 的代码,则在回调入口用 als.run(store, fn) 重建上下文。前者优先,后者用于改造既有逻辑。

CCLEE

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

合作咨询