MySQL至PostgreSQL数据迁移技术咨询:代码校验、类型处理及FDW方案
Hey Sandeep, let's dive into your questions about migrating Magento data to a PostgreSQL staging table. I'll break down each query with practical fixes and alternatives tailored to your use case.
1. Code Validation & Alternative Implementations
First, let's go over the issues in your current code and better approaches:
Key Issues in Your Current Code
- Undefined Variables:
host_mysql,user_mysql,pswd_mysql,dbna_mysql, andconn_string_psqlare missing configuration values—you'll need to define these (e.g., from environment variables or a config file). - Broken
psql_commandFunction: The function referencescur_msqlbut takesmsqlas a parameter, plus it uses an undefinedcommandvariable instead ofpsql_command. It also inserts rows one at a time, which is slow for large datasets. - Syntax Error: The second INSERT statement has a typo:
% . (created_at)sshould be%(created_at)s. - Flawed Transaction Handling: You commit after creating tables, but the final commit is in the
finallyblock—this would commit partial data if an error occurs during inserts.
Fixed & Optimized Code Example
Here's a revised version with batch inserts and proper error handling:
import psycopg2 import mysql.connector import sys from psycopg2 import sql # Configuration - fill these in or load from env vars MYSQL_CONFIG = { "host": "your_mysql_host", "user": "your_mysql_user", "passwd": "your_mysql_pass", "db": "magento_db" } PSQL_CONN_STRING = "dbname=postgres user=postgres host=your_psql_host password=your_psql_pass" def batch_insert(msql_cursor, psql_cursor, select_query, insert_query, batch_size=1000): msql_cursor.execute(select_query) while True: rows = msql_cursor.fetchmany(batch_size) if not rows: break try: psql_cursor.executemany(insert_query, rows) psql_cursor.connection.commit() print(f"Inserted {len(rows)} rows successfully") except (psycopg2.Error, mysql.connector.Error) as e: print(f"Batch insert failed: {e}") psql_cursor.connection.rollback() sys.exit("Aborting due to insert error") def db_fetch(): # MySQL Connection try: cnx_msql = mysql.connector.connect(**MYSQL_CONFIG) cur_msql = cnx_msql.cursor(dictionary=True) except mysql.connector.Error as e: print(f"MySQL Connection Failed: {e.msg}") sys.exit(1) # PostgreSQL Connection try: cnx_psql = psycopg2.connect(PSQL_CONN_STRING) cur_psql = cnx_psql.cursor() except psycopg2.Error as e: print(f"PostgreSQL Connection Failed: {e}") sys.exit(1) try: # Create staging schema and tables cur_psql.execute(sql.SQL("CREATE SCHEMA IF NOT EXISTS staging AUTHORIZATION postgres;")) cur_psql.execute(""" CREATE TABLE IF NOT EXISTS staging.sales_flat_quote ( entity_id BIGINT, store_id BIGINT, customer_email TEXT, customer_firstname TEXT, customer_middlename TEXT, customer_lastname TEXT, customer_is_guest BIGINT, customer_group_id BIGINT, created_at TIMESTAMP WITHOUT TIME ZONE, updated_at TIMESTAMP WITHOUT TIME ZONE, is_active BIGINT, items_count BIGINT, items_qty BIGINT, base_currency_code TEXT, grand_total NUMERIC(12,4), base_to_global_rate NUMERIC(12,4), base_subtotal NUMERIC(12,4), base_subtotal_with_discount NUMERIC(12,4) ); """) cur_psql.execute(""" CREATE TABLE IF NOT EXISTS staging.sales_flat_quote_item ( store_id INTEGER, row_total NUMERIC, updated_at TIMESTAMP WITHOUT TIME ZONE, qty NUMERIC, sku CHARACTER VARYING, free_shipping INTEGER, quote_id INTEGER, price NUMERIC, no_discount INTEGER, item_id INTEGER, product_type CHARACTER VARYING, base_tax_amount NUMERIC, product_id INTEGER, name CHARACTER VARYING, created_at TIMESTAMP WITHOUT TIME ZONE ); """) cnx_psql.commit() print("Staging schema and tables created successfully") # Define transfer commands transfer_commands = [ ( "SELECT entity_id, store_id, customer_email, customer_firstname, customer_middlename, " "customer_lastname, customer_is_guest, customer_group_id, created_at, updated_at, " "is_active, items_count, items_qty, base_currency_code, grand_total, base_to_global_rate, " "base_subtotal, base_subtotal_with_discount FROM sales_flat_quote WHERE is_active=1;", "INSERT INTO staging.sales_flat_quote VALUES (%(entity_id)s, %(store_id)s, %(customer_email)s, " "%(customer_firstname)s, %(customer_middlename)s, %(customer_lastname)s, %(customer_is_guest)s, " "%(customer_group_id)s, %(created_at)s, %(updated_at)s, %(is_active)s, %(items_count)s, %(items_qty)s, " "%(base_currency_code)s, %(grand_total)s, %(base_to_global_rate)s, %(base_subtotal)s, %(base_subtotal_with_discount)s);" ), ( "SELECT store_id, row_total, updated_at, qty, sku, free_shipping, quote_id, price, " "no_discount, item_id, product_type, base_tax_amount, product_id, name, created_at " "FROM sales_flat_quote_item;", "INSERT INTO staging.sales_flat_quote_item VALUES (%(store_id)s, %(row_total)s, %(updated_at)s, " "%(qty)s, %(sku)s, %(free_shipping)s, %(quote_id)s, %(price)s, %(no_discount)s, %(item_id)s, " "%(product_type)s, %(base_tax_amount)s, %(product_id)s, %(name)s, %(created_at)s);" ) ] # Execute batch transfers for select_q, insert_q in transfer_commands: batch_insert(cur_msql, cur_psql, select_q, insert_q) except Exception as error: print(f"Data transfer failed: {error}") cnx_psql.rollback() finally: # Cleanup connections cur_msql.close() cnx_msql.close() cur_psql.close() cnx_psql.close() if __name__ == '__main__': db_fetch()
Alternative Implementations
- SQLAlchemy: Simplifies connection management and data type handling. You can use its core API to build queries and handle bulk inserts with less boilerplate.
- ETL Tools: For large-scale or recurring migrations, tools like Apache Airflow or PySpark offer built-in scheduling, error recovery, and parallel processing capabilities.
2. Handling Data Type Mismatches
Data type conflicts are common between MySQL and PostgreSQL—here's how to resolve them:
Pre-Map Types: Define your staging table columns to match MySQL's types appropriately:
MySQL Type PostgreSQL Equivalent DATETIMETIMESTAMP WITHOUT TIME ZONEINT(11)BIGINT(if values exceed INT limits)DECIMAL(M,N)NUMERIC(M,N)VARCHAR(255)TEXTorVARCHAR(255)On-the-Fly Conversion:
- In MySQL Queries: Use
CASTto convert types before fetching:SELECT CAST(created_at AS CHAR) AS created_at FROM sales_flat_quote; - In Python: Convert values during batch processing (e.g., parse datetime strings into Python
datetimeobjects, which psycopg2 handles natively).
- In MySQL Queries: Use
Error Logging: Add logic to catch type mismatch errors, log the problematic rows, and skip or retry them instead of aborting the entire transfer.
3. Using Foreign Data Wrappers (FDW) like Multicorn
Yes, FDW is perfect for your use case—especially since you only need 15 columns from each table. It eliminates the need for a Python middle layer, letting you query and migrate data directly within PostgreSQL.
Step-by-Step Setup
Install Dependencies:
# Install system packages (Ubuntu/Debian example) sudo apt-get install postgresql-server-dev-15 python3-dev gcc libmysqlclient-dev # Install Multicorn via pip pip3 install multicornEnable Multicorn in PostgreSQL:
CREATE EXTENSION IF NOT EXISTS multicorn;Create a Foreign Server for MySQL:
CREATE SERVER mysql_magento FOREIGN DATA WRAPPER multicorn OPTIONS ( wrapper 'multicorn.mysqlfdw.MySQLFdw', host 'your_mysql_host', port '3306', dbname 'magento_db', user 'your_mysql_user', password 'your_mysql_password' );Create Foreign Tables (Only Include Needed Columns):
-- Foreign table for sales_flat_quote (only 15 columns) CREATE FOREIGN TABLE staging.sales_flat_quote_fdw ( entity_id BIGINT, store_id BIGINT, customer_email TEXT, customer_firstname TEXT, customer_lastname TEXT, customer_is_guest BIGINT, created_at TIMESTAMP WITHOUT TIME ZONE, updated_at TIMESTAMP WITHOUT TIME ZONE, is_active BIGINT, items_count BIGINT, base_currency_code TEXT, grand_total NUMERIC(12,4), base_to_global_rate NUMERIC(12,4), base_subtotal NUMERIC(12,4), base_subtotal_with_discount NUMERIC(12,4) ) SERVER mysql_magento OPTIONS (table_name 'sales_flat_quote'); -- Foreign table for sales_flat_quote_item CREATE FOREIGN TABLE staging.sales_flat_quote_item_fdw ( store_id INTEGER, row_total NUMERIC, updated_at TIMESTAMP WITHOUT TIME ZONE, qty NUMERIC, sku CHARACTER VARYING, quote_id INTEGER, price NUMERIC, item_id INTEGER, product_type CHARACTER VARYING, product_id INTEGER, name CHARACTER VARYING, created_at TIMESTAMP WITHOUT TIME ZONE, free_shipping INTEGER, no_discount INTEGER, base_tax_amount NUMERIC ) SERVER mysql_magento OPTIONS (table_name 'sales_flat_quote_item');Migrate Data to Your Staging Tables:
-- Insert data from foreign table to local staging table INSERT INTO staging.sales_flat_quote SELECT * FROM staging.sales_flat_quote_fdw WHERE is_active = 1; INSERT INTO staging.sales_flat_quote_item SELECT * FROM staging.sales_flat_quote_item_fdw;
Why FDW Works for Your Use Case
- Minimal Code: No Python scripts needed—all operations are done in SQL.
- Efficiency: Direct database-to-database transfer is faster than going through an application layer.
- Flexibility: You can query the foreign tables directly (for ad-hoc analysis) or set up incremental syncs with triggers or scheduled jobs.
内容的提问来源于stack exchange,提问作者Sandeep

