【发布时间】:2018-06-14 19:42:26
【问题描述】:
我有 2 个数据框,比如 df1 和 df2。
df1 数据来自数据库,df2 是我从客户那里收到的新数据。我需要处理新数据,并根据是新记录还是要更新的现有记录执行UPSERTs。
样本数据输出:
df1= sqlContext.createDataFrame([("xxx1","81A01","TERR NAME 01","NJ"),("xxx2","81A01","TERR NAME 01","NJ"),("xxx3","81A01","TERR NAME 01","NJ"),("xxx4","81A01","TERR NAME 01","CA"),("xx5","81A01","TERR NAME 01","ME")], ["zip_code","territory_code","territory_name","state"])
df2= sqlContext.createDataFrame([("xxx1","81A01","TERR NAME 55","NY"),("xxx2","81A01","TERR NAME 55","NY"),("x103","81A01","TERR NAME 01","NJ")], ["zip_code","territory_code","territory_name","state"])
df1.show()
+--------+--------------+--------------+-----+
|zip_code|territory_code|territory_name|state|
+--------+--------------+--------------+-----+
| xxx1| 81A01| TERR NAME 01| NJ|
| xxx2| 81A01| TERR NAME 01| NJ|
| xxx3| 81A01| TERR NAME 01| NJ|
| xxx4| 81A01| TERR NAME 01| CA|
| xxx5| 81A01| TERR NAME 01| ME|
+---------------------------------------------
# Print out information about this data
df2.show()
+--------+--------------+--------------+-----+
|zip_code|territory_code|territory_name|state|
+--------+--------------+--------------+-----+
| xxx1| 81A01| TERR NAME 55| NY|
| xxx2| 81A01| TERR NAME 55| NY|
| x103| 81A01| TERR NAME 01| NJ|
+---------------------------------------------
预期结果: 我需要将 df2 数据帧与 df1 进行比较。 根据上述比较创建 2 个新数据集,即要更新的记录和要附加/插入到数据库的记录。
如果 zip_code 和 territory_code 相同,则为 UPDATE,否则为 INSERT 到数据库。
例如: INSERT 的新数据帧输出:
+--------+--------------+--------------+-----+
|zip_code|territory_code|territory_name|state|
+--------+--------------+--------------+-----+
| x103| 81A01| TERR NAME 01| NJ|
+---------------------------------------------
更新的新数据框:
+--------+--------------+--------------+-----+
|zip_code|territory_code|territory_name|state|
+--------+--------------+--------------+-----+
| xxx1| 81A01| TERR NAME 55| NY|
| xxx2| 81A01| TERR NAME 55| NY|
+---------------------------------------------
有人可以帮帮我吗?我正在使用 AWS Glue。
更新:解决方案(使用连接和减去)
df3 = df1.join(df2, (df1.zip_code == df2.zip_code_new) & (df1.territory_code == df2.territory_code_new))
df5=df3.drop("zip_code", "territory_code", "territory_name", "state")
df5.show()
+------------+------------------+------------------+---------+
|zip_code_new|territory_code_new|territory_name_new|state_new|
+------------+------------------+------------------+---------+
| x103| 81A01| TERR NAME 01| NJ|
+------------+------------------+------------------+---------+
df4=df2.subtract(df5)
df4.show()
+------------+------------------+------------------+---------+
|zip_code_new|territory_code_new|territory_name_new|state_new|
+------------+------------------+------------------+---------+
| xxx1 | 81A01 | TERR NAME 55 | NY |
| xxx2 | 81A01 | TERR NAME 55 | NY |
+------------------------------------------------------------+
对于 RDS 数据库更新,我使用 pymysql/Mysqldb:
db = MySQLdb.connect("xxxx.rds.amazonaws.com", "username", "password", "databasename")
cursor = db.cursor()
#cursor.execute("REPLACE INTO table SELECT * FROM table_stg")
insertQry = "INSERT INTO table VALUES('xxx1','81A01','TERR NAME 55','NY') ON DUPLICATE KEY UPDATE territory_name='TERR NAME 55', state='NY'"
n=cursor.execute(insertQry)
db.commit()
cursor.fetchall()
db.close()
谢谢
【问题讨论】:
-
下次请把你的问题写得更好。它(也许仍然)很难理解。也许您还应该提供一个数据集(我的意思是 4-5 个样本数据),以便其他人可以测试您的代码。
-
我想,我已经提供了足够的信息,并且我的问题很清楚,不确定是什么导致了反对票。无论如何,我将编辑问题并使其更短。
标签: python pyspark amazon-rds pyspark-sql aws-glue