PySpark代码无法自动将元数据库Silver字段设为0求助
问题场景
我编写了一段PySpark代码,功能是自动检测Raw层的新列:若SQL Server元数据库的Meta.SelectiveColumnLoading表中不存在该列,则插入该列记录(插入时Silver字段设为1)。代码末尾调用IgnoreInSilver函数,本应将新插入列的Silver字段改回0,但该更新操作未生效。
核心代码
InsertInto函数(插入新列并调用IgnoreInSilver)
#function to insert (new) column names into columnname Metatable def InsertInto(Source,RawObject): # Get all columns in the latest version of the raw object sqlStmt = sqlContext.sql(f'SHOW COLUMNS IN raw.{RawObject.lower()}') allColumns = sqlStmt.toPandas().iloc[:, 0].values # print(len(allColumns)) # Get the columns that are currently in Meta.SelectiveColumnLoading Object = mid(RawObject,len(Source)+2,len(RawObject)) currentColumns = QueryMetaToDataframe("""SELECT [NewColumnName] FROM Meta.SelectiveColumnLoading WHERE [Source] LIKE '{0}' AND [Object] like '{1}'""".format(Source.lower(), Object.lower())) #get ObjectName column currentColumnList = currentColumns.toPandas().iloc[:,0].values # print(len(currentColumnList)) # Subtract all the columns that are already in Meta.SelectiveColumnLoading diff = [] for col in allColumns: if col not in currentColumnList: if col not in diff: diff.append(col) print(str(len(diff)) + " columns need to be added to Meta.SelectiveColumnLoading for " + Object) for col in diff: # Insert the (new) columns InsertStmt = 'INSERT INTO Meta.SelectiveColumnLoading ([Source],[Object],[MergeKey],[OldColumnName],[NewColumnName],[Silver]) VALUES' #start of insert string if len(diff) > 0: insertedColumns = [] insertedColumns.append(diff) #temp = col.replace("cxsrec__", "").replace("__c", "") # cleaning unnecessary parts #temp = temp.split('_') #res = temp[0] + ''.join(sub.title() for sub in temp[1:]) #remove "_" res = col # for now, the oldname is the same as the new name columns, this can change later # function to allow for different PK per table per source if CheckForMergeKeys(Source.lower(), RawObject.lower(), col) == True: MergeKey = 1 else: MergeKey = 0 silver = 1 ObjectName = RawObject.split("_", 1)[1] #.capitalize() # we don't want as a table 'source_table' but just 'Table' InsertStmt += """('{0}','{1}','{2}','{3}','{4}','{5}'),""".format(Source,ObjectName,MergeKey,col,res,silver) #add values for each column to insert string InsertStmt = InsertStmt[:-1] #remove last comma before inserting print(InsertStmt) WriteToMeta(InsertStmt) # if diff > 0 and not initial run of Object if len(diff) > 0 and len(currentColumnList) > 0: i = 0 for ins in insertedColumns: IgnoreInSilver(Source, ins[i], ObjectName) i+=1 return
IgnoreInSilver函数(更新Silver字段)
def IgnoreInSilver(Source,ColumnName,Object): WriteToMeta(f'''UPDATE Meta.SelectiveColumnLoading SET Silver = 0 WHERE Source = '{Source}' and Object = '{Object}' and OldColumnName = '{ColumnName}';''') return print('Column: ' + ColumnName + ' will be ignored in the Silver Layer.')
问题排查与修复
1. 最核心问题:insertedColumns构造逻辑错误
在for col in diff循环内部,每次都重新初始化insertedColumns = []并执行insertedColumns.append(diff),这会导致insertedColumns变成嵌套列表(例如diff = [col1, col2]时,insertedColumns = [[col1, col2]])。后续循环for ins in insertedColumns只会执行一次,且仅处理diff的第一个元素,其余新插入的列根本不会调用IgnoreInSilver。
修复代码:
将insertedColumns的初始化移到diff循环外,每次插入列后单独添加当前列到列表,最后统一调用更新:
# 初始化移到diff循环外 insertedColumns = [] print(str(len(diff)) + " columns need to be added to Meta.SelectiveColumnLoading for " + Object) for col in diff: # Insert the (new) columns InsertStmt = 'INSERT INTO Meta.SelectiveColumnLoading ([Source],[Object],[MergeKey],[OldColumnName],[NewColumnName],[Silver]) VALUES' res = col if CheckForMergeKeys(Source.lower(), RawObject.lower(), col) == True: MergeKey = 1 else: MergeKey = 0 silver = 1 ObjectName = RawObject.split("_", 1)[1] InsertStmt += """('{0}','{1}','{2}','{3}','{4}','{5}'),""".format(Source,ObjectName,MergeKey,col,res,silver) InsertStmt = InsertStmt[:-1] print(InsertStmt) WriteToMeta(InsertStmt) # 仅添加当前插入的列 insertedColumns.append(col) # 统一调用更新,移到diff循环外 if len(diff) > 0 and len(currentColumnList) > 0: for col in insertedColumns: IgnoreInSilver(Source, col, ObjectName)
2. 潜在问题:字段大小写不匹配
查询元数据表时使用了Source.lower()和Object.lower(),但插入和更新时用的是原始大小写值。如果SQL Server的排序规则区分大小写,会导致更新条件无法匹配到目标行。
修复方案:
统一所有操作的大小写,例如全部转小写:
# 插入时转小写 InsertStmt += """('{0}','{1}','{2}','{3}','{4}','{5}'),""".format(Source.lower(),ObjectName.lower(),MergeKey,col,res,silver) # 更新时也转小写 def IgnoreInSilver(Source,ColumnName,Object): WriteToMeta(f'''UPDATE Meta.SelectiveColumnLoading SET Silver = 0 WHERE Source = '{Source.lower()}' and Object = '{Object.lower()}' and OldColumnName = '{ColumnName}';''') return print('Column: ' + ColumnName + ' will be ignored in the Silver Layer.')
3. 潜在问题:ObjectName逻辑不一致
代码中Object变量用mid函数获取,而ObjectName用split获取,两者逻辑可能得到不同结果,导致更新时Object条件不匹配。
修复方案:
统一ObjectName的获取逻辑,例如都用split:
ObjectName = RawObject.split("_", 1)[1] # 后续查询元数据时使用统一的ObjectName currentColumns = QueryMetaToDataframe("""SELECT [NewColumnName] FROM Meta.SelectiveColumnLoading WHERE [Source] LIKE '{0}' AND [Object] like '{1}'""".format(Source.lower(), ObjectName.lower()))
4. 潜在问题:事务未提交
检查WriteToMeta函数,确保每次执行SQL后都提交事务(例如pyodbc连接需调用conn.commit(),Spark JDBC需配置正确的写入模式),否则插入和更新操作可能未持久化到数据库。
内容的提问来源于stack exchange,提问作者B C

