跳到主要内容

6 篇博文 含有标签「PostgreSQL」

查看所有标签

PostgreSQL 数据库迁移后留下废弃空表?用 information_schema 审计 schema 残留

· 阅读需 7 分钟

在一次广告数据源从日表切换到周表后,我发现 schema 里还躺着一张只在最早期的 migration 里建过、0 行数据、运行时代码从不引用的空表——还带着过时的字段名和没规范化的中文列。

在开发 AI运营 时遇到此问题——基于大语言模型的智能分析,自动洞察市场趋势、用户行为、销售数据,提供精准运营策略。这次数据源切换后,写入端早已迁到新的周表,旧日表的清理 migration 也补了,唯独一张只在 baseline 里 CREATE 过的月表被遗漏——它既没有对应的新写入,也没有 DROP,就那样潜伏在 schema 里,带着早已废弃的字段定义。

TL;DR

废弃表的典型特征:只在早期/baseline migration 里 CREATE、当前代码 0 引用、常带旧字段或未规范化的列名。批量迁移时它们不会被自动处理,需要主动用 information_schema.tables 列出 schema 全表,再与代码引用比对,定位 orphan 表后写一条 DROP migration 清理——而不是手动 psql 删完就了事。

问题现象

一张典型的废弃表长这样:

  • 0 行数据——业务早已不再写入它;
  • 0 运行时引用——代码里 grep 不到任何 SELECT/INSERT,只剩 migration 文件里的 CREATE
  • 旧字段残留——字段名是上一版命名(如 ad_plan_id/product_id),与当前规范不一致;
  • 未规范化的列名——甚至还有中文列名没来得及改。

它不报错、不影响线上运行,所以从「线上没出问题」的视角完全无感。但它的危害是隐性的:误导后来者以为它仍在用、占用 schema 命名空间、在跨表审计时制造噪音,还可能被某个误判的 SELECT * 意外读到脏数据。

根因

数据库迁移有一个普遍的模式:migration 是「加法」的

一次数据源切换通常这样演进:

  1. 早期 baseline migration CREATE 了一批表(日表、月表);
  2. 业务跑通后,写入端开始依赖这些表;
  3. 需求变化,引入新表(周表),写入端逐步迁移过去;
  4. 旧表的写入停了,补一条 migration DROP 旧日表;
  5. 但月表/其他只在 baseline 建过、从未被写入端直接引用的表,没有对应的 DROP

问题出在第 5 步:迁移注意力集中在「现在用到的表」上——哪些表在写入、哪些 SQL 在查。而「曾经存在、但从未进入主链路」的表既不在写入端、也不在查询端,自然不会触发任何 DROP,于是成了 orphan。这类残留和 Airflow 删除 DAG 后元数据残留 是同一类问题:「删了入口、忘了清结构」,是迁移类问题的高发区。

解决方案

核心流程:列全表 → 比对引用 → 确认空表 → 写 migration DROP → 验证

步骤 1:用 information_schema 列出 schema 下所有基础表

-- 列出某 schema 下所有基础表(排除视图)
SELECT table_name
FROM information_schema.tables
WHERE table_schema = 'your_schema'
AND table_type = 'BASE TABLE'
ORDER BY table_name;

information_schema.tables 是 SQL 标准目录视图,跨 PostgreSQL/MySQL/SQL Server 通用,字段稳定,非常适合写进审计脚本。

步骤 2:grep 代码库确认运行时引用

对每张候选表,在代码库里搜索引用,排除 migration 文件本身

# 搜索运行时代码引用,排除 migrations 目录
grep -rn "ad_product_monthly_stats" src/ --include="*.py" \
| grep -v "migrations/"
# 0 行输出 → 运行时无引用,进入候选

0 引用是判定 orphan 的关键证据。注意一定要排除 migration 目录——baseline 里的 CREATE 不算「引用」。

步骤 3:确认是空表

SELECT count(*) FROM your_schema.ad_product_monthly_stats;
-- 0 → 确认无数据,可安全清理

对有数据的表要格外谨慎:先确认它真的废弃(而非只是近期没写入),有疑问就先做逻辑备份再处理。

步骤 4:写一条 migration DROP(而非手动删)

-- db-migrations/{project}/027_drop_ad_product_monthly_stats.sql
DROP TABLE IF EXISTS your_schema.ad_product_monthly_stats;

务必走 migration 文件:它会被版本控制、在所有环境(开发/预发/生产)一致重放,留下审计轨迹。手动 psql 删一次,换台机器就又长回来了。

步骤 5:验证已删除

SELECT to_regclass('your_schema.ad_product_monthly_stats');
-- 返回 NULL 表示表已不存在

