0
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?

Spark 4.2のAuto CDCはDeltaに書けない。止めているのはif文1つ

0
Posted at

fig0-question.png

Bronze→Silver→Goldのうち、Silverに書く直前だけ検査が1つ挟まります。
4つの書き込み先のどれが通るかを見てください。

Databricksでメダリオンを組んでいて、同じ構成がOSSのSparkでどこまで成立するのか気になっている人向けです。
Bronze→Silverの中核にあるCDC処理は、2026-07-14公開のApache Spark 4.2.0でOSS側に入りました。
入ってはいるものの、それを受けられるテーブル形式が今のところ1つもありません。

データ分析基盤をDatabricksで作ったときに、Silver(ウェアハウス層)の取り扱いで悩んだことがあります。

食い違いは文書の時点から始まっていて、DatabricksのAUTO CDCのドキュメントには、いまも「The AUTO CDC APIs are not supported by Apache Spark Declarative Pipelines.」という1文が残っています。
一方でApache Spark 4.2.0のリリースノートは「Auto CDC in Spark Declarative Pipelines for declarative SCD Type 1 upserts」を目玉に挙げていて、実装はSPARK-56249。
手元のSpark 4.2.0で実行すると、どちらとも違う3つ目の答えが返ってきます。

Databricksから離れる予定が無いなら版まわりの話は当たらないので、Auto CDCがSilverテーブルに足す隠し列のところだけ拾って、残りは読み飛ばしてください。

検証環境は次のとおりです。

項目 版
実行環境 WSL2 / Ubuntu 24.04.4 LTS(kernel 6.6.87.2) / Temurin JDK 21.0.12+8 / Python 3.12.3
Spark 4.2.0(対照として4.1.3)
コネクタ delta-spark_2.13 4.3.1 / iceberg-spark-runtime-4.1_2.13 1.11.0

Spark Declarative Pipelines(以下SDP)のCLIはspark-pipelinesで、中身はSpark Connectのクライアントです。
PySpark本体に加えてgrpcio・pyarrow・pandas・zstandardが要りますが、どれもpipで入ります。

Auto CDCは動く。ただしテスト用のカタログでだけ

Bronzeはストリーミングテーブル、SilverはAuto CDCによるSCD Type 1、Goldはそれを読むマテリアライズドビューとして素直に書きます。

from pyspark import pipelines as dp                        # Spark 4.2.0 の SDP
from pyspark.sql import SparkSession
from pyspark.sql.types import (StructType, StructField, IntegerType,
                               StringType, LongType)

spark = SparkSession.active()

SCHEMA = StructType([                                      # CDC フィードのスキーマ
    StructField("order_id", IntegerType()),
    StructField("status", StringType()),
    StructField("amount", LongType()),
    StructField("seq", LongType()),                        # 順序を決める列
    StructField("op", StringType()),                       # UPSERT / DELETE
])

@dp.table(name="testcat.db1.bronze_orders")                # Bronze:来たものをそのまま貯める
def bronze_orders():
    return spark.readStream.schema(SCHEMA).json("/path/to/src")

dp.create_streaming_table("testcat.db1.silver_orders")     # Silver:書き込み先を先に宣言する

dp.create_auto_cdc_flow(                                   # Silver への CDC 適用はこれ 1 回
    target="testcat.db1.silver_orders",
    source="testcat.db1.bronze_orders",
    keys=["order_id"],                                     # 行を一意に決めるキー
    sequence_by="seq",                                     # この列の大きいほうを採用する
    apply_as_deletes="op = 'DELETE'",                      # 削除イベントの判定式
    except_column_list=["op", "seq"],                      # 出力に残さない列
    stored_as_scd_type=1,
)

@dp.materialized_view(name="spark_catalog.default.gold_orders")   # Gold:ここは普通の MV
def gold_orders():
    return spark.read.table("testcat.db1.silver_orders")

