Postgres
更多
五. 使用
5.1 执行原生SQL
import threading
from sqlalchemy import create_engine
engine = create_engine( "mysql+pymysql://zff:zff123@127.0.0.1:3306/zff" , max_overflow= 0 , pool_size= 5 )
def task ( arg) :
conn = engine. contextual_connect( )
with conn:
cur = conn. execute(
"select * from USER "
)
result = cur. fetchall( )
print ( result)
for i in range ( 20 ) :
t = threading. Thread( target= task, args= ( i, ) )
t. start( )
import datetime
from sqlalchemy import create_engine
from sqlalchemy. ext. declarative import declarative_base
from sqlalchemy import Column, Integer, String, Text, ForeignKey, DateTime, UniqueConstraint, Index
Base = declarative_base( )
class Users ( Base) :
__tablename__ = 'users'
id = Column( Integer, primary_key= True )
name = Column( String( 32 ) , index= True , nullable= False )
email = Column( String( 32 ) , unique= True )
ctime = Column( DateTime, default= datetime. datetime. now)
__table_args__ = (
)
def init_db ( ) :
"""
根据类创建数据库表
:return:
"""
engine = create_engine(
"mysql+pymysql://zff:zff123@127.0.0.1:3306/zff?charset=utf8" ,
max_overflow= 0 ,
pool_size= 5 ,
pool_timeout= 30 ,
pool_recycle= - 1
)
Base. metadata. create_all( engine)
def drop_db ( ) :
"""
根据类删除数据库表
:return:
"""
engine = create_engine(
"mysql+pymysql://zff:zff123@127.0.0.1:3306/zff?charset=utf8" ,
max_overflow= 0 ,
pool_size= 5 ,
pool_timeout= 30 ,
pool_recycle= - 1
)
Base. metadata. drop_all( engine)
if __name__ == '__main__' :
drop_db( )
init_db( )
from sqlalchemy. orm import sessionmaker
from sqlalchemy import create_engine
from models import *
engine = create_engine( "mysql+pymysql://zff:zff123@127.0.0.1:3306/zff?charset=utf8" , max_overflow= 0 , pool_size= 5 )
Session = sessionmaker( bind= engine)
session = Session( )
obj1 = Users( name= "tome" , age= 19 , email= "tome163@163.com" )
session. add( obj1)
session. commit( )
session. close( )
from sqlalchemy. orm import sessionmaker
from sqlalchemy import create_engine
from sqlalchemy. sql import text
from models import *
engine = create_engine( "mysql+pymysql://zff:zff123@127.0.0.1:3306/zff?charset=utf8" , max_overflow= 0 , pool_size= 5 )
Session = sessionmaker( bind= engine)
session = Session( )
"""
obj1 = Users(name="jack", age=19, email="jak163@163.com")
session.add(obj1)
session.add_all([
Users(name="wang", age=19, email="wang163@163.com"),
Users(name="lucy", age=19, email="lucy@163.com"),
Hosts(name="jav-pingtai03br-p002.gru1.blue.net"),
])
session.commit()
"""
"""
session.query(Users).filter(Users.id > 2).delete()
session.commit()
"""
"""
session.query(Users).filter(Users.id > 0).update({"name" : "shuke"})
session.query(Users).filter(Users.id > 0).update({Users.name: Users.name + "163"}, synchronize_session=False)
session.query(Users).filter(Users.id > 0).update({"age": Users.age + 1}, synchronize_session="evaluate")
session.commit()
# sqlalchemy 利用 session 执行 delete 时有一个 synchronize_session 参数用来说明 session 删除对象时需要执行的策略,共三个选项:
1. False
不同步 session,如果被删除的 objects 已经在 session 中存在,在 session commit 或者 expire_all 之前,这些被删除的对象都存在 session 中。
不同步可能会导致获取被删除 objects 时出错。
# 2. fetch
删除之前从 db 中匹配被删除的对象并保存在 session 中,然后再从 session 中删除,这样做是为了让 session 的对象管理 identity_map 得知被删除的对象究竟是哪些以便更新引用关系。
# 3. evaluate
# 默认值。根据当前的 query criteria 扫描 session 中的 objects,如果不能正确执行则抛出错误,这句话也可以理解为,如果 session 中原本就没有这些被删除的 objects,扫描当然不会发生匹配,相当于匹配未正确执行。
注意这里报错只会在特定 query criteria 时报错,比如 in 操作。
session.query(Users).filter(Users.id.in_([1,2,3])).delete()
sqlalchemy.exc.InvalidRequestError: Could not evaluate current criteria in Python. Specify 'fetch' or False for the synchronize_session parameter.
"""
"""
r1 = session.query(Users).all()
r2 = session.query(Users.name.label('username'), Users.age).all() # 别名
r3 = session.query(Users).filter(Users.name == "shuke").all()
r4 = session.query(Users).filter_by(name='shuke').all()
r5 = session.query(Users).filter_by(name='shuke').first()
r6 = session.query(Users).filter(text("id<:value and name=:name")).params(value=2, name='shuke').order_by(Users.id).all()
r7 = session.query(Users).from_statement(text("SELECT * FROM users where name=:name")).params(name='shuke').all()
"""
"""
filter_by用于简单的列名查询,如:
db.users.filter_by(name='Joe')
filter对于上面的代码可以这样写:
db.users.filter(db.users.name == 'Joe')
对于复杂的查询使用filter,如:
db.users.filter(or_(db.users.name == 'Ryan', db.users.country == 'England'))
注意: filter_by使用的是赋值 =, 而filter使用的是判断 ==
另外:查询时使用like这样写: items = session.query.filter(Users.name == current_user, Users.title.like('%' + keyword + '%')).all()
"""
session. close( )
favor = relationship("Favor", backref='pers')
ret3 = session. query ( Person. name, Favor. caption) . join ( Favor, isouter= True) . filter ( Favor. caption == 'blue' ) . all ( )
Person表里,写了backref='pers',就相当于在favor表里加了个字段pers,实际是不存在的
class Person ( Base) :
__tablename__ = 'person'
nid = Column( Integer, primary_key= True )
name = Column( String( 32 ) , index= True , nullable= True )
favor_id = Column( Integer, ForeignKey( "favor.nid" ) )
favor = relationship( "Favor" , backref= 'pers' )
class Favor ( Base) :
__tablename__ = 'favor'
nid = Column( Integer, primary_key= True )
caption = Column( String( 50 ) , default= 'red' , unique= True )
def __repr__ ( self) :
return "%s-%s" % ( self. nid, self. caption)
你可以直接通过Favor对象的pers字段找到跟这个颜色关联的所有person,在数据库里没有真实的字段对应的,只是帮你生成sql语句而已。
class Person ( Base) :
__tablename__ = 'person'
nid = Column( Integer, primary_key= True )
name = Column( String( 32 ) , index= True , nullable= True )
favor_id = Column( Integer, ForeignKey( "favor.nid" ) )
favor = relationship( "Favor" , backref= 'pers' )
Person对Favor 是多对一的关系,foreignkey加在了多的那端(Person表)
Person对象.favor.favor的字段:叫做正向查找
Favor对象.pers.person的字段:反向查找
M2M(基于relationship的m2m关系)
import time
import threading
from sqlalchemy. ext. declarative import declarative_base
from sqlalchemy import Column, Integer, String, ForeignKey, UniqueConstraint, Index
from sqlalchemy. orm import sessionmaker, relationship
from sqlalchemy import create_engine
from db import Users
engine = create_engine( "mysql+pymysql://root:123@127.0.0.1:3306/s6" , max_overflow= 0 , pool_size= 5 )
Session = sessionmaker( bind= engine)
def task ( arg) :
session = Session( )
obj1 = Users( name= "shuke" )
session. add( obj1)
session. commit( )
for i in range ( 10 ) :
t = threading. Thread( target= task, args= ( i, ) )
t. start( )
基于scoped_session使得线程安全
基于ThreadLocal实现
from sqlalchemy. orm import sessionmaker
from sqlalchemy import create_engine
from sqlalchemy. orm import scoped_session
from models import Users
engine = create_engine( "mysql+pymysql://root:123@127.0.0.1:3306/s6" , max_overflow= 0 , pool_size= 5 )
Session = sessionmaker( bind= engine)
"""
# 线程安全,基于本地线程实现每个线程用同一个session
# 特殊的:scoped_session中有原来方法的Session中的一下方法:
public_methods = (
'__contains__', '__iter__', 'add', 'add_all', 'begin', 'begin_nested',
'close', 'commit', 'connection', 'delete', 'execute', 'expire',
'expire_all', 'expunge', 'expunge_all', 'flush', 'get_bind',
'is_modified', 'bulk_save_objects', 'bulk_insert_mappings',
'bulk_update_mappings',
'merge', 'query', 'refresh', 'rollback',
'scalar'
)
"""
session = scoped_session( Session)
obj1 = Users( name= "shuke" )
session. add( obj1)
session. commit( )
session. close( )
参考资料:
Flask-SQLAlchemy-武沛齐-博客园
使用flask-sqlalchemy玩转MySQL | Wing's Tech Space
Flask-Migrate的使用 | Wing's Tech Space