【发布时间】:2020-06-20 14:09:30
【问题描述】:
我的数据源源不断变化。我正在通过 sqoop 提取该数据,但由于容量很大,我无法将其保留为每日截断负载。我想追加数据,但逻辑应该是更新和插入。如果通过删除先前的相同记录在源中更新记录,则应在配置单元中执行相同操作,即应删除旧记录并插入/更新新记录。 下面是一个这样的例子。
30 分钟后,数据更新如下:
现在,我的 hive 表选择了原始记录,一段时间后选择了更新的记录,但将其插入为不同的行。
我希望在不覆盖我的表的情况下,反映的数据与源中的数据相同。 (推荐使用 Pyspark 代码)
请帮忙。谢谢。
【问题讨论】:
-
您只提取新的和更新的记录,您如何识别增量记录?更新记录的数量是多少?
-
我在源中有一个时间戳列,当记录更改/更新时会更新,并且通过 sqoop 增量逻辑我将它们拉到配置单元中,因为 sqoop 总是将最后一个增量值存储在元数据中。更新记录的数量几乎是每天 8-100 万。它包括更新的 + 新条目。
-
您可能想看看其他存储格式(HBase、Kudu),因为普通 HDFS 没有更新的概念。
标签: hadoop pyspark hive apache-spark-sql pyspark-dataframes