テーブル履歴を取り扱う
Apache Iceberg と Delta Lake テーブルの場合、テーブルを変更する各オペレーションは新しいテーブルバージョンを作成します。履歴情報を使用して、オペレーションを監査したり、テーブルをロールバックしたり、タイムトラベルを使用して特定の時点のテーブルを照会したりすることができます。
データアーカイブの長期バックアップソリューションとしてテーブル履歴を使用しないでください。データとLogsの保持構成の両方をより大きな値に設定していない限り、タイムトラベルオペレーションには過去7日間のみを使用してください。
テーブル履歴の取得
テーブルへの各書き込みのオペレーション、ユーザー、タイムスタンプなどの情報を取得するには、DESCRIBE HISTORY コマンドを実行します。オペレーションは、時系列の逆順で返されます。
DESCRIBE HISTORY が返す列、operationParameters 列の値、および operationMetrics 列のオペレーションごとのメトリクスについては、テーブル履歴スキーマとオペレーションメトリクス を参照してください。
テーブル履歴の保持期間はテーブル設定logRetentionDurationによって決まります。デフォルトでは30日間です。
タイムトラベルとテーブル履歴は、異なる保持しきい値によって制御されます。See タイムトラベル.
DESCRIBE HISTORY table_name -- get the full history of the table
DESCRIBE HISTORY table_name LIMIT 1 -- get the last operation only
Spark SQL構文の詳細については、DESCRIBE HISTORYを参照してください。
Scala、Java、Python の構文の詳細については、Delta Lake API ドキュメントを参照してください。
カタログ エクスプローラーには、 「履歴」 タブでテーブルの履歴が視覚的に表示されます。
OPTIMIZE操作のタイプを特定します
自動コンパクション、リキッドクラスタリング、およびZ-Orderingはすべて、テーブル履歴にOPTIMIZE操作として表示されます。どちらが実行されたかを判断するには、operationParameters 列を調べます。
テーブル履歴内のすべての OPTIMIZE 操作を分類するには、次を実行します:
SELECT
version,
timestamp,
CASE
WHEN operationParameters.clusterBy IS NOT NULL AND operationParameters.clusterBy <> '[]' THEN 'Liquid clustering'
WHEN operationParameters.zOrderBy IS NOT NULL AND operationParameters.zOrderBy <> '[]' THEN 'Z-ordering'
WHEN operationParameters.auto = 'true' THEN 'Auto compaction'
ELSE 'Manual OPTIMIZE'
END AS optimize_type,
operationParameters.auto AS is_auto_compaction,
operationParameters.clusterBy AS cluster_by,
operationParameters.zOrderBy AS z_order_by,
operationMetrics.numRemovedFiles AS files_compacted,
operationMetrics.numAddedFiles AS files_added,
operationMetrics.numRemovedBytes AS bytes_removed,
operationMetrics.numAddedBytes AS bytes_added
FROM (DESCRIBE HISTORY table_name)
WHERE operation = 'OPTIMIZE'
ORDER BY version DESC;
次のセクションでは、各 operationParameters 値について詳しく説明します。前のクエリーで選択された operationMetrics キーの定義については、「オペレーションメトリクス」を参照してください。
自動圧縮
自動圧縮は、autoパラメーターをtrueに設定します。Databricks は、書き込み後に自動的に自動圧縮をTriggerします。autoがfalseの場合、ユーザーまたはスケジュールされたジョブがOPTIMIZEコマンドを実行しました。
たとえば、自動圧縮操作は次のように表示されます:
operationParameters: {
"auto": "true"
}
自動圧縮の詳細については、「自動圧縮」を参照してください。
リキッドクラスタリング
リキッドクラスタリングによって、clusterByパラメーターにクラスタリング列名が入力されます。空のclusterBy配列([])は、ファイル圧縮のみを示します。
たとえば、date列とregion列でデータをクラスタリングした操作は、次のように表示されます。
operationParameters: {
"clusterBy": "[\"date\",\"region\"]"
}
リキッドクラスタリングの詳細については、「テーブルにリキッドクラスタリングを使用する」を参照してください。
Z-Ordering
Z-Orderingは、zOrderByパラメーターにZ-Order列名を設定します。空の zOrderBy 配列([])は、操作が Z-Ordering を適用しなかったことを示します。
例えば、date列にZ-Orderingを適用した操作は次のようになります。
operationParameters: {
"zOrderBy": "[\"date\"]"
}
操作スコープ
predicateパラメーターは、操作がテーブル全体で実行されたか、またはその一部のみで実行されたかを示します。
- 空の
predicate配列([])は、操作がテーブル全体で実行されたことを意味します。 - データが入力された
predicate配列は、ターゲットのOPTIMIZE table_name WHERE <partition_predicate>コマンドが述語に一致するパーティションのみで実行されたことを意味します。
例えば、year = 2024に一致するパーティションを対象とする操作は次のようになります。
operationParameters: {
"predicate": "[\"'year = 2024\"]"
}
タイムトラベル
タイムトラベルは、タイムスタンプまたはテーブルバージョン(トランザクションログに記録されている)に基づいた以前のテーブルバージョンのクエリーをサポートしています。タイムトラベルは、次のようなアプリケーションに使用できます:
- 分析、レポート、または出力(機械学習モデルの出力など)を再作成します。これは、特に規制された業界でのデバッグや監査に役立つ可能性があります。
- 複雑なテンポラルクエリーを記述する。
- データの誤りを修正する。
- 急速に変化するテーブルの一連のクエリーに対してスナップショット分離を提供します。
Databricks Runtime 18.0 以降では、 deletedFileRetentionDurationテーブル プロパティ (デフォルトは 7 日) よりも古いバージョンをリクエストした場合、タイムトラベル クエリはブロックされます。 Unity Catalogマネージドテーブルの場合、これはDatabricks Runtime 12.2 以降に適用されます。
タイムトラベル構文
タイムトラベルを使用してテーブルをクエリするには、テーブル名の指定の後に句を追加します。
-
timestamp_expression次のいずれかになります:'2018-10-18T22:15:12.013Z'つまり、タイムスタンプにキャストできる文字列ですcast('2018-10-18 13:36:32 CEST' as timestamp)'2018-10-18'、つまり日付文字列ですcurrent_timestamp() - interval 12 hoursdate_sub(current_date(), 1)- タイムスタンプにキャストされる、またはタイムスタンプにキャストできるその他の式
-
versionは、DESCRIBE HISTORY table_specの出力から取得できる長い値です。
timestamp_expressionもversionもサブクエリーにすることはできません。
日付またはタイムスタンプ文字列のみが受け入れられます。たとえば、"2019-01-01"と"2019-01-01T00:00:00.000Z"です。シンタックスの例については、以下のコードを参照してください:
- SQL
- Python
SELECT * FROM people10m TIMESTAMP AS OF '2018-10-18T22:15:12.013Z';
SELECT * FROM people10m VERSION AS OF 123;
df1 = spark.read.option("timestampAsOf", "2019-01-01").table("people10m")
df2 = spark.read.option("versionAsOf", 123).table("people10m")
@構文を使用して、タイムスタンプまたはバージョンをテーブル名の一部として指定することもできます。タイムスタンプはyyyyMMddHHmmssSSS形式である必要があります。@vを使用してバージョンを指定できます。構文の例については、以下のコードを参照してください:
- SQL
- Python
-- Timestamp version
SELECT * FROM people10m@20190101000000000
-- Version number
SELECT * FROM people10m@v123
# Timestamp version
spark.read.table("people10m@20190101000000000")
# Version number
spark.read.table("people10m@v123")
タイムトラベルクエリーのデータ保持を構成する
以前のテーブルバージョンをクエリするには、そのバージョンのログファイルとデータファイルの*両方*を保持する必要があります。
VACUUMテーブルに対して実行されると、データファイルが削除されます。- テーブルバージョンのチェックポイント設定後、ログファイルは自動的に削除されます。
テーブルのデータ保持しきい値を増やすには、<format> を delta または iceberg のいずれかに置き換えて、次のテーブル プロパティを構成する必要があります。
-
<format>.logRetentionDuration = "interval <interval>":テーブルの履歴を保持する期間を制御します。デフォルトはinterval 30 daysです。- Databricks Runtime 18.0以降では、
logRetentionDurationはdeletedFileRetentionDuration以上である必要があります。Unity Catalogマネージドテーブルの場合、これはDatabricks Runtime 12.2 以降に適用されます。
- Databricks Runtime 18.0以降では、
-
<format>.deletedFileRetentionDuration = "interval <interval>":現在のテーブルバージョンで参照されなくなったデータファイルを削除するためにVACUUMが使用するしきい値を決定します。デフォルトはinterval 7 daysです。
例えば、30日間のヒストリカルデータにアクセスするには、delta.deletedFileRetentionDuration = "interval 30 days"を設定します。これは、delta.logRetentionDurationのデフォルト設定と一致します。
データ保持のしきい値を増やすと、より多くのデータファイルが保持されるため、ストレージコストが増加する可能性があります。
テーブルプロパティは、テーブル作成時に指定するか、ALTER TABLEステートメントで設定することができます。「テーブルプロパティリファレンス」を参照してください。
タイムトラベルの例
ユーザー 111 によるテーブルへの誤った削除を修正するには:
INSERT INTO my_table
SELECT * FROM my_table TIMESTAMP AS OF date_sub(current_date(), 1)
WHERE userId = 111
テーブルへの偶発的な誤った更新を修正するには:
MERGE INTO my_table target
USING my_table TIMESTAMP AS OF date_sub(current_date(), 1) source
ON source.userId = target.userId
WHEN MATCHED THEN UPDATE SET *
先週追加された新規顧客の数をクエリするには:
SELECT
(
SELECT count(distinct userId)
FROM my_table
)
-
(
SELECT count(distinct userId)
FROM my_table TIMESTAMP AS OF date_sub(current_date(), 7)
) AS new_customers
トランザクションログチェックポイント
トランザクションログは、テーブルデータとともにトランザクションログディレクトリ内のJSONファイルとしてテーブルバージョンを記録します。
チェックポイントクエリを最適化するために、テーブルバージョンはParquetチェックポイントファイルに集約されます。これにより、テーブル履歴のすべてのJSONバージョンを読み取る必要がなくなり、パフォーマンスが向上します。ユーザーはチェックポイントと直接やり取りする必要はありません。
Databricksは、データサイズとワークロードに応じてチェックポイントの頻度を最適化します。チェックポイントの頻度は予告なく変更される場合があります。
テーブルを以前の状態に復元する
次のシナリオを含め、RESTOREコマンドを使用してテーブルを以前のバージョンまたはタイムスタンプに復元します。
- すでにリストアされたテーブルをリストアすることができます。
- クローンテーブルを復元できます。
次の要件を考慮してください:
- テーブルを復元するには、テーブルに対する
MODIFY権限が必要です。 - データファイルが手動で、または
VACUUMによって削除された後、それらのファイルを参照する古いバージョンにテーブルを復元することはできません。spark.sql.files.ignoreMissingFilesがtrueに設定されている場合、このバージョンへの部分的な復元は依然として可能です。 - タイムスタンプで復元するには、
yyyy-MM-dd HH:mm:ssまたはyyyy-MM-ddの形式を使用します。
RESTORE TABLE target_table TO VERSION AS OF <version>;
RESTORE TABLE target_table TO TIMESTAMP AS OF <timestamp>;
構文の詳細については、「RESTORE」を参照してください。
ストリーミング動作
復元は、データを変更するオペレーションであり、ダウンストリームのワークロードで重複データが発生する可能性があります。RESTORE コマンドによって追加されたログエントリーには、true に設定された dataChange が含まれています。
テーブルへの更新を処理する構造化ストリーミングジョブなどのダウンストリームワークロードの場合、復元操作によって追加されたデータ変更ログエントリは新しいデータ更新と見なされ、それらを処理するとデータが重複する可能性があります。
例えば:
テーブルバージョン | オペレーション | ログの更新 | データ変更ログ更新のレコード |
|---|---|---|---|
0 |
|
| (名前 = ヴィクトル、年齢 = 29)、(名前 = ジョージ、年齢 = 55) |
1 |
|
| (名前 = ジョージ、年齢 = 39) |
2 |
|
| レコードなし。 |
3 |
|
| (名前 = ヴィクトル、年齢 = 29)、(名前 = ジョージ、年齢 = 55)、(名前 = ジョージ、年齢 = 39) |
前述の例では、RESTORE コマンドは、テーブルバージョン0および1の読み込み時に以前に確認された更新をもたらします。ストリーミングクエリがこのテーブルを再度読み込むと、これらのファイルは新規追加データと見なされ、再度処理されます。
メトリクスを復元する
完了後、RESTORE は次のメトリクスを単一行のDataFrameとしてレポートします:
-
table_size_after_restore:復元後のテーブルのサイズ。 -
num_of_files_after_restore:復元後のテーブル内のファイルの数。 -
num_removed_files:テーブルから削除された(論理的に削除された)ファイルの数。 -
num_restored_files:ロールバックによって復元されたファイルの数。 -
removed_files_size:テーブルから削除されたファイルの合計サイズ(バイト単位)。 -
restored_files_size:復元されるファイルの合計サイズ(バイト単位)。
最後のコミットバージョンを検索
全スレッド、全テーブルにわたって、現在のSparkSessionが最後に書き込んだコミットのバージョン番号を取得するには、SQL設定spark.databricks.<format>.lastCommitVersionInSessionに問い合わせます。テーブルの形式に応じて、<format>をdeltaまたはicebergに置き換えます。
例えば:
- SQL
- Python
- Scala
SET spark.databricks.delta.lastCommitVersionInSession
spark.conf.get("spark.databricks.delta.lastCommitVersionInSession")
spark.conf.get("spark.databricks.delta.lastCommitVersionInSession")
SparkSessionによってコミットが行われていない場合、キーをクエリーすると空の値が返されます。
複数のスレッド間で同じSparkSessionを共有する場合、それは複数のスレッド間で変数を共有するのと似ています。構成値への並列更新で競合状態が発生する可能性があります。