如何用Python将MySQL数据以字符串形式迁移至MongoDB?
修正现有代码实现字符串格式迁移
你的代码当前是将整列查询结果(数组)直接存入MongoDB字段,因此呈现数组格式。要改为单条字符串对应存储,需要一次性查询所有目标列,再逐行提取字段值构造文档:
import mysql.connector import pymongo from pymongo import MongoClient # 配置信息 mysql_host="**" mysql_database="**" mysql_user="**" mysql_password="**" mongodb_host = "**" mongodb_dbname = "**" # 连接MySQL mysqldb = mysql.connector.connect( host=mysql_host, database=mysql_database, user=mysql_user, password=mysql_password ) if mysqldb.is_connected(): print("successfully Connected") # 获取游标,返回字典格式结果 mycursor = mysqldb.cursor(dictionary=True) # 一次性查询Name和location列,保证行数据对应 mycursor.execute("SELECT Name, location FROM MYSQLTABLE;") rows = mycursor.fetchall() # 连接MongoDB myclient = pymongo.MongoClient(mongodb_host) mydb = myclient[mongodb_dbname] mycol = mydb["CASES"] # 构造文档列表,每个字段存储字符串值 documents = [] for row in rows: doc = { "Name": row["Name"], "location": row["location"] } documents.append(doc) # 批量插入数据 if documents: x = mycol.insert_many(documents) print(f"插入了 {len(x.inserted_ids)} 条数据") # 关闭数据库连接 mysqldb.close()
关键修改说明:
- 合并查询:避免分开查询导致的行数据对应错误,同时减少数据库请求次数
- 逐行处理:提取每行的字段原始字符串值,作为MongoDB文档的独立字段
- 移除无效判断:删除未定义的
iddata判断逻辑
推荐其他简便Python迁移方法
1. 使用Pandas(最便捷的批量迁移方案)
Pandas可快速读取MySQL数据到DataFrame,直接写入MongoDB,支持自定义字段映射:
import pandas as pd import mysql.connector from pymongo import MongoClient # MySQL配置 mysql_config = { "host": "**", "database": "**", "user": "**", "password": "**" } # MongoDB配置 mongodb_host = "**" mongodb_dbname = "**" mongodb_col = "CASES" # 读取MySQL数据到DataFrame query = "SELECT Name, location FROM MYSQLTABLE;" df = pd.read_sql(query, mysql.connector.connect(**mysql_config)) # 可选:字段重命名映射 df.rename(columns={"Name": "username", "location": "address"}, inplace=True) # 写入MongoDB client = MongoClient(mongodb_host) db = client[mongodb_dbname] db[mongodb_col].insert_many(df.to_dict("records")) print(f"成功迁移 {len(df)} 条数据")
优势:代码极简,自动处理数据类型转换,支持数据清洗、字段重命名,适合大规模数据迁移。
2. 使用SQLAlchemy + PyMongo
如果需要更灵活的MySQL操作(比如复杂查询、ORM映射),可以用SQLAlchemy:
from sqlalchemy import create_engine from pymongo import MongoClient # SQLAlchemy连接MySQL engine = create_engine("mysql+mysqlconnector://user:password@host/database") # 查询并转换为字典列表 with engine.connect() as conn: result = conn.execute("SELECT Name, location FROM MYSQLTABLE;") rows = [dict(row) for row in result] # 插入MongoDB client = MongoClient(mongodb_host) db = client[mongodb_dbname] db["CASES"].insert_many(rows)
优势:支持ORM操作,适合需要复杂数据预处理的场景。
内容的提问来源于stack exchange,提问作者Safi
相关产品推荐
相关产品推荐

