处理AWS S3文件时For循环中Try/Except失效问题排查
问题
处理AWS S3桶文件时,我希望捕获单个文件的异常:将处理失败的文件名加入列表、打印异常后继续处理剩余文件。正常流程无问题,我故意修改了一个文件的列名模拟错误,但现在所有文件都抛出异常,没有文件能正常处理。
我的代码如下:
for item in settings.keys: try: response = settings.client.get_object(Bucket=settings.source_bucket, Key=item) tmp = pd.read_csv(io.BytesIO(response['Body'].read()), encoding='unicode_escape', sep=None, engine='python') tmp['account_number'] = item.split('/')[4][:-4] tmp.columns = tmp.columns.str.strip() tmp.columns = tmp.columns.map(settings._config['balances']['columns']) df = pd.concat([df, tmp], ignore_index=False) except: settings.unprocessed.append(item) logger.exception(f'{item} Not Processed')
修改文件前所有文件都能正常处理,现在的错误日志如下:
2023-01-25 14:59:56 - ERROR - xxxx.csv Not Processed Traceback (most recent call last): File "C:\Users\xxxx\Desktop\xxxx\Python\xxxx\xxxx\xxxx.py", line 19, in balances df = pd.concat([df, tmp], ignore_index=False) File "C:\Users\xxxx\Desktop\xxxx\Python\xxxx\xxxx\xxxx\venv\lib\site-packages\pandas\util\_decorators.py", line 311, in wrapper return func(*args, **kwargs) File "C:\Users\xxxx\Desktop\xxxx\Python\xxxx\xxxx\xxxx\venv\lib\site-packages\pandas\core\reshape\concat.py", line 360, in concat return op.get_result() File "C:\Users\xxxx\Desktop\xxxx\Python\xxxx\xxxx\xxxx\venv\lib\site-packages\pandas\core\reshape\concat.py", line 591, in get_result indexers[ax] = obj_labels.get_indexer(new_labels) File "C:\Users\xxxx\Desktop\xxxx\Python\xxxx\xxxx\xxxx\venv\lib\site-packages\pandas\core\indexes\base.py", line 3721, in get_indexer raise InvalidIndexError(self._requires_unique_msg) pandas.errors.InvalidIndexError: Reindexing only valid with uniquely valued Index objects raise ValueError("cannot reindex on an axis with duplicate labels") ValueError: cannot reindex on an axis with duplicate labels
更新:处理其他类型文件时,类似逻辑的代码可以正常工作——异常被捕获打印,失败文件名加入列表,其他文件正常处理。正常运行的代码如下:
for item in settings.keys: try: tmp = pd.DataFrame() response = settings.client.get_object(Bucket=settings.source_bucket, Key=item) if item.endswith('.csv'): tmp = pd.read_csv(io.BytesIO(response['Body'].read()), encoding='unicode_escape', sep=None, engine='python') elif item.endswith('.xlsx'): tmp = pd.read_excel(io.BytesIO(response['Body'].read())) tmp['file'] = item.split('/')[4] tmp.columns = tmp.columns.map(settings._config['account statements']['columns']) tmp['row'] = tmp.index + 2 tmp.columns = tmp.columns.astype(str) tmp.rename(columns=lambda x: x.strip()) for col in tmp.columns: if col.startswith('Beløp'): settings.statement_currencies[item.split('/')[-1:][0]] = col[-3:] tmp[col] = tmp[col].astype(str) tmp[col] = tmp[col].str.replace(',', '.') tmp[col] = tmp[col].astype(float) tmp['direction'] = np.where(tmp[col] > 0, 'Incoming', 'Outgoing') df = pd.concat([df, tmp], ignore_index=False) except: settings.unprocessed.append(item) logger.exception(f'{item} Not Processed')
问题原因与解决办法
核心原因
问题出在全局df变量被异常文件污染:当第一个异常文件处理失败时,df已经因为拼接操作产生了重复索引/列标签的异常状态;后续文件处理时,依然尝试将新数据拼接到这个有问题的df上,导致所有后续拼接操作都触发相同的InvalidIndexError异常。
对比正常运行的代码,其处理的文件不会导致拼接后的df出现重复标签,因此单次异常不会影响全局;而故障代码中,被修改的文件在列映射后产生了重复列名,第一次拼接失败后df的异常状态持续存在,导致后续所有拼接都报错。
解决步骤
提前验证列映射结果:在拼接前先检查当前文件处理后的列是否存在重复,避免污染全局
df:for item in settings.keys: try: response = settings.client.get_object(Bucket=settings.source_bucket, Key=item) tmp = pd.read_csv(io.BytesIO(response['Body'].read()), encoding='unicode_escape', sep=None, engine='python') tmp['account_number'] = item.split('/')[4][:-4] tmp.columns = tmp.columns.str.strip() # 先验证列映射后是否有重复列名 mapped_cols = tmp.columns.map(settings._config['balances']['columns']) if mapped_cols.duplicated().any(): raise ValueError(f"文件{item}列映射后存在重复列名") tmp.columns = mapped_cols # 验证通过再执行拼接 df = pd.concat([df, tmp], ignore_index=False) except: settings.unprocessed.append(item) logger.exception(f'{item} Not Processed')初始化干净的
df:在循环开始前显式初始化df,避免之前的运行残留异常状态:df = pd.DataFrame()批量拼接优化(可选):用列表收集所有处理成功的数据框,最后一次性拼接,彻底避免单次失败影响全局:
processed_dfs = [] for item in settings.keys: try: # 处理当前文件的逻辑... # 验证通过后加入列表 processed_dfs.append(tmp) except: settings.unprocessed.append(item) logger.exception(f'{item} Not Processed') # 最后统一拼接所有成功的数据 df = pd.concat(processed_dfs, ignore_index=False)
内容的提问来源于stack exchange,提问作者pyth0nEiken
相关产品推荐
相关产品推荐

