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

Pandas数据更新时OperationType标记错误问题求助

场景描述

初始输入CSV:

Date,ProductID,Price,Quantity
2023-01-01,1001,10,1
2023-01-02,1001,10,1
2023-01-02,1011,10,6

对应的初始output_data.csv:

ProductID,TotalSales
1001,20
1011,60

当输入CSV更新为:

Date,ProductID,Price,Quantity
2023-01-01,1001,10,2
2023-01-02,1001,10,1
2023-01-02,1011,10,6
2023-01-02,1012,10,6

期望的output_data.csv输出:

ProductID,TotalSales,OperationType
1001,30,Updated
1011,60,No Change
1012,60,Insert

OperationType规则:

  • TotalSales变化标记为Updated
  • 无变化标记为No Change
  • 新增ProductID标记为Insert

问题

运行以下Pandas代码后,实际输出的OperationType全部为No Change,与预期不符:

import pandas as pd
import logging

logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s')

def extract(file_path):
    """Extract data from a CSV file."""
    return pd.read_csv(file_path)

def transform(data):
    """Transform data by calculating total sales."""
    data['TotalSales'] = data['Price'] * data['Quantity']
    return data.groupby('ProductID')['TotalSales'].sum().reset_index()

def load(data, output_file_path):
    """Load data into a new CSV file."""
    data.to_csv(output_file_path, index=False)

def insert_data(data, new_data):
    """Insert new data into the existing dataset."""
    updated_data = pd.concat([data, new_data], ignore_index=True)
    return updated_data

def delete_data(data, condition):
    """Delete rows based on a condition."""
    deleted_rows = data[condition]
    updated_data = data.drop(data[condition].index)
    return updated_data, deleted_rows

def update_data(data, condition, update_values):
    """Update data based on a condition."""
    updated_rows = data.loc[condition].copy()
    data.loc[condition, update_values.columns] = update_values.values
    return data, updated_rows

def identify_operation(original_data, updated_data):
    """Identify operation type for each row."""
    merged = pd.merge(original_data, updated_data, on='ProductID', how='outer', suffixes=('_old', '_new'), indicator=True)
    
    operations = []
    for index, row in merged.iterrows():
        if row['_merge'] == 'left_only':
            operations.append('No Change')
        elif row['_merge'] == 'right_only':
            operations.append('Insert')
        else:
            if row['TotalSales_old'] != row['TotalSales_new']:
                operations.append('Updated')
            else:
                operations.append('No Change')
    
    return pd.Series(operations, name='OperationType')  # Return as a Series

def run_pipeline(input_file_path, output_file_path):
    """Run the data pipeline."""
    data = extract(input_file_path)
    transformed_data = transform(data)
    original_data = transformed_data.copy()

    # Identify operation types
    operations = identify_operation(original_data, transformed_data)
    transformed_data = pd.concat([transformed_data, operations], axis=1)  # Concatenate as a new column
    load(transformed_data[['ProductID', 'TotalSales', 'OperationType']], output_file_path)

if __name__ == '__main__':
    run_pipeline('data.csv', 'output_data.csv')

实际输出:

ProductID,TotalSales,OperationType
1001,30,No Change
1011,60,No Change
1012,60,No Change

问题排查与修复

核心问题

代码中run_pipeline函数的逻辑错误:original_data直接从当前输入CSV转换后的transformed_data复制而来,相当于拿同一份数据和自己对比,自然所有操作类型都会被判定为No Change。正确逻辑应该是加载上一次生成的output_data.csv作为原始数据,再和当前输入转换后的新数据对比。

修复后的完整代码

import pandas as pd
import logging
import os

logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s')

def extract(file_path):
    """Extract data from a CSV file."""
    return pd.read_csv(file_path)

def transform(data):
    """Transform data by calculating total sales."""
    data['TotalSales'] = data['Price'] * data['Quantity']
    return data.groupby('ProductID')['TotalSales'].sum().reset_index()

def load(data, output_file_path):
    """Load data into a new CSV file."""
    data.to_csv(output_file_path, index=False)

def identify_operation(original_data, updated_data):
    """Identify operation type for each row."""
    merged = pd.merge(original_data, updated_data, on='ProductID', how='outer', suffixes=('_old', '_new'), indicator=True)
    
    operations = []
    for index, row in merged.iterrows():
        if row['_merge'] == 'right_only':
            operations.append('Insert')
        elif row['_merge'] == 'left_only':
            # 可根据需求处理删除场景,当前场景暂不涉及
            operations.append('Deleted')
        else:
            if row['TotalSales_old'] != row['TotalSales_new']:
                operations.append('Updated')
            else:
                operations.append('No Change')
    
    # 合并操作类型与最新数据,避免索引对齐问题
    merged['OperationType'] = operations
    result = merged[['ProductID', 'TotalSales_new', 'OperationType']].rename(columns={'TotalSales_new': 'TotalSales'})
    return result

def run_pipeline(input_file_path, output_file_path):
    """Run the data pipeline."""
    # 加载当前输入并转换
    current_data = extract(input_file_path)
    transformed_current = transform(current_data)
    
    # 加载历史输出作为原始数据,首次运行则初始化空DataFrame
    if os.path.exists(output_file_path):
        original_data = extract(output_file_path)
        # 仅保留核心对比列,避免历史OperationType干扰
        original_data = original_data[['ProductID', 'TotalSales']]
    else:
        original_data = pd.DataFrame(columns=['ProductID', 'TotalSales'])
    
    # 识别操作类型并生成结果
    result_data = identify_operation(original_data, transformed_current)
    
    # 保存最终结果
    load(result_data, output_file_path)

if __name__ == '__main__':
    run_pipeline('data.csv', 'output_data.csv')

关键修改点

  1. 原始数据来源修正:改为读取历史output_data.csv作为对比基准,首次运行时初始化空DataFrame。
  2. 合并逻辑优化:在identify_operation中直接生成包含最新TotalSales的结果集,避免后续拼接时的索引对齐问题。
  3. 历史数据清理:读取历史数据时仅保留ProductID和TotalSales列,防止旧的OperationType干扰对比逻辑。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 10:57:49