testcatはSparkのテスト用jarに入っているSharedTablesInMemoryRowLevelOperationTableCatalogで、DSv2のSupportsRowLevelOperationsを実装しています。
本物のテーブル形式ではなくこれを使っている理由は、次の節。

投入するCDCフィードは6行で、order_id=1についてはseq=3のSHIPPEDをseq=2のPAIDより前に置いて順序を崩してあります。
order_id=2はseq=2で削除され、order_id=3はseq=1で作られたまま更新されません。

spark-pipelines run --jars spark-catalyst_2.13-4.2.0-tests.jar,spark-sql_2.13-4.2.0-tests.jar \
  --conf spark.driver.host=127.0.0.1
2026-08-02 13:32:03: Flow testcat.db1.bronze_orders has COMPLETED.
2026-08-02 13:32:03: Flow testcat.db1.silver_orders has COMPLETED.
2026-08-02 13:32:04: Flow spark_catalog.default.gold_orders has COMPLETED.
2026-08-02 13:32:05: Run is COMPLETED.

Goldに落ちたparquetを開くと、SCD Type 1として正しい2行が入っています。

{'order_id': 1, 'status': 'SHIPPED', 'amount': 100,
 '__spark_autocdc_metadata': {'deleteSequence': None, 'upsertSequence': 3}}
{'order_id': 3, 'status': 'NEW', 'amount': 300,
 '__spark_autocdc_metadata': {'deleteSequence': None, 'upsertSequence': 1}}

ファイルの並び順ではなくseqの大きいほうが採用されていて、削除も効いています。
目を引くのは__spark_autocdc_metadataという増えた列で、upsertSequenceとdeleteSequenceを持っています。
自分で書くならターゲット側に足すことになる順序判定用の列を、Auto CDCは隠し列として自前で抱える設計。
ここは素直に良い作りだと思いました。

fig1-silver-shape.png

図の右側が実際に出来たSilverの1行で、ユーザーが定義した列の後ろに構造体が1つ足されているところを見てください。

定義では1箇所つまずいていて、spark-pipeline.ymlにcatalog: testcatとdatabase: db1を書いたところ[SCHEMA_NOT_FOUND] The schema testcat.db1 cannot be found.で止まりました。
このカタログはCREATE TABLE testcat.db1.t1をネームスペース無しで受け付けるのに、SDPは起動時にデフォルトスキーマの存在だけ先に確認します。
同じ組み合わせを試すなら、specからcatalogとdatabaseを落として、データセット名を完全修飾で書くほうが早いです。

止めているのはif文1つ

同じパイプラインを、テーブル名からtestcat.db1.を外してSparkのデフォルトカタログに向けます。
spark-warehouseにparquetで書かれる、いちばん普通の構成です。

org.apache.spark.sql.AnalysisException: [AUTOCDC_TARGET_DOES_NOT_SUPPORT_MERGE]
Cannot start AutoCDC flow: the target table `spark_catalog`.`default`.`silver_orders`
(format: parquet) does not support row-level operations.
AutoCDC requires a target backed by a connector that supports MERGE. SQLSTATE: 0A000
	at org.apache.spark.sql.pipelines.graph.AutoCdcMergeWriteBase
	    .requireDestinationSupportsRowLevelOps(FlowExecution.scala:493)
	at org.apache.spark.sql.pipelines.graph.Scd1MergeStreamingWrite.<init>(FlowExecution.scala:623)

parquetがMERGEできないのは当たり前なので、ここまでは筋書き通りです。
引っかかったのは検査の名前のほうで、requireDestinationSupportsRowLevelOpsは「MERGEできるか」ではなく「SupportsRowLevelOperationsというDSv2のインターフェースを実装しているか」を見ています。
スタックトレースのScd1MergeStreamingWrite.<init>が示すとおり、この検査が走るのはストリーミング書き込みを組み立てる前、つまりデータが1行も流れる前。