to_regclass() 是验证表是否存在的标准手段,返回 NULL 即确认删除成功。

批量审计:按前缀一次性排查同类遗漏

单张表清掉后,按前缀把同类表全部列出来逐个核对,避免「清了一张、漏了兄弟」:

-- 列出某前缀下所有表,逐个走 步骤2-5
SELECT table_name
FROM information_schema.tables
WHERE table_schema = 'your_schema'
AND table_name LIKE 'ad_%'
ORDER BY table_name;

注意事项

  • DROP 前先备份/快照:生产库删表不可逆。对任何有数据的表,先确认废弃再做逻辑备份(如 CREATE TABLE ... AS SELECT 导出到归档库)。
  • 外键依赖要排查:如果有其他表的外键指向它,DROP TABLE 会失败。确认依赖已解除或有意 CASCADE——但 CASCADE 会连带删除依赖对象,生产环境慎用。
  • 走 migration,不要手动 psql:手动删除只在当前环境生效,迁移文件才能保证多环境一致并留下记录。
  • 用前缀批量审计:一次切换通常涉及一组同前缀的表(如 ad_*),清完一张后用 LIKE 'ad_%' 把兄弟表都过一遍,主动发现同类遗漏。

常见问题

怎么列出 PostgreSQL 数据库中的所有表?

information_schema.tables,过滤 table_schematable_type = 'BASE TABLE',即可列出某 schema 下所有基础表。它比 psql 的 \dt 更适合写进脚本做自动化审计,且是 SQL 标准、跨数据库通用,代码可移植性更好。

PostgreSQL 怎么找出没被使用的废弃表?

information_schema.tables 列出全部表,再与代码库或查询日志的引用做比对,运行时代码 0 引用且无写入的表即为废弃候选;空表可进一步用 SELECT count(*) 确认行数,确认无数据、无外键依赖后再写 migration DROP 清理。

information_schema 和 pg_catalog 有什么区别?

information_schema 是 SQL 标准定义的目录视图,跨 PostgreSQL/MySQL/SQL Server 通用、字段稳定不易变,适合写可移植的审计脚本;pg_catalog 是 PostgreSQL 专有目录,信息更全更细(如精确行数估算、存储细节),但版本间可能调整。做通用 schema 审计优先用 information_schema

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

合作咨询

PostgreSQL ON CONFLICT 报 there is no unique constraint?改唯一键后 INSERT 必须同步

· 阅读需 6 分钟

在为某张表收紧唯一键(移除一个不再区分数据的列)之后,原本正常的 UPSERT 写入立刻批量报错——there is no unique or exclusion constraint matching the ON CONFLICT specification

在开发 AI运营 时遇到此问题——基于大语言模型的智能分析,自动洞察市场趋势、用户行为、销售数据,提供精准运营策略。

TL;DR

PostgreSQL 的 ON CONFLICT (cols) 要求 cols 精确匹配一个已存在的唯一约束或唯一索引(列与顺序都要一致,否则错误码 42P10)。一旦你 ALTER 了唯一键,所有引用它的 INSERT ... ON CONFLICT 必须同步修改;而且 migration 跑完后写入端要立刻部署,中间窗口会持续报错。

问题现象

唯一键改造一上线,定时导入任务全量失败,写库日志只剩这一条:

ERROR: there is no unique or exclusion constraint matching the ON CONFLICT specification
SQL state: 42P10

业务表 0 行写入,但同一张表的其它纯 SELECT 查询完全正常——问题只出在带 ON CONFLICT 的写入路径上。

根因

ON CONFLICT (cols) 里指定的列集叫仲裁器(arbiter)。PostgreSQL 要求它精确匹配表上某个 UNIQUE 约束或唯一索引:

  • 列的集合必须相同;
  • 列的顺序也要相同;
  • 如果是带 WHERE 的部分唯一索引(partial unique index),ON CONFLICT 还要带上相同的 WHERE

找不到匹配项时,PostgreSQL 不知道用哪个索引来判断"冲突",于是抛出 42P10。

典型触发场景是收缩唯一键:原先唯一键含 3 列,你发现其中一列(比如 audience)的 4 个取值对应的指标行 100% 全等、纯属冗余,于是把唯一键降到 2 列。这是对的优化方向,但旧的 INSERT 仍写着 ON CONFLICT (c1, c2, c3),而表上只剩 (c1, c2) 的唯一约束——仲裁器找不到落点,报错。

旧唯一键: UNIQUE (store_id, metric_key, audience)
新唯一键: UNIQUE (store_id, metric_key)

旧 INSERT: ON CONFLICT (store_id, metric_key, audience) ← 找不到匹配

