数据操作
约 2521 字大约 8 分钟
2026-08-30
Core
Raw SQL
# model/main.py
import asyncio
from sqlalchemy import text
from app.model.engine import get_engine
from app.model.base import Base
import app.model.category # noqa: F401
import app.model.product # noqa: F401
import app.model.sku # noqa: F401
async def init_db():
engine = get_engine()
async with engine.begin() as conn:
await conn.run_sync(Base.metadata.drop_all)
await conn.run_sync(Base.metadata.create_all)
print("✅ 数据库表创建完成")
async def raw_sql_insert():
"""直接用 SQL 字符串插入数据"""
engine = get_engine()
async with engine.begin() as conn:
await conn.execute(
text(
"INSERT INTO product (name, description, brand) VALUES ('iPhone 16', '最新款智能手机', 'Apple')"
)
)
await conn.execute(
text(
"INSERT INTO product (name, description, brand) VALUES ('iPhone 15', '上一代旗舰', 'Apple')"
)
)
await conn.execute(
text(
"INSERT INTO product (name, description, brand) VALUES ('Galaxy S25', '三星旗舰', 'Samsung')"
)
)
await conn.execute(
text(
"INSERT INTO product (name, description, brand) VALUES ('Xiaomi 15', '性价比之选', 'Xiaomi')"
)
)
print("✅ 数据插入完成")
async def search_products(keyword: str):
"""按关键字搜索商品——接受用户输入,拼接 SQL"""
engine = get_engine()
async with engine.connect() as conn:
sql = f"SELECT * FROM product WHERE name LIKE '%{keyword}%'"
result = await conn.execute(text(sql))
rows = result.fetchall()
for row in rows:
print(f" [{row.id}] {row.name} - {row.brand}")
async def delete_product_by_id(product_id: str):
"""按 ID 删除商品——接受用户输入,拼接 SQL"""
engine = get_engine()
async with engine.begin() as conn:
sql = f"DELETE FROM product WHERE id = {product_id}"
await conn.execute(text(sql))
print(f"✅ 删除完成,id={product_id}")
async def main():
await init_db()
await raw_sql_insert()
# 正常查询
print("=== 搜索 'iPhone' ===")
await search_products("iPhone")
# 正常删除
print("=== 删除 id=1 ===")
await delete_product_by_id("1")
print("=== 再次搜索 'iPhone' ===")
await search_products("iPhone")
if __name__ == "__main__":
asyncio.run(main())运行后输出:
✅ 数据库表创建完成
✅ 数据插入完成
=== 搜索 'iPhone' ===
[1] iPhone 16 - Apple
[2] iPhone 15 - Apple
=== 删除 id=3 ===
✅ 删除完成,id=3
=== 再次搜索 'iPhone' ===
[1] iPhone 16 - Apple
[2] iPhone 15 - Apple到这里,数据能插入、能查询、能按 ID 删除了,看起来一切正常。
SQL 注入
上面的 search_products 直接把用户输入拼进了 SQL 字符串。来验证一下:
# 正常搜索
await search_products("iPhone")输出符合预期:
=== 搜索 'iPhone' ===
[1] iPhone 16 - Apple
[2] iPhone 15 - Apple现在换攻击者的输入:
await search_products("%' OR 1=1 --")拼接后的 SQL:
SELECT * FROM product WHERE name LIKE '%%' OR 1=1 -- %'OR 1=1 永远为真,-- 注释掉后续内容——所有数据被暴露:
[1] iPhone 16 - Apple
[2] iPhone 15 - Apple
[3] Galaxy S25 - Samsung
[4] Xiaomi 15 - Xiaomi这还只是查询。同样的拼接方式在按 ID 删除里一样有效:
# 正常删除,删掉 id=3(Galaxy S25)
await delete_product_by_id("3")
# 输出:✅ 删除完成,id=3现在攻击者传入:
await delete_product_by_id("1 OR 1=1")拼接后的 SQL:
DELETE FROM product WHERE id = 1 OR 1=1WHERE id = 1 OR 1=1 匹配所有行——整张表被清空:
# 再查询,数据已经没了
await search_products("iPhone")
# 输出:(无结果,数据已被删光)这就是经典的 SQL 注入——攻击者通过篡改输入改变 SQL 语句的意图,轻则偷数据,重则删库跑路。
修复:参数化绑定
解决方案是:SQL 结构与数据分离。用 :参数名 作为占位符,通过字典传入值:
async def search_products(keyword: str):
"""按关键字搜索商品——参数化绑定,安全"""
engine = get_engine()
async with engine.connect() as conn:
result = await conn.execute(
text("SELECT * FROM product WHERE name LIKE :keyword"),
{"keyword": f"%{keyword}%"},
)
rows = result.fetchall()
for row in rows:
print(f" [{row.id}] {row.name} - {row.brand}")再用恶意输入试一次:
await search_products("%' OR 1=1 --")
# 输出:(无结果,%' OR 1=1 -- 被当作普通字符串匹配,不会命中任何商品名)原理:参数化绑定将 SQL 模板和参数值分开发送给数据库。数据库先解析 SQL 结构,再把参数值作为纯数据填入。无论攻击者输入什么,都不会被当作 SQL 指令执行。
手写 SQL 能跑,但两个核心问题摆在这里:一是注入风险,二是写字符串容易出错、没有 IDE 检查和补全。下一节引入 Core 表达式来解决这两个问题。
总结
使用RAW SQL会带来以下问题:
- 开发成本高
- 易出错
- 有注入风险
- 无法适配不同数据库
但由于其高性能和极大的灵活性,仍然会有小概率用到
SQL表达式
SQLAlchemy Core 提供了一套 Pythonic 的 API 来构建 SQL 语句,底层自动完成参数化绑定,彻底杜绝注入风险。
# model/main.py
import asyncio
from sqlalchemy import select, insert, update, delete
from app.model.engine import get_engine
from app.model.base import Base
from app.model.product import Product
from app.model.category import Category
import app.model.sku # noqa: F401
async def init_db():
engine = get_engine()
async with engine.begin() as conn:
await conn.run_sync(Base.metadata.drop_all)
await conn.run_sync(Base.metadata.create_all)
print("✅ 数据库表创建完成")
async def core_insert():
"""Core 表达式插入"""
engine = get_engine()
async with engine.begin() as conn:
# insert() 返回 Insert 对象,values() 设置列值
stmt = insert(Product).values(
name="iPhone 16",
description="最新款智能手机",
brand="Apple",
)
await conn.execute(stmt)
stmt = insert(Product).values(
name="iPhone 15",
description="上一代旗舰",
brand="Apple",
)
await conn.execute(stmt)
stmt = insert(Category).values(
name="手机", description="移动通信设备"
)
await conn.execute(stmt)
print("✅ Core 插入完成")
async def core_query():
"""Core 表达式查询"""
engine = get_engine()
async with engine.connect() as conn:
# select() 返回 Select 对象
# where() 用 Python 表达式构建条件——Python 的 == 而非 SQL 的 =
stmt = select(Product).where(
Product.name.like("%iPhone%")
)
result = await conn.execute(stmt)
for row in result:
print(f" [{row.id}] {row.name} - {row.brand}")
async def core_update():
"""Core 表达式更新"""
engine = get_engine()
async with engine.begin() as conn:
stmt = (
update(Product)
.where(Product.id == 1)
.values(name="iPhone 16 Pro")
)
await conn.execute(stmt)
print("✅ Core 更新完成")
async def core_delete():
"""Core 表达式删除"""
engine = get_engine()
async with engine.begin() as conn:
stmt = delete(Product).where(Product.id == 2)
await conn.execute(stmt)
print("✅ Core 删除完成")
async def main():
await init_db()
await core_insert()
print("=== 查询 ===")
await core_query()
await core_update()
print("=== 查询(更新后) ===")
await core_query()
await core_delete()
print("=== 查询(删除后) ===")
await core_query()
if __name__ == "__main__":
asyncio.run(main())总结
使用SQL表达式有以下好处:
- 类型化,不易出错
- 无SQL注入风险
- 适配不同方言
会话
# model/main.py
import asyncio
from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker
from sqlalchemy import select, insert
from app.model.engine import get_engine
from app.model.base import Base
from app.model.product import Product
import app.model.category # noqa: F401
import app.model.sku # noqa: F401
session_factory = async_sessionmaker(
get_engine(),
class_=AsyncSession, # 明确指定使用异步 Session
expire_on_commit=False,
)
async def init_db():
engine = get_engine()
async with engine.begin() as conn:
await conn.run_sync(Base.metadata.drop_all)
await conn.run_sync(Base.metadata.create_all)
print("✅ 数据库表创建完成")
async def insert_with_session():
session = session_factory()
try:
stmt = insert(Product).values(
name="iPhone 16",
description="最新款智能手机",
brand="Apple",
)
await session.execute(stmt)
await session.commit()
except:
await session.rollback()
raise
finally:
await session.close()
print("✅ 提交完成")
async def insert_with_manual_commit():
"""手动提交"""
async with session_factory() as session:
try:
stmt = insert(Product).values(
name="iPhone 16",
description="最新款智能手机",
brand="Apple",
)
await session.execute(stmt)
await session.commit()
except:
await session.rollback()
raise
print("✅ 手动提交完成")
async def insert_with_auto_commit():
"""自动提交"""
async with session_factory.begin() as session:
stmt = insert(Product).values(
name="iPhone 15",
description="上一代旗舰",
brand="Apple",
)
await session.execute(stmt)
print("✅ 自动提交完成")
async def main():
await init_db()
await insert_with_session()
await insert_with_manual_commit()
await insert_with_auto_commit()
if __name__ == "__main__":
asyncio.run(main())优化代码结构
新建core/database.py
from sqlalchemy.ext.asyncio import (
AsyncEngine,
AsyncSession,
create_async_engine,
async_sessionmaker,
)
from app.core.config import db_settings
_engine: AsyncEngine | None = None
def get_engine() -> AsyncEngine:
global _engine
if _engine is None:
url = f"postgresql+asyncpg://{db_settings.user}:{db_settings.password}@{db_settings.host}:{db_settings.port}/{db_settings.name}"
_engine = create_async_engine(
url, pool_size=10, max_overflow=20, pool_pre_ping=True, echo=False
)
return _engine
_session_factory: async_sessionmaker[AsyncSession] | None = None
def get_session_factory() -> async_sessionmaker[AsyncSession]:
global _session_factory
if _session_factory is None:
_session_factory = async_sessionmaker(
get_engine(),
class_=AsyncSession, # 明确指定使用异步 Session
expire_on_commit=False,
)
return _session_factory现在可以删除掉model/engine.py了
ORM
前面两节无论是 Raw SQL 还是 Core 表达式,操作的都是表(Table)。而 ORM 的核心思想是:操作的是 Python 对象,SQLAlchemy 负责把对象的变化同步到数据库。
# model/main.py
import asyncio
from sqlalchemy import select
from app.core.database import get_session_factory, get_engine
from app.model.base import Base
from app.model.product import Product
from app.model.category import Category
import app.model.sku # noqa: F401
session_factory = get_session_factory()
async def init_db():
engine = get_engine()
async with engine.begin() as conn:
await conn.run_sync(Base.metadata.create_all)
print("✅ 数据库表创建完成")
async def orm_insert():
"""ORM 方式插入——创建对象,add 到 session"""
async with session_factory.begin() as session:
product = Product(
name="iPhone 16",
description="最新款智能手机",
brand="Apple",
)
session.add(product)
product2 = Product(
name="iPhone 15",
description="上一代旗舰",
brand="Apple",
)
session.add(product2)
category = Category(
name="手机", description="移动通信设备"
)
session.add(category)
print("✅ ORM 插入完成")
async def orm_query():
"""ORM 方式查询——execute + scalars() 拿到对象列表"""
async with session_factory() as session:
stmt = select(Product).where(
Product.name.like("%iPhone%")
)
result = await session.execute(stmt)
# scalars() 返回模型实例,可直接访问属性
products = result.scalars().all()
for p in products:
print(f" [{p.id}] {p.name} - {p.brand}")
async def orm_update():
"""ORM 方式更新——查出对象,改属性,自动同步"""
async with session_factory.begin() as session:
stmt = select(Product).where(Product.id == 1)
result = await session.execute(stmt)
product = result.scalar_one()
# 直接修改 Python 对象属性,session 提交时自动生成 UPDATE
product.name = "iPhone 16 Pro"
print("✅ ORM 更新完成")
async def orm_delete():
"""ORM 方式删除——查出对象,调用 session.delete()"""
async with session_factory.begin() as session:
stmt = select(Product).where(Product.id == 2)
result = await session.execute(stmt)
product = result.scalar_one()
await session.delete(product)
print("✅ ORM 删除完成")
async def main():
await init_db()
await orm_insert()
print("=== 查询 ===")
await orm_query()
await orm_update()
print("=== 查询(更新后) ===")
await orm_query()
await orm_delete()
print("=== 查询(删除后) ===")
await orm_query()
if __name__ == "__main__":
asyncio.run(main())执行时机
session 在下面的情况下会触发执行 SQL
session.flush(),执行记录的操作,但不提交session.commit(),先自动flush,再提交session.execute(),执行传入的 SQL 语句,不提交- 当
autoflush=True,则会先将 session 中 pending 的操作flush到数据库,再执行传入的 SQL,这是默认值 - 当
autoflush=False,仅执行传入的 SQL,pending 操作保留在 session 中不动
- 当
模型对象.字段,当expire_on_commit=True时,可能会触发执行,不提交当 session 提交后,它会把 session 当中关联的模型对象的每一个字段都标记为过期,因为这些字段可能已经跟数据库不一样了。后续读取这些字段的时候,它会重新触发查询,将该对象的所有字段同步为数据库的字段值。
这一切的前提条件是要开启
expire_on_commit=True开启这个配置会损耗性能,同时会带来一些其他问题,所以说往往会关闭。
session.refresh(模型对象),同步模型对象,不提交- 将数据库真实的值同步到模型对象。当
expire_on_commit=False时,才可能需要使用该方法
- 将数据库真实的值同步到模型对象。当
总结
- 简单 CRUD 用 ORM 更顺手,复杂查询(聚合、批量更新)可混用 Core 表达式,前两者都难以处理的情况下,可以考虑使用 RAW SQL,但要注意 SQL 注入问题
- 耗时操作(比如I/O)尽量不要放到一次连接中
session的能力总结:- 模型结合
- 统一管理事务
- 记录操作、统一执行
作业
SQLAlchemy提供了几种操作数据的方式?分别适用于什么场景?- 什么是SQL注入?如何避免?
SQLAlchemy的会话功能有什么用?expire_on_commit这个配置有什么用?flush和commit有什么区别?
