参考
メモ
Delta LakeのDelta LogのIDがいつ確定するのか、というのが気になり確認した。
前提
- Delta Lakeバージョン:0.8.0
createRelationから確認する
org.apache.spark.sql.delta.sources.DeltaDataSource#createRelation
をエントリポイントとする。
ポイントは、DeltaLog
がインスタンス化されるときである。
まず最初にインスタンス化されるのは以下。
org/apache/spark/sql/delta/sources/DeltaDataSource.scala:141
1 | val deltaLog = DeltaLog.forTable(sqlContext.sparkSession, path) |
org.apache.spark.sql.delta.DeltaLog は、
org.apache.spark.sql.delta.SnapshotManagement
トレイトをミックスインしている。 当該トレイトには、
currentSnapshot というメンバ変数があり、これは
org.apache.spark.sql.delta.SnapshotManagement#getSnapshotAtInit
メソッドを利用し得られる。
org/apache/spark/sql/delta/SnapshotManagement.scala:47
1 | protected var currentSnapshot: Snapshot = getSnapshotAtInit |
このメソッドは以下のように定義されている。
org/apache/spark/sql/delta/SnapshotManagement.scala:184
1 | protected def getSnapshotAtInit: Snapshot = { |
ポイントは、スナップショットを作る際に用いられるセグメントである。 セグメントにバージョン情報が持たれている。
ここでは3行目の
1 | val segment = getLogSegmentFrom(lastCheckpoint) |
にて
org.apache.spark.sql.delta.SnapshotManagement#getLogSegmentFrom
メソッドを用いて、
前回チェックポイントからセグメントの情報が生成される。
なお、参考までにLogSegmentクラスの定義は以下の通り。
org/apache/spark/sql/delta/SnapshotManagement.scala:392
1 | case class LogSegment( |
上記の通り、コンストラクタ引数にバージョン情報が含まれていることがわかる。
インスタンス化の例は以下の通り。
org/apache/spark/sql/delta/SnapshotManagement.scala:140
1 | LogSegment( |