DSv2には行レベル操作の作法が定められていて、SupportsRowLevelOperationsを実装したテーブルはRowLevelOperationTableでくるまれ、RewriteMergeIntoTableというアナライザルールがMERGEを書き換える経路に乗ります。
Sparkから見れば、この線の内側にいるものだけがMERGEできるテーブル。
Auto CDCの検査はその線をそのまま流用しています。

線の外側にもMERGEできるテーブルがあるところが、この記事でいちばん時間を使った部分です。
DeltaのMERGEはDSv2の行レベル操作ではなく、Delta自身がセッションに差し込むアナライザルールで実現されていて、spark.sql.extensions=io.delta.sql.DeltaSparkSessionExtensionと設定するあの1行がその差し込み口になっています。
機能としてはMERGEできるのに、インターフェースの上では「行レベル操作に対応していないテーブル」に見えます。

fig2-gate.png

図の縦の破線がAuto CDCの検査が見ている境界で、下側の経路はMERGEに届いているのに検査は上側しか通しません。

この検査はMERGEできるかどうかではなく、MERGEをどうやって実装したかで通す相手を選んでいます。
Delta側が作法から外れているという言い方もできますが、DSv2の行レベル操作が入るより前からDeltaのMERGEは動いていて、拡張で差し込む形はSparkが公式に用意した口です。
どちらかが規約違反というより、2つの正しいやり方が並んでいて、新しい検査が片方しか知らない形になっています。

同じ検査を3つの実装に当てる

インターフェースの実装はjarのバイトコードを読めば分かるので、.classの定数プールに文字列があるかを見たうえで、直接実装しているクラスだけを数えました。

jar インターフェースへの参照 直接実装しているクラス
Spark 4.2.0 配布jar(全ファイル) あり 0
Spark 4.2.0 テストjar あり 1(InMemoryRowLevelOperationTable)
Delta Lake 4.3.1 無し 0
Iceberg 1.11.0(spark-runtime 4.1) あり 1(SparkTable)

DeltaのjarはSupportsRowLevelOperationsという文字列を1度も含まず、DeltaTableV2が実装しているのはTable・SupportsWrite・V2TableWithV1Fallbackの3つだけです。
Sparkの配布jarに至っては、検査を通せる本番用のテーブルが1つも入っていません。
枠組みだけがあって、実装はテストjarの中にあります。

Icebergは実装しているので道が開けそうなところ、Maven Centralにiceberg-spark-runtime-4.2_2.13がまだありません。

iceberg-spark-runtime-4.0_2.13  HTTP 200
iceberg-spark-runtime-4.1_2.13  HTTP 200
iceberg-spark-runtime-4.2_2.13  HTTP 404

Deltaも同じで、公開されている最新の4.3.1はPOMでspark-sql_2.13の4.1.0をprovidedに取っていて、Spark 4.2向けのビルドは出ていません。
実際にSpark 4.2.0へ載せると、テーブルを作る時点で落ちます。

java.lang.NoSuchMethodError: 'org.apache.spark.sql.catalyst.catalog.CatalogStorageFormat
 org.apache.spark.sql.catalyst.catalog.CatalogStorageFormat.copy(scala.Option, scala.Option,
 scala.Option, scala.Option, boolean, scala.collection.immutable.Map)'

同じSQLをSpark 4.1.3で流すと通るので、これはコネクタ側の対応待ちです。

手元のコネクタがゲートを通るかは、jarを1つ渡せば分かります。
追試するなら--packagesが落としたjarを~/.ivy2*/cache/から拾ってください。

import sys, zipfile                                        # jar は zip なので展開せずに読める

IFACE = b"org/apache/spark/sql/connector/catalog/SupportsRowLevelOperations"

for jar in sys.argv[1:]:
    z = zipfile.ZipFile(jar)
    hit = any(IFACE in z.read(n) for n in z.namelist() if n.endswith(".class"))
    print(f"{'通る' if hit else '通らない'} {jar.split('/')[-1]}")

版が合うかはcurl -s .../delta-spark_2.13-4.3.1.pom | grep -A2 spark-sql_2.13で足りて、<version>4.1.0</version>が出てくるのでNoSuchMethodErrorを踏む前に判断できます。

