灾难恢复
灾难恢复对于任何软件都至关重要,尤其是对于数据系统,其影响可能非常严重,会导致业务决策延迟,甚至在某些情况下产生错误的业务决策。Apache Hudi 提供了两个操作来帮助您从之前的状态恢复数据:savepoint 和 restore。
Savepoint(保存点)
顾名思义,savepoint 会在指定的提交时间点保存表的状态,以便您在需要时将表恢复到该保存点。系统会确保清理器(cleaner)不会清理任何已保存点的文件。同样地,已经清理过的提交无法触发 savepoint。简单来说,这类似于进行备份,只是我们不会创建表的新副本,而是优雅地保存表的状态,以便在需要时进行恢复。
Restore(恢复)
此操作可将您的表恢复到某个保存点的提交状态。该操作无法撤销(或回滚),因此在执行恢复之前务必谨慎。Hudi 会删除大于目标恢复保存点提交的所有数据文件和提交文件(时间线文件)。在执行恢复操作时,应暂停对表的所有写入操作,因为这些写入很可能在恢复过程中失败。此外,读取操作也可能会失败,因为快照查询会命中最新的文件,而这些文件在恢复过程中极有可能被删除。
操作手册
Savepoint 和 restore 可以通过 Hudi CLI 和 SQL 存储过程 触发。下面我们通过一个示例来了解如何创建保存点以及之后如何恢复表的状态。
注意: 使用 Hudi CLI 时,我们需要指定 表路径;而使用 SQL 存储过程时,我们需要提供 表名称。
让我们通过 spark-shell 创建一个 Hudi 表,并触发一批插入操作。
import org.apache.hudi.QuickstartUtils._
import scala.collection.JavaConversions._
import org.apache.spark.sql.SaveMode._
import org.apache.hudi.DataSourceReadOptions._
import org.apache.hudi.DataSourceWriteOptions._
import org.apache.hudi.config.HoodieWriteConfig._
val tableName = "hudi_trips_cow"
val basePath = "file:///tmp/hudi_trips_cow"
val dataGen = new DataGenerator
val inserts = convertToStringList(dataGen.generateInserts(10))
val df = spark.read.json(spark.sparkContext.parallelize(inserts, 2))
df.write.format("hudi").
options(getQuickstartWriteConfigs).
option("hoodie.table.ordering.fields", "ts").
option("hoodie.datasource.write.recordkey.field", "uuid").
option("hoodie.datasource.write.partitionpath.field", "partitionpath").
option("hoodie.table.name", tableName).
mode(Overwrite).
save(basePath)让我们再插入四批数据。
for (_ <- 1 to 4) {
val inserts = convertToStringList(dataGen.generateInserts(10))
val df = spark.read.json(spark.sparkContext.parallelize(inserts, 2))
df.write.format("hudi").
options(getQuickstartWriteConfigs).
option("hoodie.table.ordering.fields", "ts").
option("hoodie.datasource.write.recordkey.field", "uuid").
option("hoodie.datasource.write.partitionpath.field", "partitionpath").
option("hoodie.table.name", tableName).
mode(Append).
save(basePath)
}记录总数应为 50。
val tripsSnapshotDF = spark.
read.
format("hudi").
load(basePath)
tripsSnapshotDF.createOrReplaceTempView("hudi_trips_snapshot")
spark.sql("select count(partitionpath, uuid) from hudi_trips_snapshot").show()
+--------------------------+
|count(partitionpath, uuid)|
+--------------------------+
| 50|
+--------------------------+让我们看看插入 5 批数据之后的时间线。
ls -ltr /tmp/hudi_trips_cow/.hoodie
total 128
drwxr-xr-x 2 nsb wheel 64 Jan 28 16:00 archived
-rw-r--r-- 1 nsb wheel 546 Jan 28 16:00 hoodie.properties
-rw-r--r-- 1 nsb wheel 0 Jan 28 16:00 20220128160040171.commit.requested
-rw-r--r-- 1 nsb wheel 2594 Jan 28 16:00 20220128160040171.inflight
-rw-r--r-- 1 nsb wheel 4374 Jan 28 16:00 20220128160040171.commit
-rw-r--r-- 1 nsb wheel 0 Jan 28 16:01 20220128160124637.commit.requested
-rw-r--r-- 1 nsb wheel 2594 Jan 28 16:01 20220128160124637.inflight
-rw-r--r-- 1 nsb wheel 4414 Jan 28 16:01 20220128160124637.commit
-rw-r--r-- 1 nsb wheel 0 Jan 28 16:02 20220128160226172.commit.requested
-rw-r--r-- 1 nsb wheel 2594 Jan 28 16:02 20220128160226172.inflight
-rw-r--r-- 1 nsb wheel 4427 Jan 28 16:02 20220128160226172.commit
-rw-r--r-- 1 nsb wheel 0 Jan 28 16:02 20220128160229636.commit.requested
-rw-r--r-- 1 nsb wheel 2594 Jan 28 16:02 20220128160229636.inflight
-rw-r--r-- 1 nsb wheel 4428 Jan 28 16:02 20220128160229636.commit
-rw-r--r-- 1 nsb wheel 0 Jan 28 16:02 20220128160245447.commit.requested
-rw-r--r-- 1 nsb wheel 2594 Jan 28 16:02 20220128160245447.inflight
-rw-r--r-- 1 nsb wheel 4428 Jan 28 16:02 20220128160245447.commitSavepoint 示例
可以通过 Hudi CLI 或 SQL 存储过程创建 Savepoint。
让我们针对最新的一次提交创建一个 Savepoint。
使用 Hudi CLI
- 启动 Hudi CLI。
- 如果尚未指定 SPARK_HOME,则指定它。
cd hudi-cli
./hudi-cli.sh
set --conf SPARK_HOME=<SPARK_HOME>- 使用表路径连接到表,例如
/tmp/hudi_trips_cow/。 - 运行
commits show命令以显示表中的提交记录。 - 运行
savepoint create命令并指定commit_time以创建 Savepoint。
connect --path /tmp/hudi_trips_cow/
commits show
savepoint create --commit 20220128160245447 --sparkMaster local[2]注意:
请将 20220128160245447 替换为你表中最新的提交。
使用 Spark SQL 存储过程
- 通过指定 Spark 版本和 Hudi 版本启动
spark-sqlshell。例如:
export SPARK_VERSION=3.5
export HUDI_VERSION=1.2.1
spark-sql --packages org.apache.hudi:hudi-spark$SPARK_VERSION-bundle_2.12:$HUDI_VERSION \
--conf 'spark.serializer=org.apache.spark.serializer.KryoSerializer' \
--conf 'spark.sql.extensions=org.apache.spark.sql.hudi.HoodieSparkSessionExtension' \
--conf 'spark.sql.catalog.spark_catalog=org.apache.spark.sql.hudi.catalog.HoodieCatalog' \
--conf 'spark.kryo.registrator=org.apache.spark.HoodieSparkKryoRegistrar'- 运行
show_commits命令,显示表中的提交记录。 - 运行
create_savepoint命令并指定 commit_time,以创建保存点。
call show_commits(table => 'hudi_trips_cow');
call create_savepoint(table => 'hudi_trips_cow', commit_time => '20220128160245447');注意:
请确保将 20220128160245447 替换为你表中最新的 commit。
我们来查看 savepoint 之后的时间线。
ls -ltr /tmp/hudi_trips_cow/.hoodie
total 136
drwxr-xr-x 2 nsb wheel 64 Jan 28 16:00 archived
-rw-r--r-- 1 nsb wheel 546 Jan 28 16:00 hoodie.properties
-rw-r--r-- 1 nsb wheel 0 Jan 28 16:00 20220128160040171.commit.requested
-rw-r--r-- 1 nsb wheel 2594 Jan 28 16:00 20220128160040171.inflight
-rw-r--r-- 1 nsb wheel 4374 Jan 28 16:00 20220128160040171.commit
-rw-r--r-- 1 nsb wheel 0 Jan 28 16:01 20220128160124637.commit.requested
-rw-r--r-- 1 nsb wheel 2594 Jan 28 16:01 20220128160124637.inflight
-rw-r--r-- 1 nsb wheel 4414 Jan 28 16:01 20220128160124637.commit
-rw-r--r-- 1 nsb wheel 0 Jan 28 16:02 20220128160226172.commit.requested
-rw-r--r-- 1 nsb wheel 2594 Jan 28 16:02 20220128160226172.inflight
-rw-r--r-- 1 nsb wheel 4427 Jan 28 16:02 20220128160226172.commit
-rw-r--r-- 1 nsb wheel 0 Jan 28 16:02 20220128160229636.commit.requested
-rw-r--r-- 1 nsb wheel 2594 Jan 28 16:02 20220128160229636.inflight
-rw-r--r-- 1 nsb wheel 4428 Jan 28 16:02 20220128160229636.commit
-rw-r--r-- 1 nsb wheel 0 Jan 28 16:02 20220128160245447.commit.requested
-rw-r--r-- 1 nsb wheel 2594 Jan 28 16:02 20220128160245447.inflight
-rw-r--r-- 1 nsb wheel 4428 Jan 28 16:02 20220128160245447.commit
-rw-r--r-- 1 nsb wheel 0 Jan 28 16:05 20220128160245447.savepoint.inflight
-rw-r--r-- 1 nsb wheel 1168 Jan 28 16:05 20220128160245447.savepoint你会注意到系统新增了 savepoint 元数据文件,其中记录了构成最新表快照的文件。
现在,让我们继续插入三批数据。
for (_ <- 1 to 3) {
val inserts = convertToStringList(dataGen.generateInserts(10))
val df = spark.read.json(spark.sparkContext.parallelize(inserts, 2))
df.write.format("hudi").
options(getQuickstartWriteConfigs).
option("hoodie.table.ordering.fields", "ts").
option("hoodie.datasource.write.recordkey.field", "uuid").
option("hoodie.datasource.write.partitionpath.field", "partitionpath").
option("hoodie.table.name", tableName).
mode(Append).
save(basePath)
}总记录数将为 80,因为我们总共执行了 8 批次(保存点之前 5 批,保存点之后 3 批)。
val tripsSnapshotDF = spark.
read.
format("hudi").
load(basePath)
tripsSnapshotDF.createOrReplaceTempView("hudi_trips_snapshot")
spark.sql("select count(partitionpath, uuid) from hudi_trips_snapshot").show()
+--------------------------+
|count(partitionpath, uuid)|
+--------------------------+
| 80|
+--------------------------+恢复示例
假设发生了某些意外情况,你需要将表恢复到较旧的快照。我们可以通过 Hudi CLI 或 SQL Procedures 执行恢复操作。请记住,在执行恢复期间要关闭所有的写入进程。
在触发恢复操作之前,我们先来看一下 timeline。
ls -ltr /tmp/hudi_trips_cow/.hoodie
total 208
drwxr-xr-x 2 nsb wheel 64 Jan 28 16:00 archived
-rw-r--r-- 1 nsb wheel 546 Jan 28 16:00 hoodie.properties
-rw-r--r-- 1 nsb wheel 0 Jan 28 16:00 20220128160040171.commit.requested
-rw-r--r-- 1 nsb wheel 2594 Jan 28 16:00 20220128160040171.inflight
-rw-r--r-- 1 nsb wheel 4374 Jan 28 16:00 20220128160040171.commit
-rw-r--r-- 1 nsb wheel 0 Jan 28 16:01 20220128160124637.commit.requested
-rw-r--r-- 1 nsb wheel 2594 Jan 28 16:01 20220128160124637.inflight
-rw-r--r-- 1 nsb wheel 4414 Jan 28 16:01 20220128160124637.commit
-rw-r--r-- 1 nsb wheel 0 Jan 28 16:02 20220128160226172.commit.requested
-rw-r--r-- 1 nsb wheel 2594 Jan 28 16:02 20220128160226172.inflight
-rw-r--r-- 1 nsb wheel 4427 Jan 28 16:02 20220128160226172.commit
-rw-r--r-- 1 nsb wheel 0 Jan 28 16:02 20220128160229636.commit.requested
-rw-r--r-- 1 nsb wheel 2594 Jan 28 16:02 20220128160229636.inflight
-rw-r--r-- 1 nsb wheel 4428 Jan 28 16:02 20220128160229636.commit
-rw-r--r-- 1 nsb wheel 0 Jan 28 16:02 20220128160245447.commit.requested
-rw-r--r-- 1 nsb wheel 2594 Jan 28 16:02 20220128160245447.inflight
-rw-r--r-- 1 nsb wheel 4428 Jan 28 16:02 20220128160245447.commit
-rw-r--r-- 1 nsb wheel 0 Jan 28 16:05 20220128160245447.savepoint.inflight
-rw-r--r-- 1 nsb wheel 1168 Jan 28 16:05 20220128160245447.savepoint
-rw-r--r-- 1 nsb wheel 0 Jan 28 16:06 20220128160620557.commit.requested
-rw-r--r-- 1 nsb wheel 2594 Jan 28 16:06 20220128160620557.inflight
-rw-r--r-- 1 nsb wheel 4428 Jan 28 16:06 20220128160620557.commit
-rw-r--r-- 1 nsb wheel 0 Jan 28 16:06 20220128160627501.commit.requested
-rw-r--r-- 1 nsb wheel 2594 Jan 28 16:06 20220128160627501.inflight
-rw-r--r-- 1 nsb wheel 4428 Jan 28 16:06 20220128160627501.commit
-rw-r--r-- 1 nsb wheel 0 Jan 28 16:06 20220128160630785.commit.requested
-rw-r--r-- 1 nsb wheel 2594 Jan 28 16:06 20220128160630785.inflight
-rw-r--r-- 1 nsb wheel 4428 Jan 28 16:06 20220128160630785.commit使用 Hudi CLI
- 启动 Hudi CLI,或使用已打开的 Hudi CLI。
- 如果尚未指定 SPARK_HOME,请指定其路径。
cd hudi-cli
./hudi-cli.sh
set --conf SPARK_HOME=<SPARK_HOME>- 使用指定的表路径连接到 Hudi 表,例如
/tmp/hudi_trips_cow/。 - 执行
refresh命令,将表状态更新到其最新版本。 - 运行
savepoints show命令以显示所有 savepoint。 - 运行
savepoint rollback并指定 savepoint 的 instant_time,以执行回滚操作。 - (可选)运行
savepoint delete命令,从现有 savepoint 中删除指定 instant_time 的 savepoint。
connect --path /tmp/hudi_trips_cow/
refresh
savepoints show
╔═══════════════════╗
║ SavepointTime ║
╠═══════════════════╣
║ 20220128160245447 ║
╚═══════════════════╝
savepoint rollback --savepoint 20220128160245447 --sparkMaster local[2]
savepoint delete --commit 20220128160245447 --sparkMaster local[2]注意:
请确保将 20220128160245447 替换为你表中最新的 savepoint。
使用 Spark SQL 过程
- 通过指定 Spark 版本和 Hudi 版本启动
spark-sqlshell,或使用现有的spark-sqlshell。
export SPARK_VERSION=3.5
export HUDI_VERSION=1.2.1
spark-sql --packages org.apache.hudi:hudi-spark$SPARK_VERSION-bundle_2.12:$HUDI_VERSION \
--conf 'spark.serializer=org.apache.spark.serializer.KryoSerializer' \
--conf 'spark.sql.extensions=org.apache.spark.sql.hudi.HoodieSparkSessionExtension' \
--conf 'spark.sql.catalog.spark_catalog=org.apache.spark.sql.hudi.catalog.HoodieCatalog' \
--conf 'spark.kryo.registrator=org.apache.spark.HoodieSparkKryoRegistrar'- 运行
show_savepoints命令,显示表中所有的保存点。 - 运行
rollback_to_savepoint命令,并指定要回滚到的保存点的 instant_time。 - (可选)运行
delete_savepoint命令,从现有的保存点中删除指定 instant_time 的保存点。
call show_savepoints(table => 'hudi_trips_cow');
call rollback_to_savepoint(table => 'hudi_trips_cow', instant_time => '20220128160245447');
call delete_savepoint(table => 'hudi_trips_cow', instant_time => '20220128160245447');注意:
请确保将 20220128160245447 替换为你表中最新的 savepoint(保存点)。
Hudi 表应已恢复到保存点所指向的提交 20220128160245447。数据文件和时间线文件都应已被删除。
ls -ltr /tmp/hudi_trips_cow/.hoodie
total 152
drwxr-xr-x 2 nsb wheel 64 Jan 28 16:00 archived
-rw-r--r-- 1 nsb wheel 546 Jan 28 16:00 hoodie.properties
-rw-r--r-- 1 nsb wheel 0 Jan 28 16:00 20220128160040171.commit.requested
-rw-r--r-- 1 nsb wheel 2594 Jan 28 16:00 20220128160040171.inflight
-rw-r--r-- 1 nsb wheel 4374 Jan 28 16:00 20220128160040171.commit
-rw-r--r-- 1 nsb wheel 0 Jan 28 16:01 20220128160124637.commit.requested
-rw-r--r-- 1 nsb wheel 2594 Jan 28 16:01 20220128160124637.inflight
-rw-r--r-- 1 nsb wheel 4414 Jan 28 16:01 20220128160124637.commit
-rw-r--r-- 1 nsb wheel 0 Jan 28 16:02 20220128160226172.commit.requested
-rw-r--r-- 1 nsb wheel 2594 Jan 28 16:02 20220128160226172.inflight
-rw-r--r-- 1 nsb wheel 4427 Jan 28 16:02 20220128160226172.commit
-rw-r--r-- 1 nsb wheel 0 Jan 28 16:02 20220128160229636.commit.requested
-rw-r--r-- 1 nsb wheel 2594 Jan 28 16:02 20220128160229636.inflight
-rw-r--r-- 1 nsb wheel 4428 Jan 28 16:02 20220128160229636.commit
-rw-r--r-- 1 nsb wheel 0 Jan 28 16:02 20220128160245447.commit.requested
-rw-r--r-- 1 nsb wheel 2594 Jan 28 16:02 20220128160245447.inflight
-rw-r--r-- 1 nsb wheel 4428 Jan 28 16:02 20220128160245447.commit
-rw-r--r-- 1 nsb wheel 0 Jan 28 16:05 20220128160245447.savepoint.inflight
-rw-r--r-- 1 nsb wheel 1168 Jan 28 16:05 20220128160245447.savepoint
-rw-r--r-- 1 nsb wheel 0 Jan 28 16:07 20220128160732437.restore.inflight
-rw-r--r-- 1 nsb wheel 4152 Jan 28 16:07 20220128160732437.restore让我们检查一下表中的总记录数。它应该与我们触发 savepoint 之前的记录数一致。
val tripsSnapshotDF = spark.
read.
format("hudi").
load(basePath)
tripsSnapshotDF.createOrReplaceTempView("hudi_trips_snapshot")
spark.sql("select count(partitionpath, uuid) from hudi_trips_snapshot").show()
+--------------------------+
|count(partitionpath, uuid)|
+--------------------------+
| 50|
+--------------------------+如你所见,整个表的状态会被恢复到被保存点(savepoint)所标记的那个提交。用户可以选择按固定周期触发保存点,并在创建新保存点时删除旧的保存点。请记住,清理器(cleaner)不会清理被保存点标记的文件,因此用户需要不时地删除这些保存点。否则,存储空间可能无法得到回收。
注意: MOR 表的保存点与恢复功能仅从 0.11 版本开始提供。
视频
评论
登录后参与评论
KnowForge