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

基于Python的MySQL/Teradata与Snowflake/GCP分区对比步骤及检查项咨询

跨数据库分区对比全流程(Python实现)

一、前置准备

  • 安装依赖库:
    pip install mysql-connector-python teradatasql snowflake-connector-python google-cloud-bigquery pandas
    
  • 配置数据库连接:将各数据库的连接参数存入字典,敏感信息建议用环境变量传入
    # 示例:MySQL连接配置
    mysql_config = {
        'host': 'xxx',
        'user': 'xxx',
        'password': 'xxx',
        'database': 'xxx'
    }
    

二、元数据采集:获取源/目标端分区信息

1. MySQL 分区信息采集

通过INFORMATION_SCHEMA.PARTITIONS表获取分区的键、类型、边界、行数等信息:

import mysql.connector
from mysql.connector import Error

def get_mysql_partitions(db_config, table_name):
    try:
        conn = mysql.connector.connect(**db_config)
        cursor = conn.cursor(dictionary=True)
        query = """
            SELECT 
                PARTITION_NAME,
                PARTITION_METHOD,
                PARTITION_EXPRESSION,
                PARTITION_DESCRIPTION,
                TABLE_ROWS
            FROM INFORMATION_SCHEMA.PARTITIONS
            WHERE TABLE_SCHEMA = %s AND TABLE_NAME = %s
            ORDER BY PARTITION_ORDINAL_POSITION
        """
        cursor.execute(query, (db_config['database'], table_name))
        return cursor.fetchall()
    except Error as e:
        print(f"MySQL查询错误: {e}")
    finally:
        if conn.is_connected():
            cursor.close()
            conn.close()

2. Teradata 分区信息采集

结合DBC.TablesV和DBC.PartitioningColumnsV视图获取分区元数据:

import teradatasql

def get_teradata_partitions(db_config, table_name):
    try:
        with teradatasql.connect(**db_config) as conn:
            with conn.cursor() as cursor:
                query = """
                    SELECT 
                        t.TableName,
                        p.ColumnName AS PARTITION_EXPRESSION,
                        t.PartitionType AS PARTITION_METHOD,
                        t.PartitionCount,
                        s.CurrentPerm AS STORAGE_SIZE
                    FROM DBC.TablesV t
                    JOIN DBC.PartitioningColumnsV p 
                        ON t.DatabaseName = p.DatabaseName AND t.TableName = p.TableName
                    JOIN DBC.TableSizeV s 
                        ON t.DatabaseName = s.DatabaseName AND t.TableName = s.TableName
                    WHERE t.DatabaseName = ? AND t.TableName = ?
                """
                cursor.execute(query, (db_config['database'], table_name))
                partitions = []
                for row in cursor.fetchall():
                    partitions.append({
                        'PARTITION_NAME': f"p_{row[3]}",  # Teradata无直接分区名,需自定义
                        'PARTITION_METHOD': row[2],
                        'PARTITION_EXPRESSION': row[1],
                        'STORAGE_SIZE': row[4],
                        'TABLE_ROWS': None  # Teradata需单独查询分区行数
                    })
                return partitions
    except Exception as e:
        print(f"Teradata查询错误: {e}")

3. Snowflake 分区信息采集

通过INFORMATION_SCHEMA.PARTITIONS或SHOW PARTITIONS命令获取数据:

import snowflake.connector

def get_snowflake_partitions(sf_config, table_name):
    try:
        conn = snowflake.connector.connect(**sf_config)
        cursor = conn.cursor()
        # 拆分表的库、模式、表名
        db, schema, tbl = table_name.split('.') if '.' in table_name else (sf_config['database'], sf_config['schema'], table_name)
        query = f"""
            SELECT 
                PARTITION_NAME,
                PARTITION_TYPE,
                START_VALUE,
                END_VALUE,
                ROW_COUNT,
                BYTES
            FROM {db}.INFORMATION_SCHEMA.PARTITIONS
            WHERE TABLE_SCHEMA = '{schema}' AND TABLE_NAME = '{tbl}'
            ORDER BY PARTITION_ORDINAL_POSITION
        """
        cursor.execute(query)
        partitions = []
        for row in cursor.fetchall():
            partitions.append({
                'PARTITION_NAME': row[0],
                'PARTITION_METHOD': row[1],
                'PARTITION_BOUND': f"{row[2]}-{row[3]}",
                'TABLE_ROWS': row[4],
                'STORAGE_SIZE': row[5]
            })
        return partitions
    except Exception as e:
        print(f"Snowflake查询错误: {e}")
    finally:
        conn.close()

4. GCP BigQuery 分区信息采集

利用BigQuery客户端API获取分区的时间范围、行数、存储大小:

from google.cloud import bigquery

