【发布时间】:2019-08-20 11:14:51
【问题描述】:
我们正在使用 spark 处理大数据,最近有了新的用例,我们需要使用 spark 更新 Hive 表中的数据。
下面是一个简单的例子: 数据驻留在 Hive 表中,应用程序使用 PySpark 读入数据帧(例如 df1)。 例如:数据框有下面的列。
EmpNo Name 年龄薪水
1 aaaa 28 30000
2 bbbb 38 20000
3 cccc 26 25000
4 dddd 30 32000
需要使用 spark 向表中添加更多记录。
例如:
Action EmpNo Name 年龄薪水
添加 5 dddd 30 32000
应用程序可以通过剥离 Action 列并将新数据附加到表中,将新数据读入第二个数据帧(例如 df2)。它很简单,而且效果很好。
df.write.format('镶木地板') \ .mode('追加') \ .saveAsTable(canonical_hive_table)
在某些情况下,我们需要删除现有记录或根据 Action 列更新它们。
例如:
Action EmpNo Name 年龄薪水
删除 2 bbbb 38 20000
更新 4 dddd 30 42000
在上面的例子中,应用程序需要删除 EmpNo:2 并更新 EmpNo:4。
最终输出应如下所示:
EmpNo Name 年龄薪水
1 aaaa 28 30000
3 cccc 26 25000
4 dddd 30 42000
5 dddd 30 32000
据我了解,更新操作在 Spark Sql 中不可用,而且数据框是不可变的,无法更改记录。
有人遇到过这种情况吗?或知道使用 PySpark 更新 Hive 表中现有记录的任何选项?
请:应用程序需要定期处理数百万条记录的数千条更新。
提前致谢。
【问题讨论】:
-
您是否尝试在日常增量负载上使用 Hive 执行和 UPSERT(插入 + 更新)?如果是这样的话,它每天都会覆盖。我们可能需要比这里的逻辑更多的东西来起草解决方案。
-
你是怎么到这里来的?
标签: hive pyspark-sql