上流も「厳しすぎる」と書いている

Auto CDCが使えるのは、今のところSparkのテストjarに入っているテーブルだけです。
その状態でJIRAを見にいくと、同じ壁の話が先に書かれていました。

SPARK-57623の題は「Drop AutoCDC requireDestinationSupportsRowLevelOps Check」で、本文にはDSv2の契約ではなく独自のアナライザルールでMERGEを実装しているコネクタ(Deltaが名指しされています)があるため、この検査は不必要に厳しい、と書いてあります。
提案されているのは検査を落として実行時の失敗に任せることで、起票は2026-06-22、affects versionは4.3.0、この記事を書いている時点では未解決のままです。

つまりDelta側がDSv2に寄せてくるのを待つ話ではなく、Spark側が検査を緩める話として進んでいます。
Deltaが4.2に対応しただけでは、このif文が残っている限りAUTOCDC_TARGET_DOES_NOT_SUPPORT_MERGEは消えません。
塞がる穴が2つあって、両方が要ります。

SCD Type 2のほうも見ておくと、4.2.0のPython APIはstored_as_scd_typeをLiteral[1, "1"]しか受け付けず、2を渡した時点で弾かれます。

pyspark.errors.exceptions.base.PySparkTypeError: [NOT_EXPECTED_TYPE]
Argument `stored_as_scd_type` should be Literal[1, '1'], got int.

値が範囲外なのに型の名前で怒られるので、メッセージだけでは何を直せばいいのか分かりません。
Type 2そのものはSPARK-58247で2026-07-23に解決済みで、fix versionは4.3.0。
DatabricksがType 1とType 2の両方に加えてSQLのAUTO CDC INTOまで持っているのに対し、OSS側は今のところType 1のPython APIだけですが、3週間前に入ったばかりなので順番として妥当な後回しに見えます。

今日Silverを組むなら、手書きのMERGEになる

待てない場合の現実解はSpark 4.1.3とDelta 4.3.1で、CDCの適用は自分で書きます。
Auto CDCがcreate_streaming_tableとcreate_auto_cdc_flowの10行で済ませる範囲を数えたところ、手書きは31行でした(空行とコメントを除いた行数)。

SILVER = f"{WAREHOUSE}/silver_orders"
(DeltaTable.createIfNotExists(spark)                        # 順序判定用の列は自分で持つ
 .location(SILVER)
 .addColumn("order_id", "INT").addColumn("status", "STRING")
 .addColumn("amount", "BIGINT").addColumn("_seq", "BIGINT")
 .execute())

def upsert(batch_df, batch_id):                             # マイクロバッチごとに呼ばれる
    newest = Window.partitionBy("order_id").orderBy(F.col("seq").desc())
    latest = (batch_df                                      # 同一バッチ内の順序崩れをここで潰す
              .withColumn("_rn", F.row_number().over(newest))
              .filter("_rn = 1").drop("_rn"))
    (DeltaTable.forPath(spark, SILVER).alias("t")
     .merge(latest.alias("s"), "t.order_id = s.order_id")
     .whenMatchedDelete(condition="s.op = 'DELETE' AND s.seq > t._seq")
     .whenMatchedUpdate(condition="s.op <> 'DELETE' AND s.seq > t._seq",
                        set={"status": "s.status", "amount": "s.amount", "_seq": "s.seq"})
     .whenNotMatchedInsert(condition="s.op <> 'DELETE'",
                           values={"order_id": "s.order_id", "status": "s.status",
                                   "amount": "s.amount", "_seq": "s.seq"})
     .execute())

(spark.readStream.schema(SCHEMA).json(SRC)                  # Bronze を読む
 .writeStream.foreachBatch(upsert)                          # MERGE はバッチ関数の中でしか書けない
 .option("checkpointLocation", CKPT)
 .trigger(availableNow=True).start().awaitTermination())

