Python遍历DataFrame传参Pipeline处理时触发索引越界异常排查
问题场景
需要对接源端共300张数据表,构建数据管道完成数据连接与发布,当前实现代码如下:
#Create a pandas DF from a Json file and extract certain columns like tabname,colname,col length and data type df = metadata_extract(file_name) #Create a instance of an object to run a pipeline and after passing the connection details pipeline = Pipeline(server,source_type,source_name,database_name,user_name,password) tables = pd.unique(df["tab_name"]) #create a function to extract and process the tables from the source system def get_table_data(): for i in range(len(tables)): current_tab = df[df["tab_name"] == tables[i]].reset_index() print (current_tab) data = pipeline.process_table(current_tab) print (data) get_table_data()
逻辑说明:首先通过metadata_extract读取JSON文件生成pandas DataFrame,提取表名、列名、列长度、数据类型等元数据字段;初始化传入数据库连接参数的Pipeline实例后,提取tab_name字段的去重唯一值列表,循环遍历逐表筛选对应元数据生成current_tab对象,调用pipeline.process_table方法完成单表处理,预期在控制台输出所有源表的处理结果。
异常表现
执行过程中,部分经源端确认有效的表名触发IndexError: list index out of range异常,控制台仅能成功输出1张表的处理结果,其余表均执行失败。打印current_tab可正常输出多表的元数据内容,示例输出如下:
282 15940 Schema1.TABLE1 Colname CHAR 283 15941 Schema1.TABLE1 Colname CHAR 284 15942 Schema1.TABLE1 Colname CHAR 285 15943 Schema1.TABLE1 Colname CHAR 286 15944 Schema1.TABLE1 Colname CHAR [287 rows x 5 columns] index TableName ColumnName ColumnDataType ColumnLength 0 15945 Schema1.table2 Colname CHAR 2 1 15946 Schema1.table2 Colname CUKY 5 [154 rows x 5 columns] index TableName ColumnName ColumnDataType ColumnLength 0 15504 Schema1.TABLE3 Colname1 CHAR 3 1 15505 Schema1.TABLE3 Colname2 CHAR 3 2 15506 Schema1.TABLE3 Colname3 CHAR 3 3 15507 Schema1.TABLE3 Colname4 CHAR 3 4 15508 Schema1.TABLE3 Colname5 CHAR 3
根据报错堆栈定位,异常触发在pipeline包内部的generate_select_query方法,对应代码行为table_name = str(df["table_nm"][0]).split(".")[1],完整报错栈如下:
System.Private.CoreLib: Exception while executing function: Functions.metadata_trigger. System.Private.CoreLib: Result: Failure Exception: IndexError: list index out of range Stack: File "C:\Program Files\Microsoft\Azure Functions Core Tools\workers\python\3.9/WINDOWS/X64\azure_functions_worker\dispatcher.py", line 407, in _handle__invocation_request call_result = await self._loop.run_in_executor( File "C:\Users\AppData\Local\Programs\Python\Python39\lib\concurrent\futures\thread.py", line 58, in run result = self.fn(*self.args, **self.kwargs) File "C:\Program Files\Microsoft\Azure Functions Core Tools\workers\python\3.9/WINDOWS/X64\azure_functions_worker\dispatcher.py", line 649, in _run_sync_func return ExtensionManager.get_sync_invocation_wrapper(context, File "C:\Program Files\Microsoft\Azure Functions Core Tools\workers\python\3.9/WINDOWS/X64\azure_functions_worker\extension.py", line 215, in _raw_invocation_wrapper result = function(**args) File "C:\Users\metadata_trigger\__init__.py", line 31, in main get_tables() File "C:\Users\metadata_trigger\__init__.py", line 26, in get_tables data = pipeline.process_tab(current_tab) File "C:\Users\.venv\lib\site-packages\mdt\pipeline.py", line 78, in process_tab select = self.generate_select_query(df) File "C:\Users\.venv\lib\site-packages\mdt\pipeline.py", line 68, in generate_select_query table_name = str(df["table_nm"][0]).split(".")[1]
已确认打印的current_tab中存在对应表名数据,不清楚索引越界的触发原因,需要实现所有源表数据正常在控制台打印输出的效果。
根因分析
报错直接原因是pipeline内部硬编码解析逻辑的前置条件不满足:
- 该行代码默认传入的元数据DataFrame存在
table_nm列,且列中第0行值为半角点分隔的schema名.表名格式,通过split(".")[1]截取纯表名 - 触发
IndexError说明split(".")返回的列表长度小于2,即拿到的字符串中不存在半角.,无法取到索引为1的元素
结合日志可定位到两个具体诱因:
- 列名不匹配:打印出的
current_tab表名列名为TableName(大写开头),不符合pipeline内部硬编码读取的table_nm列名要求,导致内部映射逻辑没有正确取到带schema前缀的表名值 - 缺少前置校验:去重得到的表名列表没有提前过滤空值、格式非法值,部分表名可能存在缺失schema前缀、分隔符为全角句号、带不可见前后空格、值为NaN等问题,传入后直接触发解析失败
修复方案
按以下步骤调整代码即可解决:
- 优化循环写法,提前过滤无效表名,重置索引时丢弃旧索引避免多余列干扰
- 传入pipeline前统一重命名列,完全匹配pipeline要求的字段名
- 增加表名格式前置校验,提前拦截非法格式的表名,避免进入pipeline内部逻辑才报错
调整后的可运行代码如下:
df = metadata_extract(file_name) pipeline = Pipeline(server,source_type,source_name,database_name,user_name,password) tables = pd.unique(df["tab_name"]) def get_table_data(): for tab in tables: # 跳过空值、空字符串类无效表名 if pd.isna(tab) or not str(tab).strip(): print(f"跳过无效空表名条目") continue current_tab = df[df["tab_name"] == tab].reset_index(drop=True) # 重命名列适配pipeline的硬编码字段要求,根据实际字段映射补充即可 current_tab = current_tab.rename(columns={ "TableName": "table_nm", "ColumnName": "col_name", "ColumnDataType": "data_type", "ColumnLength": "col_length" }) # 提前校验表名格式 first_table_name = str(current_tab["table_nm"].iloc[0]).strip() if "." not in first_table_name: print(f"表[{first_table_name}]不符合schema.表名格式,跳过处理") continue print(current_tab) data = pipeline.process_table(current_tab) print(data) get_table_data()
内容的提问来源于stack exchange,提问作者user3764303
相关产品推荐
相关产品推荐