解决方案

下面是最小复现,建表、触发、修复一条龙,可直接在 psql 里跑:

-- 1. 带 3 列唯一键的表
CREATE TABLE daily_metric (
store_id TEXT NOT NULL,
metric_key TEXT NOT NULL,
audience TEXT NOT NULL,
value NUMERIC,
CONSTRAINT daily_metric_unique UNIQUE (store_id, metric_key, audience)
);

-- 2. 旧 UPSERT:ON CONFLICT 含 audience
INSERT INTO daily_metric (store_id, metric_key, audience, value)
VALUES ('s1', 'revenue', 'visitor', 100)
ON CONFLICT (store_id, metric_key, audience)
DO UPDATE SET value = EXCLUDED.value;

-- 3. 收缩唯一键:移除 audience
ALTER TABLE daily_metric
DROP CONSTRAINT daily_metric_unique,
ADD CONSTRAINT daily_metric_unique_new UNIQUE (store_id, metric_key);

-- 4. 再跑第 2 步的 INSERT,立刻报错 ↓
-- ERROR: there is no unique or exclusion constraint matching the ON CONFLICT specification

修复就是把 INSERTON CONFLICT 列同步收缩到 2 列;既然 audience 不再区分数据,写入端干脆把它的值固定为字面量,避免按入参凭空拼出多行:

INSERT INTO daily_metric (store_id, metric_key, audience, value)
VALUES ('s1', 'revenue', 'visitor', 100)
ON CONFLICT (store_id, metric_key) -- ← 同步收缩
DO UPDATE SET value = EXCLUDED.value;

真正容易踩的是部署顺序,不是 SQL 本身:

  1. 先发 migration(DROP 旧约束 + ADD 新约束);
  2. 紧接着发布写入端代码(INSERTON CONFLICT 改为 2 列);
  3. 两步之间不要留间隔——旧代码撞新 schema 必报 42P10,新代码撞旧 schema 同样报 42P10(找不到 2 列的唯一约束)。

如果你用 Drizzle 这类 ORM,ON CONFLICT 的列一旦在 sql 模板里写死,改 schema 时极易漏改——schema 与写入端不同步的代价,在另一篇 Drizzle + PostgreSQL 的坑里也领教过。

注意事项

注意事项

  • 列顺序敏感ON CONFLICT (a, b)UNIQUE (b, a) 不算匹配,顺序必须一致。
  • 部分唯一索引要带 WHERE:若仲裁器是 UNIQUE ... WHERE activeINSERT 里要写成 ON CONFLICT (cols) WHERE active DO ...,否则同样报 42P10。
  • 只想"冲突就跳过":用不带列的 ON CONFLICT DO NOTHING,它不指定仲裁器、无需匹配任何具体索引,能捕获所有冲突。
  • 灰度并存:新老版本写入端可能短暂共存,确保两套代码都能匹配当前 schema,或让 migration 与代码同步上线、不留窗口。

常见问题

PostgreSQL ON CONFLICT 必须有唯一约束吗?

只有指定列时才必须。ON CONFLICT (cols)cols 要精确对应一个已存在的 UNIQUE 约束或唯一索引,否则报 42P10。如果你只想"有任何冲突就跳过"、不关心具体哪个约束,用不带列的 ON CONFLICT DO NOTHING,它不需要匹配特定索引。

PostgreSQL ON CONFLICT 可以指定多个唯一约束吗?

不能。单条 INSERTON CONFLICT 只能指定一个仲裁约束(一个列集,或一个索引名)。表上可以有多个唯一键,但一条语句只能选其一做冲突判定。需要按不同唯一键分别处理时,要么拆成多次写入,要么在应用层先查再决定 INSERT 还是 UPDATE

报 there is no unique or exclusion constraint matching the ON CONFLICT specification 怎么办?

这是错误码 42P10,含义是 ON CONFLICT 指定的列集在表上找不到匹配的唯一索引。按顺序排查:确认存在覆盖这些列的 UNIQUE 约束、列与顺序完全一致、最近改过唯一键后 INSERT 已同步更新;如果用的是带 WHERE 的部分唯一索引,ON CONFLICT 还要补上相同的 WHERE 子句。

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

合作咨询

UPSERT 写入全零?Drizzle sql 模板混用参数化值与 SQL 表达式的坑

· 阅读需 4 分钟

在为客户构建电商数据分析平台时遇到此问题,记录根因与解法。

TL;DR

Drizzle ORM 的 sql 模板标签中,sql.join(values.map(v => sql(v))) 会把所有值参数化传递。如果 values 数组里混入了 SQL 表达式(如 date_trunc('week', '2026-05-17'::date)::date),PostgreSQL 会把它当成普通字符串解析,报 invalid input syntax for type date 错误。SQL 表达式必须用 sql.raw() 或单独写在模板外部