同じ6行を流すとAuto CDCと同じ2行になるので、order_id=1がSHIPPED(seq=3)になったあとにseq=2のPAIDを遅れて投入しました。

実装 1回目 遅れてseq=2のPAIDが到着
Auto CDC(Spark 4.2.0) SHIPPED / upsertSequence=3 SHIPPED / upsertSequence=3
手書き(s.seq > t._seqあり) SHIPPED / _seq=3 SHIPPED / _seq=3
手書き(条件をtrueにしたもの) SHIPPED / _seq=3 PAID / _seq=2

3つ目は、バッチ内の重複排除だけ書いてバッチ間の比較を忘れた場合です。
row_number()で最新1件に絞る処理は自然に思いつくのに対して、既存行より古いかの判定はターゲット側に列を足さないと書けないので、抜けやすいのはこちら。
whenMatchedUpdateとwhenMatchedDeleteの両方に条件が要る点も見落としやすい箇所です。

MERGE1回の値段は、テーブルの大きさでは決まらない

条件を足すぶんのコストも測っていて、ここで一度、測り方を間違えました。
guardedとnaiveを常にguarded→naiveの順で実行していたところ、比の中央値は0.89、つまり条件があるほうが速いという結果。
筋が通らないので順番を疑い、ラウンドごとに実行順を入れ替えて50回取り直すと、原因が出ました。

guarded を先に実行した回: 比の中央値 0.95(0.76〜1.22、n=25)
naive   を先に実行した回: 比の中央値 0.76(0.53〜0.90、n=25)

後に実行したほうが2割ほど速くなるので、この精度では条件のコスト自体を測れていません。
言えるのは、大きくはない、というところまで。

規模のほうは素直で、Silverの行数を1,000から1,000万まで4桁振っても、3行のCDCバッチ1回の中央値は647〜762msに収まります(各10回、単発では565〜1,331msまで振れます)。

fig3-scale.png

横軸は1目盛りごとに行数が10倍になります。
右へ4つ進んでも棒の高さが変わっていないところを見てください。

理由はDESCRIBE HISTORYのoperationMetricsに出ていて、どの規模でも書き換わるファイルは1〜2個です。
行数を100万に固定してファイル数だけを10から1,992まで振っても、中央値は462〜873msでした。

だから一番損なのは、CDCフィードを細かく分けてforeachBatchを何度も呼ぶ組み方です。
1万行でも1,000万行でも1回あたり0.7秒前後を払うので、バッチを10分割すれば7秒に増えます。
テーブルが育っただけでSilverの更新が重くなることは、この範囲では起きません。

所感

3つの情報源のうち、いちばん実態に近かったのはDatabricksのドキュメントでした。
APIは入っているのに使える書き込み先が無いので、「サポートされていない」という記述は更新漏れというより、まだ間違っていません。

requireDestinationSupportsRowLevelOpsまで降りて分かったのは、塞いでいるのが機能の不足ではないことです。
MERGEできるテーブルは2種類あって、新しく入った検査が片方の作法しか知らない。
待つ対象はDeltaの4.2対応だけだと思って調べ始めたので、Spark側のif文も同時に外れないと通らないと分かった時点で見積もりを書き直しました。

測れていないことも書いておきます。
ゲートを通る唯一のテーブルがインメモリのテスト実装なので、Auto CDC側の性能はまったく測っていません。
扱ったCDCフィードは6行でキーの重複も1件だけなので、数百万行のフィードでrow_number()のシャッフルが効いてくる領域は見ていません。
SCD Type 2は4.2.0に無いので動かしておらず、挙動の確認は4.3.0が出てからです。

Silver層の話は設計論として語られることが多く、3つの箱を描くところまでは誰でもできます。
箱1つを実際に動かそうとすると、通るか通らないかを決めていたのはインターフェースの実装方法という、図には出てこない粒度でした...

0
0
0

Register as a new user and use Qiita more conveniently

  1. You get articles that match your needs
  2. You can efficiently read back useful information
  3. You can use dark theme
What you can do with signing up
0
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?