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

PySpark代码无法自动将元数据库Silver字段设为0求助

PySpark代码调用IgnoreInSilver更新Silver字段未生效的排查与修复

问题场景

我编写了一段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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 03:05:12