问题现象

电商数据采集流程:Chrome 扩展采集 → CCLHub 转发 → Analytics 写库。现象:

  1. CCLHub 日志显示采集数据正常(uv: 403, payAmt: 19478.47
  2. Analytics 返回 200 成功
  3. 但数据库查询结果全是 0uv: 0, pay_amt: 0.00
-- 数据库实际数据
report_date | uv | pay_amt | reveal_cnt
-------------+-----+----------+------------
2026-05-12 | 392 | 7333.67 | 11879 -- 旧数据正常
2026-05-13 | 0 | 0.00 | 0 -- 新数据全零!

同时 Analytics 错误日志有:

PostgresError: invalid input syntax for type date:
"date_trunc('week', '2026-05-17'::date)::date"

根因

原始代码混用了参数化值和 SQL 表达式:

// ❌ 问题代码
const insertVals: (string | number | null)[] = [
String(shop_id),
String(platform_id),
reportDate,
tenant_id,
`date_trunc('week', '${reportDate}'::date)::date`, // ← SQL 表达式
];

// sql.join 会把所有值参数化,包括 date_trunc 表达式
await db.execute(sql`
INSERT INTO table (..., week_start_date)
VALUES (${sql.join(insertVals.map(v => sql`${v}`), sql`,`)})
...
`);

生成的 SQL:

-- PostgreSQL 收到的 $5 参数值是字面字符串
INSERT INTO table (..., week_start_date)
VALUES ($1, $2, $3, $4, $5, ...)
-- $5 = "date_trunc('week', '2026-05-17'::date)::date" ← 被当字符串!

PostgreSQL 尝试把 "date_trunc('week', '2026-05-17'::date)::date" 解析为 date 类型 → 报错。

为什么数据是 0 而不是报错? 因为同一张表有独立的询盘写入(PARTIAL UPSERT),询盘 INSERT 成功创建了行(看板列默认值 0),日报 UPSERT 失败但没有回滚已存在的行。

解决方案

把 SQL 表达式从参数化数组中分离出来,用 sql.raw() 或直接写在模板中:

// ✅ 修复:参数化值和 SQL 表达式分开
const insertCols = ['shop_id', 'platform_id', 'report_date', 'tenant_id'];
const insertVals: (string | number | null)[] = [
String(shop_id), String(platform_id), reportDate, tenant_id,
];

// 19 个数据列正常参数化
for (const [apiKey, dbCol] of Object.entries(DAILY_COLUMNS)) {
insertCols.push(dbCol);
insertVals.push(row[apiKey] != null ? String(row[apiKey]) : '0');
}

// week_start_date 用 SQL 表达式,不进参数化数组
await db.execute(sql`
INSERT INTO table (${sql.raw(insertCols.join(', '))}, week_start_date)
VALUES (
${sql.join(insertVals.map(v => sql`${v}`), sql`,`)},
date_trunc('week', ${reportDate}::date)::date -- ← 直接写在模板里
)
...
`);

关键区别:

写法Drizzle 处理方式PostgreSQL 收到
sql 模板插值参数化($N字符串字面量
sql.raw(expression)原样拼入 SQLSQL 表达式
直接写在 sql 模板中作为模板的一部分SQL 表达式

注意事项

注意事项

  • sql.raw() 存在 SQL 注入风险,不要用于用户输入。本例中 reportDate 来自内部 API,格式可控
  • Drizzle 的 sql 模板标签会自动参数化所有插值——这是安全特性,但 SQL 函数调用不该被参数化
  • 如果整条 SQL 都是动态构建的,考虑用 Drizzle 的 query builder API 代替 raw SQL
  • 数据库连接配置也容易踩坑——如果你遇到连接到了错误的 PostgreSQL 实例,可能是端口被 Docker 静默占用
  • 环境变量加载时序也是常见坑源,JWT 签名静默失败就是 dotenv 在 import 链之后才执行的典型例子

WSL2 + Docker 两个网络坑:端口被静默占用 & host 模式 localhost 不通

· 阅读需 5 分钟

TL;DR

WSL2 + Docker Desktop 有两个常见的网络坑:

  1. 端口被静默占用:Docker 容器映射 5432 后,SSH 隧道 localhost:5432 连到的是容器内的 PostgreSQL 而非远程服务器——密码没错,连的实例错了
  2. host 模式 localhost 不通network_mode: host 共享的是 Docker 工具 VM 网络,不是 WSL2 网络——curl localhost:8080 失败