You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

MySQL至PostgreSQL数据迁移技术咨询:代码校验、类型处理及FDW方案

Magento MySQL to PostgreSQL Staging Table: Answers to Your Technical Questions

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, and conn_string_psql are missing configuration values—you'll need to define these (e.g., from environment variables or a config file).
  • Broken psql_command Function: The function references cur_msql but takes msql as a parameter, plus it uses an undefined command variable instead of psql_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)s should be %(created_at)s.
  • Flawed Transaction Handling: You commit after creating tables, but the final commit is in the finally block—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 TypePostgreSQL Equivalent
    DATETIMETIMESTAMP WITHOUT TIME ZONE
    INT(11)BIGINT (if values exceed INT limits)
    DECIMAL(M,N)NUMERIC(M,N)
    VARCHAR(255)TEXT or VARCHAR(255)
  • On-the-Fly Conversion:

    • In MySQL Queries: Use CAST to 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 datetime objects, which psycopg2 handles natively).
  • 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

  1. 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 multicorn
    
  2. Enable Multicorn in PostgreSQL:

    CREATE EXTENSION IF NOT EXISTS multicorn;
    
  3. 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'
    );
    
  4. 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');
    
  5. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.13 09:03:22