def get_bigquery_partitions(project_id, dataset_id, table_name):
    client = bigquery.Client(project=project_id)
    table_ref = client.dataset(dataset_id).table(table_name)
    table = client.get_table(table_ref)
    
    query = f"""
        SELECT 
            partition_id,
            total_rows,
            total_bytes,
            partition_time
        FROM `{project_id}.{dataset_id}.INFORMATION_SCHEMA.PARTITIONS`
        WHERE table_name = '{table_name}'
        ORDER BY partition_id
    """
    results = client.query(query).result()
    partitions = []
    for row in results:
        partitions.append({
            'PARTITION_NAME': row.partition_id,
            'PARTITION_METHOD': 'TIME' if table.time_partitioning else 'RANGE',
            'PARTITION_EXPRESSION': table.time_partitioning.field if table.time_partitioning else None,
            'TABLE_ROWS': row.total_rows,
            'STORAGE_SIZE': row.total_bytes,
            'PARTITION_TIME': row.partition_time
        })
    return partitions

三、核心分区对比检查项

基础一致性检查

  • 分区键一致性:对比源和目标表的分区列名称、顺序是否完全一致
  • 分区类型一致性:例如源端是范围分区,目标端需同为范围分区;哈希分区需对比哈希列和分区数量
  • 分区总数一致性:统计两边的分区数量是否相等

分区边界/值检查

  • 范围分区:对比每个分区的上下边界(如MySQL的PARTITION_DESCRIPTION、Snowflake的START_VALUE/END_VALUE)
  • 列表分区:对比每个分区的枚举值集合是否完全匹配
  • 时间分区:对比分区的时间范围(如BigQuery的partition_time)是否对齐

分区数据状态检查

  • 分区行数对比:验证同分区的源/目标端记录数是否一致
  • 存储大小对比:对比同分区的存储占用(部分数据库支持)
  • 分区有效性检查:检查源/目标端是否存在空分区、离线分区或未同步的分区

四、对比逻辑实现

将源和目标的分区信息标准化后,逐维度对比:

def compare_partitions(source_parts, target_parts, source_type, target_type):
    diffs = []
    # 检查分区数量
    if len(source_parts) != len(target_parts):
        diffs.append(f"分区数量不匹配:{source_type}({len(source_parts)}) vs {target_type}({len(target_parts)})")
    
    # 检查分区键和类型
    if source_parts and target_parts:
        source_key = source_parts[0]['PARTITION_EXPRESSION']
        target_key = target_parts[0]['PARTITION_EXPRESSION']
        if source_key != target_key:
            diffs.append(f"分区键不匹配:{source_type}({source_key}) vs {target_type}({target_key})")
        
        source_type_val = source_parts[0]['PARTITION_METHOD']
        target_type_val = target_parts[0]['PARTITION_METHOD']
        if source_type_val != target_type_val:
            diffs.append(f"分区类型不匹配:{source_type}({source_type_val}) vs {target_type}({target_type_val})")
    
    # 逐个分区对比数据
    # 针对不同数据库的分区命名规则,优先按边界/值匹配,而非名称
    for src_p in source_parts:
        # 匹配逻辑:根据分区类型找对应目标分区
        match_target = None
        for tgt_p in target_parts:
            # 范围/时间分区按边界匹配
            if src_p.get('PARTITION_BOUND') and tgt_p.get('PARTITION_BOUND'):
                if src_p['PARTITION_BOUND'] == tgt_p['PARTITION_BOUND']:
                    match_target = tgt_p
                    break
            # 列表分区按枚举值匹配(需提前标准化)
            elif src_p.get('PARTITION_VALUES') and tgt_p.get('PARTITION_VALUES'):
                if set(src_p['PARTITION_VALUES']) == set(tgt_p['PARTITION_VALUES']):
                    match_target = tgt_p
                    break
        
        if not match_target:
            diffs.append(f"{source_type}分区{src_p['PARTITION_NAME']}在{target_type}中无匹配")
            continue
        
        # 对比行数
        if int(src_p['TABLE_ROWS']) != int(match_target['TABLE_ROWS']):
            diffs.append(f"分区{src_p['PARTITION_NAME']}行数不匹配:{source_type}({src_p['TABLE_ROWS']}) vs {target_type}({match_target['TABLE_ROWS']})")
        
        # 对比存储大小
        if src_p.get('STORAGE_SIZE') and match_target.get('STORAGE_SIZE'):
            if int(src_p['STORAGE_SIZE']) != int(match_target['STORAGE_SIZE']):
                diffs.append(f"分区{src_p['PARTITION_NAME']}存储大小不匹配:{source_type}({src_p['STORAGE_SIZE']}) vs {target_type}({match_target['STORAGE_SIZE']})")
    
    return diffs

五、生成对比报告

用pandas将差异结果导出为Excel或CSV:

import pandas as pd

def generate_report(diffs, output_path):
    df = pd.DataFrame({'差异描述': diffs})
    df.to_excel(output_path, index=False)
    print(f"对比报告已生成:{output_path}")

内容的提问来源于stack exchange,提问作者Alia

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 21:45:38