034AI Agent阅读记录
第 034 卷

从零开始 AI Agent 实战(五):当 JSON 文件撑不住时换上 PostgreSQL

复现 JSON 存储的并发撞号问题,使用 PostgreSQL、异步 SQLAlchemy、序列和唯一约束完成可替换 Repository,并提前装好 pgvector。

乌漆嘛黑和 Ahri
第 034 期

当 JSON 文件撑不住时换上 PostgreSQL

先复现两个请求拿到同一个编号

JSON Repository 的编号算法是 len(rows) + 1。两个请求同时读取 41 条记录,都会计算出 TICKET-000042,最后写入两张同号工单。即使加 Python 锁,也只对单进程有效;多副本部署后锁各自存在。

async def race(repo):
    return await asyncio.gather(
        service.create('VPN A', '...', 'u-1'),
        service.create('VPN B', '...', 'u-2'),
    )

这不是“换个更快的 JSON 库”能解决的问题,而是需要数据库原子操作和唯一约束。

Preview
修复前后的两个事务时序对照

表结构先服务于查询

九张核心表不会在本篇一次建完。先看它们的关系,再决定这一篇动哪几张:

Preview
九张表分三条业务线,深色的是本篇落地的部分
create sequence ticket_number_seq start 1;

create table tickets (
  id uuid primary key,
  number text not null unique,
  title text not null,
  description text not null,
  status text not null check (status in ('new','triaging','need_info','resolved','workaround','closed')),
  reporter_id uuid not null references users(id),
  app_version text,
  os_name text,
  error_code text,
  created_at timestamptz not null default now(),
  deleted_at timestamptz
);
create index tickets_reporter_created_idx on tickets(reporter_id, created_at desc)
  where deleted_at is null;

工单号可由序列生成:

select format('TICKET-%s', lpad(nextval('ticket_number_seq')::text, 6, '0'));

即使应用重试,number unique 也会让错误显性化,而不是悄悄产生重复数据。

序列有空洞,什么时候需要 advisory lock

序列的递增不参与事务回滚:事务 A 取到 42 后回滚,42 这个号就永久消失了。对内部工单编号这通常无所谓,但如果业务要求编号连续无空洞(发票、合同这类场景),序列就不能用。

那时的替代方案是咨询锁——让同一类编号的分配串行化:

async def next_number_serialized(session: AsyncSession) -> str:
    # 同一个 key 上的并发调用会排队,锁在事务结束时自动释放(42 为工单锁专属标识)
    await session.execute(text('select pg_advisory_xact_lock(:key)'), {'key': 42})
    row = await session.execute(text('select coalesce(max(seq), 0) + 1 from tickets'))
    return f'TICKET-{row.scalar():06d}'

代价也很明确:所有创建工单的请求在这把锁上排成一队,写入吞吐被锁的持有时间限制。所以默认选序列,只有在「不许有空洞」是硬需求时才换成 advisory lock——不要因为它听起来更严谨就先用它

Repository 第二个实现

class SqlTicketRepository:
    def __init__(self, session_factory):
        self.session_factory = session_factory

    async def next_number(self):
        async with self.session_factory() as s:
            value = await s.scalar(text("select nextval('ticket_number_seq')"))
            return f'TICKET-{value:06d}'

    async def save(self, ticket):
        async with self.session_factory.begin() as s:
            await s.execute(insert(TicketRow).values(
                id=ticket.id, number=ticket.number, title=ticket.title,
                description=ticket.description, status=ticket.status.value,
                reporter_id=ticket.reporter_id,
            ))
        return ticket

application 层只依赖 Protocol,第 1 篇的领域测试无需修改。需要注意 next_numbersave 最好放在同一个事务用例里;更稳妥的写法是让 Repository 提供 create,在一次事务内取序列并插入。

Preview
接口不变,实现整体替换,只改一行组装代码

这里用的是 SQLAlchemy 2.0 的异步 Session。2.0 相比 1.x 的关键差异是 select() 风格的统一查询构造和原生 async 支持——旧版 Query API 在异步下有很多陷阱,网上大量教程仍是 1.x 写法,照抄会遇到 MissingGreenlet 之类的报错。

为什么现在就装 pgvector

第 9 篇要把文档 chunk 存到向量列,第 10 篇再加上群聊。如果现在才改镜像、迁移和连接池,读者会在中途遇到一轮与 Agent 无关的环境重建。先安装扩展但不使用它,成本很低:

create extension if not exists vector;
alter table document_chunks add column embedding vector(1536);
create index document_chunks_embedding_idx on document_chunks
  using hnsw (embedding vector_cosine_ops);

维度必须和 embedding 模型固定,换模型时新建列或重建索引,不能把 1536 维向量塞进 768 维列。

迁移、软删除与分页

Alembic 迁移要纳入 CI;不要直接在生产数据库手工执行 SQL。分页使用游标而不是 offset

select * from tickets
where reporter_id = :reporter
  and deleted_at is null
  and (created_at, id) < (:cursor_time, :cursor_id)
order by created_at desc, id desc
limit 20;

软删除字段参与所有查询条件,否则“删除”的工单会被 Agent 再次检索出来。

测试和验收

docker compose up -d postgres
uv run alembic upgrade head
uv run pytest tests/repositories/test_sql_ticket_repo.py -q
uv run pytest tests/repositories/test_concurrent_ticket_number.py -q

并发测试要真正启动两个数据库事务,而不是在同一个 Fake 上调用两次函数。断开数据库后应返回结构化 database_unavailable,API 不得泄露 SQL 语句。

本篇验收标准:并发创建 100 次无重复 number;旧 application 测试全部通过;docker compose up 能加载 vector 扩展;文章列表查询有索引命中(用 EXPLAIN (ANALYZE, BUFFERS) 留下证据)。

下一篇补上真实用户。没有认证,检索和工具权限都只是注释里的愿望。