До сих пор я устанавливал переменные для каждого столбца в начале кода.
Есть ли другой вариант сделать код более динамичным, чтобы мне не приходилось устанавливать переменные для каждого столбца поверх а также для типов данных?
Я также определил в коде, что он будет проверять метаданные перед этим и конвертировать мои типы данных Hana в типы данных Spark
from pyspark.sql.functions import col
# Define column names as variables
COLUMN_KENNZAHL = 'kennzahl'
COLUMN_RECORDMODE = 'recordmode'
COLUMN_REQTSN = 'reqtsn'
COLUMN_REQTSN2 ="REQTSN"
COLUMN_DATAPAKID = 'datapakid'
COLUMN_RECORD = 'record'
COLUMN_NAME = 'name'
# Update Lakehouse Table
if initialLoad:
df = pd.DataFrame(BWDataInit, columns=columns)
print(f"DataFrame columns: {df.columns}") # Debug output
if COLUMN_KENNZAHL in df.columns:
df[COLUMN_KENNZAHL] = pd.to_numeric(df[COLUMN_KENNZAHL], errors='coerce') # Convert datatype of 'kennzahl' to float
else:
print(f"Error: '{COLUMN_KENNZAHL}' column not found in the data.")
if COLUMN_RECORDMODE in df.columns:
df = df.drop([COLUMN_RECORDMODE], axis=1)
# Save initial load to SQL table
append_records(df, DBWRITESCHEMA, DBWRITETABLE)
# Get requests from BWData
dfRequests = pd.DataFrame(BWDataRequest, columns=[COLUMN_REQTSN2])
# Update last Request table
update_last_request_table(dfRequests, DBWRITESCHEMA, DBWRITETABLE)
else:
df = pd.DataFrame(BWDataChange, columns=columns)
print(f"DataFrame columns: {df.columns}") # Debug output
if COLUMN_KENNZAHL in df.columns:
df[COLUMN_KENNZAHL] = pd.to_numeric(df[COLUMN_KENNZAHL], errors='coerce') # Convert datatype of 'kennzahl' to float
else:
print(f"Error: '{COLUMN_KENNZAHL}' column not found in the data.")
# Get list of all new Requests in Change Log
dfNewReq = df[COLUMN_REQTSN].unique()
dfNewReq = pd.DataFrame(dfNewReq, columns=[COLUMN_REQTSN2])
# Set Key Columns for Join Condition
colsKey = [COLUMN_NAME]
# Set KPI Columns for Join
colsKPI = [COLUMN_KENNZAHL]
if not dfNewReq.empty:
# Read content of SQL Table
sdfTable = spark.read.table(f"{DBWRITESCHEMA}.{DBWRITETABLE}")
# Ensure 'kennzahl' column in Spark DataFrame has the same datatype as in Pandas DataFrame
sdfTable = sdfTable.withColumn(COLUMN_KENNZAHL, sdfTable[COLUMN_KENNZAHL].cast('float'))
# Update SQL table for each request and record mode
for request in dfNewReq[COLUMN_REQTSN2]:
# Select only new records (Recordmode = N)
dfNew = df[(df[COLUMN_REQTSN] == request) & (df[COLUMN_RECORDMODE] == "N")]
if not dfNew.empty:
# Drop support columns
dfNew = dfNew.drop([COLUMN_REQTSN, COLUMN_DATAPAKID, COLUMN_RECORD, COLUMN_RECORDMODE], axis=1)
# Convert pandas df to Spark df
sdfNew = spark.createDataFrame(dfNew)
# Append new data to old data
sdfTable = sdfTable.union(sdfNew)
# Update changed records (Recordmode = "")
dfUpdate = df[(df[COLUMN_REQTSN] == request) & (df[COLUMN_RECORDMODE] == "")]
if not dfUpdate.empty:
# Drop unnecessary data left over from changelog structure
dfUpdate = dfUpdate.drop([COLUMN_REQTSN, COLUMN_DATAPAKID, COLUMN_RECORD, COLUMN_RECORDMODE], axis=1)
# Convert pandas to spark df
sdfUpdate = spark.createDataFrame(dfUpdate)
# Inner join sdfTable with sdfUpdate
sdfTable = sdfTable.alias('old').join(
sdfUpdate.alias('new'), colsKey, how='inner'
)
for kpi_col in colsKPI:
sdfTable = sdfTable.withColumn(f"old.{kpi_col}", col(f"new.{kpi_col}"))
# Delete old records (Recordmode = R)
dfDelete = df[(df[COLUMN_REQTSN] == request) & (df[COLUMN_RECORDMODE] == "R")]
if not dfDelete.empty:
# Drop unnecessary data left over from changelog structure
dfDelete = dfDelete.drop([COLUMN_REQTSN, COLUMN_DATAPAKID, COLUMN_RECORD, COLUMN_RECORDMODE], axis=1)
# Convert pandas to spark df
sdfDelete = spark.createDataFrame(dfDelete)
# Anti join sdfTable with sdfDelete to remove all records
sdfTable = sdfTable.alias('old').join(
sdfDelete.alias('new'), colsKey, how='left_anti'
)
# Overwrite old content of SQL Table with updated Spark df
update_last_request_table(dfNewReq, DBWRITESCHEMA, DBWRITETABLE)
sdfTable.write.mode("overwrite").format("delta").saveAsTable(f"{DBWRITESCHEMA}.{DBWRITETABLE}")
#print result table
print("Result Table:")
spark.read.table(f"{DBWRITESCHEMA}.{DBWRITETABLE}").sort(COLUMN_NAME).show()
Подробнее здесь: https://stackoverflow.com/questions/785 ... t-fix-in-t