I'm trying to perform merge(upsert) operation on MySql using Pyspark DataFrames and JDBC connection. Followed the below article to do same which is in scala. ( https://medium.com/@thomaspt748/how-to-upsert-data-into-relational-database-using-apache-spark-part-2-45a9d49d0f43 ).
But I need to perform upsert operations using pyspark got stuck at iterating through Pyspark Dataframe to call upsert function like below. Need to pass the source dataframe as input to read as parameters and perform sql upsert. (In simple terms performing the sql upsert using pyspark dataframe)
def upsertToDelta(id, name, price, purchase_date):
try:
connection = mysql.connector.connect(host='localhost',
database='Electronics',
user='pynative',
password='pynative@#29')
cursor = connection.cursor()
mySql_insert_query = "MERGE INTO targetTable USING
VALUES (%s, %s, %s, %s) as INSROW((Id, Name, Price, Purchase_date)
WHEN NOT MATCHED THEN INSERT VALUES (INSROW.Id,INSROW.Price,INSROW.Purchase,INSROW.Purchase_date)
WHEN MATCHED THEN UPDATE SET set Name=INSROW.Name"
recordTuple = (id, name, price, purchase_date)
cursor.execute(mySql_insert_query, recordTuple)
connection.commit()
print("Record inserted successfully into test table")
except mysql.connector.Error as error:
print("Failed to insert into MySQL table {}".format(error))
**
dataFrame.writeStream \
.format("delta") \
.foreachBatch(upsertToDelta) \
.outputMode("update") \
.start()
**
Any help on this is highly appreciated.