from sqlalchemy import create_engine, text
from sqlalchemy.orm import sessionmaker

# 配置数据库连接(示例为PostgreSQL->MySQL)
SRC_DB_URL = 'postgresql://user:pass@source_host:5432/source_db'
DST_DB_URL = 'mysql+pymysql://user:pass@dest_host:3306/dest_db'

# 创建引擎和会话
src_engine = create_engine(SRC_DB_URL)
dst_engine = create_engine(DST_DB_URL)
SrcSession = sessionmaker(bind=src_engine)
DstSession = sessionmaker(bind=dst_engine)

def migrate_with_raw_sql():
    with SrcSession() as src_session, DstSession() as dst_session:
        # 从源数据库查询数据(使用原生SQL)
        query = text("SELECT id, name, age FROM users WHERE updated_at > :last_update")
        src_data = src_session.execute(query, {"last_update": "2025-01-01"}).fetchall()
        
        # 构建批量更新SQL语句
        update_sql = text("""
            UPDATE users 
            SET name = :name, age = :age 
            WHERE id = :id
        """)
        
        # 执行批量更新
        for row in src_data:
            dst_session.execute(
                update_sql, 
                {"id": row.id, "name": row.name, "age": row.age}
            )
        
        dst_session.commit()
        print(f"通过原生SQL语句完成 {len(src_data)} 条记录更新")

if __name__ == '__main__':
    migrate_with_raw_sql()
 

Logo

DAMO开发者矩阵,由阿里巴巴达摩院和中国互联网协会联合发起,致力于探讨最前沿的技术趋势与应用成果,搭建高质量的交流与分享平台,推动技术创新与产业应用链接,围绕“人工智能与新型计算”构建开放共享的开发者生态。

更多推荐