Python数据库篇:sqlite3、pymysql、aiomysql、sqlalchemy、pydantic
一:sqlite3
import sqlite3
conn = sqlite3.connect("test.db")
cursor = conn.cursor()
cursor.execute("create table user (id varchar(20) primary key, name varchar(20))")
cursor.execute("insert into user (id, name) values (\'1\', \'Michael\')")
print(cursor.rowcount)
conn.commit()
cursor.execute("select * from user where id=?", ("1",))
rows = cursor.fetchall()
print(rows)
cursor.close()
conn.close()
二:MySQL
pip install pymysql
import pymysql
from pymysql.cursors import DictCursor
# 1.创建连接
# cursorclass: 全局定义,表示查询时是返回的数据类型,字典还是元组(默认)
conn = pymysql.connect(host='127.0.0.1',
port=3306,
user='root',
password='123456',
database='test',
charset='utf8',
autocommit=True,
cursorclass=DictCursor)
# 2.获取游标,局部定义cursor
cursor = conn.cursor(cursor=pymysql.cursors.DictCursor)
# 3.定义sql,最好使用"""来定义字符串,这样关键字会变色
create_table_sql = """
CREATE TABLE IF NOT EXISTS sys_user2 (
`id` bigint(20) unsigned NOT NULL AUTO_INCREMENT,
`username` varchar(15) NOT NULL COMMENT '用户名',
`gender` tinyint(3) unsigned DEFAULT '0' COMMENT '性别(0: 女 1:男)',
`amount` decimal(10,2) DEFAULT NULL,
`create_time` datetime DEFAULT CURRENT_TIMESTAMP COMMENT '注册时间',
PRIMARY KEY (`id`),
UNIQUE KEY `uk` (`username`) USING BTREE
) ENGINE=InnoDB AUTO_INCREMENT=1 DEFAULT CHARSET=utf8
"""
# 4.执行
cursor.execute(create_table_sql)
delete_sql = """
delete from sys_user2
"""
cursor.execute(delete_sql)
insert_sql = """
insert into sys_user2(username, gender, amount, create_time) values(%s, %s, %s, now())
"""
# 插入1条
cursor.execute(insert_sql, ['monday', 0, 9999999])
# 插入多条
cursor.executemany(insert_sql, [('test', 0, 9999999), ('modely', 0, 66666666)])
update_sql = """
update sys_user2 set amount = %s where username = %s
"""
cursor.execute(update_sql, ['88888888', 'modely'])
# 注意:先查询再获取数据,不像其它语言一样直接返回结果
select_sql = """
select * from sys_user2
"""
cursor.execute(select_sql)
fetchone = cursor.fetchone()
fetchmany = cursor.fetchmany(2)
# 注意:总共有三条,上面游标已经走到最后了,所以下面游标再走就没有数据了
fetchall = cursor.fetchall()
# insert, update, delete 操作都需要进行提交,如果配置了autocommit=True就不需要显式提交了会自动提交
# conn.commit()
# conn.rollback()
# 5.关闭游标
cursor.close()
# 关闭连接
conn.close()
三、aiomysql
aiomysql 是 Python 中基于 asyncio 的异步 MySQL 客户端,它在 PyMySQL 的基础上做了异步封装,API 风格几乎与 PyMySQL 一致,但所有 I/O 操作都改成了 async/await 的协程形式。
1. 简介
核心特性
| 特性 | 说明 |
|---|---|
| 异步非阻塞 | 基于 asyncio,不阻塞事件循环 |
| 类 PyMySQL API | 学习成本低,迁移容易 |
| 连接池 | 内置 create_pool,支持高并发 |
| 游标支持 | DictCursor、SSCursor、SSDictCursor 等 |
| 事务 | 支持 begin/commit/rollback |
| SQLAlchemy 集成 | 可作为 aiomysql.sa 后端用于异步 ORM |
适用场景
- 高并发 Web 服务(FastAPI、aiohttp、Sanic)的 MySQL 访问
- 异步爬虫的数据落库
- 大量 I/O 密集型数据库操作(批量查询、轮询、长连接)
- 需要与
aiohttp、aioredis等异步生态搭配使用
2. 安装
# aiomysql 会自动安装依赖 `PyMySQL`。
pip install aiomysql
# 如果需要 SQLAlchemy Core 异步支持
pip install aiomysql[sa]
2. 快速入门
推荐使用
async with,可保证连接和游标在退出时自动释放,避免泄漏。
import asyncio
import aiomysql
async def main():
async with aiomysql.connect(
host="127.0.0.1", port=3306,
user="root", password="123456",
db="test",
) as conn:
async with conn.cursor() as cur:
await cur.execute("SELECT NOW()")
result = await cur.fetchone()
print(result)
# 自动关闭连接和游标
asyncio.run(main())
3. CRUD 操作
1. 创建表 execute
async def create_table(conn):
async with conn.cursor() as cur:
await cur.execute("""
CREATE TABLE IF NOT EXISTS users (
id INT AUTO_INCREMENT PRIMARY KEY,
name VARCHAR(50) NOT NULL,
age INT,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4
""")
await conn.commit()
2. 插入数据(参数化查询防 SQL 注入) execute
async def insert_user(conn, name, age):
async with conn.cursor() as cur:
# 使用 %s 占位符,值通过参数传入(切勿用 f-string 拼接 SQL!)
sql = "INSERT INTO users (name, age) VALUES (%s, %s)"
await cur.execute(sql, (name, age))
await conn.commit()
print(f"插入成功,新 ID: {cur.lastrowid}")
3. 批量插入 executemany
async def insert_many(conn, users):
async with conn.cursor() as cur:
sql = "INSERT INTO users (name, age) VALUES (%s, %s)"
await cur.executemany(sql, users) # users 是 [(name, age), ...]
await conn.commit()
print(f"批量插入 {cur.rowcount} 条")
# 调用
await insert_many(conn, [
("Alice", 25),
("Bob", 30),
("Charlie", 28),
])
4. 查询数据 fetchone、fetchmany、fetchall
async def query_users(conn):
async with conn.cursor() as cur:
await cur.execute("SELECT id, name, age FROM users WHERE age > %s", (20,))
# 取一行
row = await cur.fetchone()
print("第一行:", row) # (1, 'Alice', 25)
# 取多行
rows = await cur.fetchmany(10)
print(f"取了 {len(rows)} 行")
# 取全部
all_rows = await cur.fetchall()
print("剩余:", all_rows)
5. 字典游标(返回 dict 而不是 tuple) aiomysql.DictCursor
async def query_as_dict(conn):
# 在创建游标时指定 DictCursor
async with conn.cursor(aiomysql.DictCursor) as cur:
await cur.execute("SELECT id, name, age FROM users LIMIT 3")
rows = await cur.fetchall()
for row in rows:
print(f"{row['name']} ({row['age']})")
# 输出:Alice (25), Bob (30), ...
6. 更新与删除
async def update_user(conn, user_id, new_age):
async with conn.cursor() as cur:
await cur.execute(
"UPDATE users SET age = %s WHERE id = %s",
(new_age, user_id),
)
await conn.commit()
print(f"更新行数: {cur.rowcount}")
async def delete_user(conn, user_id):
async with conn.cursor() as cur:
await cur.execute("DELETE FROM users WHERE id = %s", (user_id,))
await conn.commit()
4. 连接池(高并发必备)
频繁建立和关闭连接代价很高,生产环境必须使用连接池。
1. 创建连接池
import asyncio
import aiomysql
async def main():
pool = await aiomysql.create_pool(
host="127.0.0.1", port=3306,
user="root", password="123456",
db="test",
minsize=1, # 最小连接数
maxsize=10, # 最大连接数
autocommit=False,
charset="utf8mb4",
)
# 从池中获取连接
async with pool.acquire() as conn:
async with conn.cursor() as cur:
await cur.execute("SELECT 1")
result = await cur.fetchone()
print(result)
# 关闭连接池
pool.close()
await pool.wait_closed()
asyncio.run(main())
连接池常用参数
| 参数 | 说明 | 默认值 |
|---|---|---|
minsize |
池中最少保持的空闲连接数 | 1 |
maxsize |
池中最大连接数 | 10 |
pool_recycle |
连接最长复用时间(秒),超过则重建,防 wait_timeout | -1(不回收) |
echo |
是否打印 SQL | False |
autocommit |
是否自动提交 | False |
charset |
字符集 | None |
重要:MySQL 默认
wait_timeout是 8 小时,但运维常配置成几分钟。生产环境建议设置pool_recycle=3600(1 小时),避免使用到被服务端关闭的连接。
2. 并发查询(连接池真正发挥作用的场景)
import asyncio
import aiomysql
import time
async def query_one(pool, user_id):
async with pool.acquire() as conn:
async with conn.cursor() as cur:
await cur.execute("SELECT * FROM users WHERE id = %s", (user_id,))
return await cur.fetchone()
async def main():
pool = await aiomysql.create_pool(
host="127.0.0.1", port=3306,
user="root", password="123456",
db="test", minsize=5, maxsize=20,
)
start = time.time()
# 并发发起 100 个查询
results = await asyncio.gather(*(query_one(pool, i) for i in range(1, 101)))
print(f"100 次查询耗时: {time.time() - start:.2f}s")
pool.close()
await pool.wait_closed()
asyncio.run(main())
5.事务
1. 基本事务
async def transfer_money(pool, from_id, to_id, amount):
async with pool.acquire() as conn:
async with conn.cursor() as cur:
try:
await conn.begin() # 开启事务
await cur.execute(
"UPDATE accounts SET balance = balance - %s WHERE id = %s",
(amount, from_id),
)
await cur.execute(
"UPDATE accounts SET balance = balance + %s WHERE id = %s",
(amount, to_id),
)
await conn.commit() # 提交
print("转账成功")
except Exception as e:
await conn.rollback() # 回滚
print(f"转账失败,已回滚: {e}")
raise
2. 显式控制 autocommit
aiomysql 提供 三种 控制 autocommit 的方式,可根据场景灵活切换:
方式 ①:创建池/连接时通过 autocommit 参数设定
# 全局关闭自动提交(默认行为):事务必须手动 commit
pool = await aiomysql.create_pool(
host="127.0.0.1", port=3306,
user="root", password="123456",
db="test",
autocommit=False, # ← 显式声明,所有从池中取出的连接默认都关闭自动提交
)
async with pool.acquire() as conn:
async with conn.cursor() as cur:
await conn.begin() # 开启事务
try:
await cur.execute("INSERT INTO users(name) VALUES (%s)", ("Alice",))
await cur.execute("INSERT INTO users(name) VALUES (%s)", ("Bob",))
await conn.commit() # 必须手动提交
except Exception:
await conn.rollback()
raise
方式 ②:运行时调用 conn.autocommit() 动态切换
async with pool.acquire() as conn:
# 运行时把当前连接切到 autocommit=True
await conn.autocommit(True) # ← 显式打开自动提交
async with conn.cursor() as cur:
await cur.execute("INSERT INTO logs(msg) VALUES (%s)", ("hello",))
# 不需要 commit,执行即落库
# 临时切回手动模式做事务
await conn.autocommit(False) # ← 显式关闭自动提交
async with conn.cursor() as cur:
await conn.begin()
try:
await cur.execute("UPDATE balance SET v = v - 100 WHERE id = 1")
await cur.execute("UPDATE balance SET v = v + 100 WHERE id = 2")
await conn.commit()
except Exception:
await conn.rollback()
raise
方式 ③:创建池时打开 autocommit,临时事务用 begin() 包裹
pool = await aiomysql.create_pool(
host="127.0.0.1", port=3306,
user="root", password="123456",
db="test",
autocommit=True, # ← 默认每条 SQL 立即落库,适合查询多/写少的场景
)
async with pool.acquire() as conn:
async with conn.cursor() as cur:
# 普通写入:无需 commit
await cur.execute("INSERT INTO logs(msg) VALUES (%s)", ("auto-committed",))
# 需要事务时,用 begin() 临时关闭 autocommit,直到 commit/rollback
async with conn.cursor() as cur:
await conn.begin()
try:
await cur.execute("UPDATE accounts SET balance = balance - 100 WHERE id = 1")
await cur.execute("UPDATE accounts SET balance = balance + 100 WHERE id = 2")
await conn.commit() # ← 同时恢复 autocommit=True
except Exception:
await conn.rollback()
raise
三种方式对比
| 方式 | 适用场景 | 优点 | 缺点 |
|---|---|---|---|
① 池级 autocommit=False |
业务以事务为主(转账、订单) | 安全,逻辑清晰 | 简单查询也要 commit |
② 运行时 conn.autocommit() 切换 |
同一连接需在两种模式间切换 | 灵活 | 心智负担高,易乱 |
③ 池级 autocommit=True + begin() |
写少读多、日志类落库 | 简化代码 | 忘了 begin() 就没事务保护 |
当
autocommit=False(默认)时,任何修改 SQL 都需要commit()才会生效——这是 aiomysql 新手最常犯的"插了数据但库里没有"问题的根源。一条经验法则:新手默认用方式 ①(池级关闭 + 显式 commit),最安全。
6. SSCursor:大结果集流式读取
普通游标会一次性把所有结果加载到内存,百万行结果会撑爆内存出现OOM。SSCursor(Server-Side Cursor)逐行从服务端拉取。或使用 SSDictCursor 同时拿到字典格式。
async def stream_large_table(conn):
async with conn.cursor(aiomysql.SSCursor) as cur:
await cur.execute("SELECT * FROM big_table")
while True:
row = await cur.fetchone()
if row is None:
break
# 处理每一行,内存占用恒定
process(row)
7. 常见陷阱与最佳实践
1. 忘记 commit()
# ❌ 数据没插进去
async with conn.cursor() as cur:
await cur.execute("INSERT INTO users (name) VALUES ('Alice')")
# autocommit=False 时,没 commit 就退出 = 自动 rollback
# ✅ 必须 commit
async with conn.cursor() as cur:
await cur.execute("INSERT INTO users (name) VALUES ('Alice')")
await conn.commit()
或在创建连接/池时设置 autocommit=True(只适合简单场景)。
2. 千万不要用 f-string 拼接 SQL
# ❌ SQL 注入风险!
name = request.args.get("name")
await cur.execute(f"SELECT * FROM users WHERE name = '{name}'")
# ✅ 使用参数化查询
await cur.execute("SELECT * FROM users WHERE name = %s", (name,))
3. 注意占位符是 %s,不要加引号
aiomysql 沿用 PyMySQL 的占位符:所有类型(字符串、整数、日期)都用 %s,不需要加引号。
# ✅ 正确
await cur.execute("SELECT * FROM users WHERE id = %s AND name = %s", (1, "Alice"))
# ❌ 错误:不要加引号
await cur.execute("SELECT * FROM users WHERE name = '%s'", ("Alice",))
4. 处理连接失效
生产环境网络抖动 / MySQL 重启会导致连接失效。建议:
- 设置
pool_recycle让连接定期重建 - 在关键操作外层加重试装饰器(配合 tenacity)
from tenacity import retry, stop_after_attempt, wait_exponential
@retry(stop=stop_after_attempt(3), wait=wait_exponential(min=1, max=5))
async def safe_query(pool, sql, params):
async with pool.acquire() as conn:
async with conn.cursor() as cur:
await cur.execute(sql, params)
return await cur.fetchall()
5. aiomysql vs 其他 MySQL 库
| 库 | 类型 | 性能 | 维护状态 | 备注 |
|---|---|---|---|---|
| aiomysql | 异步(asyncio) | 中等 | 活跃 | 基于纯 Python 的 PyMySQL |
| asyncmy | 异步(asyncio) | 高 | 活跃 | C 扩展实现,性能领先 |
| PyMySQL | 同步 | 中等 | 活跃 | 纯 Python,跨平台 |
| mysqlclient | 同步 | 高 | 活跃 | C 扩展,Django 默认 |
| mysql-connector-python | 同步 | 中等 | 官方 | Oracle 官方提供 |
| databases | 异步封装层 | 取决于驱动 | 活跃 | 统一封装多种异步驱动 |
选型建议:
- 入门、对性能不极端要求 → aiomysql
- 极致性能、大规模并发 → asyncmy
- 同步项目 → PyMySQL 或 mysqlclient
- 想用 ORM → SQLAlchemy 2.0 + asyncmy 或 Tortoise ORM
9. 完整示例:异步任务从 MySQL 读 → 处理 → 写回
import asyncio
import aiomysql
async def worker(pool, worker_id):
"""每个 worker 不停从待处理表取一条记录,处理后写回。"""
while True:
async with pool.acquire() as conn:
async with conn.cursor(aiomysql.DictCursor) as cur:
# 锁一行待处理任务(SKIP LOCKED 避免多 worker 抢同一行)
await conn.begin()
await cur.execute("""
SELECT id, payload FROM tasks
WHERE status = 'pending'
LIMIT 1 FOR UPDATE SKIP LOCKED
""")
task = await cur.fetchone()
if not task:
await conn.commit()
await asyncio.sleep(1) # 没任务就休息一下
continue
# 标记为处理中
await cur.execute(
"UPDATE tasks SET status='running' WHERE id=%s",(task["id"],),
)
await conn.commit()
# 模拟业务处理(放在锁外)
print(f"worker-{worker_id} 处理任务 {task['id']}")
await asyncio.sleep(0.5)
# 写回结果
async with pool.acquire() as conn:
async with conn.cursor() as cur:
await cur.execute(
"UPDATE tasks SET status='done' WHERE id=%s",
(task["id"],),
)
await conn.commit()
async def main():
pool = await aiomysql.create_pool(
host="127.0.0.1", port=3306,
user="root", password="123456",
db="test", minsize=5, maxsize=20,
autocommit=False,
)
# 启动 5 个 worker 并发处理
workers = [asyncio.create_task(worker(pool, i)) for i in range(5)]
try:
await asyncio.gather(*workers)
finally:
pool.close()
await pool.wait_closed()
asyncio.run(main())
10.总结
aiomysql是 asyncio 生态下最常用的 MySQL 客户端,API 与 PyMySQL 几乎一致- 生产环境必用连接池
create_pool,搭配pool_recycle防连接失效 - 修改 SQL 后别忘
commit()(autocommit=False是默认值) - 永远使用参数化查询
%s,防 SQL 注入 - 大结果集用
SSCursor流式读取 - 极致性能场景可考虑 asyncmy(C 扩展)
- 与
aiohttp/ FastAPI / asyncio 完美配合,构建高并发数据服务
四:ORM sqlalchemy(pymysql)
4.1 简介
在Python中最著名的ORM(Object Relationship Mapping)对象关系映射)框架是SQLAlchemy,类似于Java中的Hibernate, 在Java中Hibernate已经被淘汰多年了,原因是Hibernate属于重量级框架SQL是框架自动生成的不能手动写SQL来优化SQL语句。在Java中一般都使用MyBatis,自己写sql语句,然后映射到对象上。
SQLAlchemy只是一种ORM框架,它并不能直接操作数据库,直接操作数据库还需要通过pymysql模块来操作。
ORM最重要的映射有两个:一是表名和实体类的映射;另一个是表的字段和类的属性之间的映射;
ORM基类:
获取ORM基类通过 Base = declarative_base() 来
获取,所有实体类都要继承Base类。
类与表的映射通过__tablename__属性来指定,如 __tablename__="user"
属性与字段的映射通过Column类来实现。可以指定列的数据类型、是否允许为空、默认值、是否为主键、是否自增、是否唯一、 注释等
sqlalchemy针对MySQL方言提供了专门的数据类型sqlalchemy.dialects.mysql, 如果使用的数据库是MySQL建议使用这套数据类型,这套数据类型和MySQL的数据类型一一对应。
4.2 语法
1. query()函数中参数是要查询的字段,如果要查询所有字段只需要将实体类型作为参数,如果要查询多个字段就通过
# select * from user
session.query(User).all()
# select id, username from user
session.query(User.id, User.username).all()
2. 关联关系:relationship('引用的类名',backref='反向关联的属性')
# 一篇文章的作者对应一个用户
author = relationship('User',backref='articles')
# 一个作者有多篇文章
articles = relationship("Article")
3. 删除修改都行将对象先查询出来,然后再操作,这样就是执行2次,还不如直接写sql执行1次。
4. 实际开发中表一般是不设置外键的,使用ORM就必须设置外键了。
4.3 示例
1. 安装依赖
pip install pymysql
pip install sqlalchemy
2. 建表和初始化数据
#!/usr/bin/env python
# -*- coding:utf-8 -*-
author = 'suncity'
import sqlalchemy
from sqlalchemy import create_engine
from sqlalchemy.ext.declarative import declarative_base
from sqlalchemy import Table, Column, types, Enum, ForeignKey, and_, or_
from sqlalchemy.dialects.mysql import VARCHAR, TEXT, BIGINT, INTEGER, SMALLINT, TINYINT, DECIMAL, FLOAT, DOUBLE, DATETIME, \
TIMESTAMP
from sqlalchemy.orm import sessionmaker, relationship, backref
from sqlalchemy.sql import func
from enum import Enum, unique
from random import randint
# 初始化数据库连接
HOST = '127.0.0.1'
PORT = 3306
USERNAME = 'root'
PASSWORD = 'root123'
DATABASE = 'test'
DATASOURCE_URL = "mysql+pymysql://{username}:{password}@{host}:{port}/" \
"{db}?charset=utf8".format(username=USERNAME, password=PASSWORD, host=HOST, port=PORT, db=DATABASE)
engine = create_engine(DATASOURCE_URL, encoding='utf-8', echo=True)
# 判断有没有链接成功
conn = engine.connect()
result = conn.execute("select 1").fetchone()
print(result)
# ORM基类
Base = declarative_base()
# 多对多配置
user_teacher = Table("user_teacher",
Base.metadata,
Column("user_id", BIGINT(unsigned=True), ForeignKey("user.id"), primary_key=True),
Column("teacher_id", BIGINT(unsigned=True), ForeignKey("teacher.id"), primary_key=True))
@unique
class StatusEnum(Enum):
CLOSE = 0
OPEN = 1
class User(Base):
tablename = "user"
id = Column(BIGINT(unsigned=True), primary_key=True, autoincrement=True)
username = Column(VARCHAR(15), unique=True, nullable=False, comment="用户名")
# 数字类型的默认值需要写成字符串
gender = Column(TINYINT(unsigned=True), server_default="0", comment="性别(0: 女 1:男)")
age = Column(TINYINT(unsigned=True), server_default="0", comment="年龄")
# name用于指定数据库中的字段名,如不指定和属性名保持一致
# 数据库命名规范一般是小写,每个单词用下划线分隔,如果Python属性也采用同样的命名规则就不需要显式指定列明。
# 只有当属性名和列明不一样时才显式指定
price = Column(DECIMAL(10, 2), name="amount", nullable=True)
# 枚举实际中不建议使用(只有高版本的MySQL才支持enum类型),这里只是演示一下,实际情况下一般使用tinyint
status = Column(types.Enum(StatusEnum))
# DateTime类型的默认值使用func.now()
create_time = Column(DATETIME, server_default=func.now(), comment="注册时间")
# onupdate 当更新数据时会自动修改值
update_time = Column(TIMESTAMP, onupdate=func.now())
# 正向一对一(关闭一对多就是一对一了)
# cascade=save-update默认值,当添加和更新的时候相关联的数据也会级联更新
detail = relationship("UserDetail", uselist=False, cascade="save-update,delete")
# 一对多
tags = relationship("Tag")
# 多对多, secondary用于指定中间表
teachers = relationship("Teacher", secondary=user_teacher)
def __init__(self, username, gender):
self.username = username
self.gender = gender
def __str__(self):
return ",\n".join([str(item) for item in self.__dict__.items()])
class UserDetail(Base):
tablename = "user_detail"
id = Column(BIGINT(unsigned=True), primary_key=True, autoincrement=True)
city = Column(VARCHAR(15), nullable=False, comment="地址")
description = Column(TEXT, nullable=True, comment="介绍")
user_id = Column(BIGINT(unsigned=True), ForeignKey("user.id"))
# 反向一对一
user = relationship("User", uselist=False)
def __init__(self, city, description, user_id = None):
self.city = city
self.description = description
self.user_id = user_id
class Tag(Base):
tablename = "tag"
id = Column(BIGINT(unsigned=True), primary_key=True, autoincrement=True)
tag = Column(VARCHAR(15), nullable=False, comment="标签")
user_id = Column(BIGINT(unsigned=True), ForeignKey("user.id"))
def __init__(self, tag, user_id):
self.tag = tag
self.user_id = user_id
def __str__(self):
return ",\n".join([str(item) for item in self.__dict__.items()])
class Teacher(Base):
tablename = "teacher"
id = Column(BIGINT(unsigned=True), primary_key=True, autoincrement=True)
name = Column(VARCHAR(15), nullable=False, comment="姓名")
def __init__(self, name):
self.name = name
# 删除所有表结构(一般不用,这里只是为了每次都初始化数据)
Base.metadata.drop_all(engine)
# 创建表结构(如果映射已经发生了改变不会重复创建)
Base.metadata.create_all(engine)
# 创建session实例
session = sessionmaker(engine)()
users = [("xiaoming", "shanghai", "description1", ["活泼", "热情", "美丽"], [1, 2, 3]),
("xiaohong", "beijing", "description2", ["开朗"], [1, 2]),
("wangwu", "hangzhou", "description3", ["机敏"], [2, 3]),
("suncity", "suzhou", "description4", ["健谈"], [1, 3])
]
teacher_zhang = Teacher("张老师")
teacher_wang = Teacher("王老师")
teacher_li = Teacher("李老师")
teacher_list = [teacher_zhang, teacher_wang, teacher_li]
# 批量插入
session.add_all(teacher_list)
for username, city, description, tags, teachers in users:
# 执行原生SQL
result = session.execute("insert into user(username, gender, age, status)values(:username, :gender, :age, :status)",
params={"username": username, "gender": randint(0, 1), "age": randint(0, 150),
"status": StatusEnum.OPEN.value})
session.commit()
user_id = result.lastrowid
session.add(UserDetail(city, description, user_id))
for tag in tags:
session.add(Tag(tag, user_id))
for teacher_id in teachers:
session.execute("insert into user_teacher(user_id, teacher_id)values(:userId, :teacherId)",
params={"userId": user_id, "teacherId": teacher_id})
else:
session.commit()
3. execute
# 执行原生SQL
result = session.execute("select * from user where id > 1 order by create_time desc limit 1, 10").fetchall()
print(result)
4. filter
# select * from user
session.query(User).all()
# select * from user where id = 1
session.query(User).get(1)
# select * from user where username ='suncity' limit 1
session.query(User).filter_by(username="suncity").first()
# select * from user where username !='suncity' limit 1
session.query(User).filter(User.username!="suncity").first()
# select * from user where username like 'xiao%'
session.query(User).filter(User.username.like("xiao%")).all()
# ilike忽略大小写,其实是统一转为小写
# select * from user where lower(user.username) LIKE lower("%xiao")
session.query(User).filter(User.username.ilike("xiao%")).all()
# select * from user where id in(1, 2)
session.query(User).filter(User.id.in_([1, 2])).all()
# select * from user where id not in(1, 2)
session.query(User).filter(User.id.notin_([1, 2])).all()
# select * from user where update_time IS NULL
session.query(User).filter(User.update_time == None).all()
# select * from user where update_time IS NOT NULL
session.query(User).filter(User.update_time != None).all()
# filter_by指定属性时不需要指定类名, query中可以指定要查询的列, label给列起别名就是SQL中的as
# select id, username, create_time as join_time from user where username = "suncity"
session.query(User.id, User.username, User.create_time.label("join_time")).filter_by(username="suncity").all()
# and方式一:多个filter使用and拼接
# select * from user where id in(1, 2) and username like 'xiao%'
session.query(User).filter(User.id.in_([1, 2])).filter(User.username.like('xiao%')).all()
# and方式二:将多个条件写在一个filter中
# select * from user where id = 1 and username = 'xiaoming'
session.query(User).filter(User.id == 1, User.username == 'xiaoming').first()
# and方式三:使用and_来指定
# select * from user where id = 1 and username = 'xiaoming'
session.query(User).filter(and_(User.id == 1, User.username == 'xiaoming')).first()
# or_
# select * from user where id = 1 or username = 'xiaohong'
session.query(User).filter(or_(User.id == 1, User.username == 'xiaohong')).all()
# select * from user where id = 1 and (username='xiaoming' or gender = 0)
session.query(User).filter(User.id == 1, or_(User.username == 'xiaoming', User.gender == 0)).all()
5. 聚合函数
# func类有常用的聚合函数,如:count()、avg()、max()、min()、sum()
# SELECT count(user.id) FROM user LIMIT 1
session.query(func.count(User.id)).first()
# SELECT avg(user.age) FROM user LIMIT 1
session.query(func.avg(User.age)).first()
# SELECT max(user.age) FROM user LIMIT 1
session.query(func.max(User.age)).first()
# SELECT min(user.age) FROM user LIMIT 1
session.query(func.max(User.age)).first()
# SELECT sum(user.age) FROM user LIMIT 1
session.query(func.sum(User.age)).first()
6. 排序
# SELECT * FROM user ORDER BY create_time
session.query(User).order_by(User.create_time).all()
# SELECT * FROM user ORDER BY create_time DESC
session.query(User).order_by(User.create_time.desc()).all()
# SELECT * FROM user ORDER BY create_time ASC
session.query(User).order_by(User.create_time.asc()).all()
7. 分组
# 分组
# SELECT gender, count(user.id) AS count_1 FROM user GROUP BY gender
session.query(User.gender, func.count(User.id)).group_by(User.gender).all()
# SELECT gender, count(user.id) AS count_1 FROM user GROUP BY gender HAVING gender > 2
session.query(User.gender, func.count(User.id)).group_by(User.gender).having(User.gender > 2).all()
8. 分页
# select * from user limit 3
session.query(User).limit(3).all()
# select * from user limit 2, 18446744073709551615
session.query(User).offset(2).all()
# select * from user limit 2, 3
session.query(User).offset(2).limit(3).all()
# select * from user limit 1, 2
session.query(User).slice(1, 3).all()
# select * from user limit 1, 2
session.query(User)[1:3]
9. 关联查询和子查询
# 关联查询
# SELECT user.username, user_detail.city FROM user INNER JOIN user_detail ON user.id = user_detail.user_id
session.query(User.username, UserDetail.city).join(UserDetail, User.id == UserDetail.user_id).all()
# SELECT user.username, user_detail.city FROM user LEFT OUTER JOIN user_detail ON user.id = user_detail.user_id
session.query(User.username, UserDetail.city).outerjoin(UserDetail, User.id == UserDetail.user_id).all()
# 子查询subquery
# SELECT
# user.id AS user_id, user.username AS user_username
# FROM user, (SELECT user_detail.user_id AS user_id FROM user_detail WHERE user_detail.city = 'shanghai') AS anon_1
# WHERE user.id IN (anon_1.user_id)
sq = session.query(UserDetail.user_id).filter(UserDetail.city == "shanghai").subquery()
session.query(User.id, User.username).filter(User.id.in_(sq.c)).all()
# 所有的子查询都转换为多表连接查询
# SELECT user.id, user.username, anon_1.city
# FROM user, (SELECT user_detail.city FROM user_detail, user WHERE user_detail.id = user.id) AS anon_1
sq = session.query(UserDetail.city).filter(UserDetail.id == User.id).subquery()
session.query(User.id, User.username, sq.c.city).all()
10. 懒加载
# 懒加载 backref(lazy="select")
# select * from user where id = 1
user = session.query(User).get(1)
# 当获取user.detail.city时会执行 select * from user_detail where user_id = 1
print(user.detail.city)
# select * from user_detail where user_id = 1
user_detail = session.query(UserDetail).filter_by(user_id=1).first()
# select * from user where id = 1
print(user_detail.user.username)
user3 = session.query(User).filter_by(id=1).first()
print(str([item.tag for item in user3.tags]))
# <class 'sqlalchemy.orm.collections.InstrumentedList'>
print(type(user3.tags))
11. 级联添加
# 级联添加
# insert into user(username, gender, amount, status, update_time)values("admin", 0, null, null, null)
# insert into user_detail(city, description, user_id)values('shenzhen', 'admin description', 5)
# INSERT INTO user_teacher (user_id, teacher_id) VALUES(5, 1)(5, 3)
admin_user = User("admin", 0)
detail = UserDetail(city="shenzhen", description="admin description")
admin_user.detail = detail
admin_user.teachers = [teacher_zhang, teacher_li]
session.add(admin_user)
session.commit()
# 反向级联添加
# insert into user(username, gender, amount, status, update_time)values("root", 0, null, null, null)
# insert into user_detail(city, description, user_id)values('wuhan', 'root description', 5)
rootdetial = UserDetail(city="wuhan", description="root description")
rootdetial.user = User("root", 1)
session.add(rootdetial)
session.commit()
12. update
# 更新: 一般先查询出来,然后修改属性的值然后提交即可完成修改
obj = session.query(User).filter_by(id=1).first()
obj.gender = 0
session.commit()
13. delete
# 删除主表时如果有其它表外键引用主表主键orm会先处理掉使用外键的表(删除记录、设置外键为null),然后最后再处理主表
# 多对多:DELETE FROM user_teacher WHERE user_teacher.user_id = 1
# 一对多:会将外键设置为null, UPDATE tag SET user_id=null WHERE tag.id in (1, 2, 3)
# 一对一:会将外键设置为null, UPDATE user_detail SET user_id=null WHERE user_detail.id = 1
# 最后才会删除主表: DELETE FROM user WHERE user.id = 1
obj = session.query(User).get(1)
session.delete(obj)
session.commit()
五:ORM sqlalchemy(aiomysql)
六:pydantic
pydantic使用场景:
- 和sqlalchemy结合使用,用于定义数据库表结果对应的entity
- 用于配置类
class定义通常还使用pydantic进行设置一些默认值或者参数合法检查。
Field:
- default:默认值,
default=[ ]:所有模型实例会共享同一个列表对象,修改一个实例的 会影响其他实例(经典的可变默认值陷阱),注意列表要用default_factory。 - default_factory:动态生成,接收任何无参数的 callable,list表示每次自动执行 list(),生成全新空列表。常用于dict、set、datetime.now、uuid4
| 参数 | 作用 | 示例 |
|---|---|---|
... (省略号) |
必填字段 | message: str = Field(...) |
default |
默认值 | session_id: str = Field(default=None) |
default_factory |
默认工厂 | decisions: list = Field(default_factory=list) |
description |
字段描述 | description="用户消息" |
ge / le |
数值范围 | confidence: float = Field(ge=0.0, le=1.0) |
min_length / max_length |
字符串长度 | message: str = Field(max_length=8000) |
class MeetingSummary(BaseModel):
"""会议纪要"""
title:str=Field(...,description="会议标题")
date: str = Field(...,description="会议日期")
attendees:List[str]= Field(...,description="参会人员")
summary: str=Field(...,description="会议摘")
decisions: List[str]=Field(default_factory=list, description="决议事项")
action_items: List[MeetingAction]= Field(default_factory=list, description="行动顶")
notes:str=Field(default=None,description="备注")
from enum import Enum
from pydantic import BaseModel, Field,ValidationError
class Intent(str, Enum):
"""用户意图枚举"""
CHAT = "chat" # 闲聊问答
MEETING = "meeting" #会议纪要
UNKNOWN = "unknown" # 未知
class RouterDecision(BaseModel):
"""路由决策结果"""
intent: Intent # 识别到的意图
confidence:float = Field(ge=0.0,le=1.0) # 置信度0-1
reasoning: str # 推理过程
extracted_info:dict = None # 提取的信息
decision =RouterDecision(
intent=Intent.CHAT,
confidence-0.95,
reasoning="用户在进行日常对话"
)
ConfigDict
from pydantic import BaseModel, ConfigDict, Field, ValidationError
class User:
def __init__(self, username: str, age: int, password: str):
self.username = username
self.age = age
self.password = password
def __repr__(self) -> str:
return (
f"User(username={self.username!r}, "
f"age={self.age!r}, "
f"password={self.password!r})"
)
class UserModel(BaseModel):
username: str = Field(..., description="姓名")
age: int = Field(..., gt=0, description="年龄")
password: str = Field(..., min_length=6, description="密码")
# from_attributes=True 允许 model_validate 从普通对象属性中读取字段值。object -> model
model_config = ConfigDict(from_attributes=True)
# 0. dict -> object
user_dict = {
"username": "melong",
"age": "18",
"password": "123456",
}
print(type(User(**user_dict)))
# 1. model_validate: dict -> model(UserModel)
user_model = UserModel.model_validate(user_dict)
# dict -> UserModel: username='melong' age=18 password='123456'
print(type(user_model))
print(user_model)
# 2. model -> object
user_dict2 = user_model.model_dump()
# <class 'dict'>
print(type(user_dict2))
# {'username': 'melong', 'age': 18, 'password': '123456'}
print(user_dict2)
user = User(**user_dict2)
print(type(user))
# 3. user -> model
user_model = UserModel.model_validate(user)
print(type(user_model))
# 4. 校验失败:字段类型或规则不满足时会抛出 ValidationError。
try:
UserModel.model_validate(
{
"username": "melong",
"age": -1,
"password": "123",
}
)
except ValidationError as exc:
print("\nvalidate error:")
print(exc)
AtomGit 是由开放原子开源基金会联合 CSDN 等生态伙伴共同推出的新一代开源与人工智能协作平台。平台坚持“开放、中立、公益”的理念,把代码托管、模型共享、数据集托管、智能体开发体验和算力服务整合在一起,为开发者提供从开发、训练到部署的一站式体验。
更多推荐



所有评论(0)