差分同期でデータを欠落させない設計|watermark・checkpoint・削除伝播

「差分同期でデータを欠落させない設計|watermark・checkpoint・削除伝播」の内容を表す技術イラスト

データ連携を毎回の全件コピーから差分同期へ変えると、処理時間やAPI負荷を抑えられます。一方で、updated_at > 前回実行時刻だけを条件にすると、同時刻に更新された複数レコード、処理途中の失敗、遅れて確定した更新、削除済みデータを取りこぼすことがあります。

この記事では、watermark、checkpoint、複合キー、overlap、tombstoneを組み合わせ、ETL・API・SQLite・Pythonで再実行できる差分同期を設計します。

目次

差分同期を構成する4つの要素

差分同期では、似た役割の値を分けて考えます。

要素役割代表例
watermarkどこまでの変更を読み取ったかを表す境界(updated_at, id)
checkpoint正常完了したwatermarkを次回用に永続化した値2026-09-30T01:30:00Zとasset-0042
tombstone削除された事実を連携するデータdeleted_at、削除イベント
overlap境界の少し前から読み直す安全幅前回時刻の5分前など

watermarkは処理中にも進みますが、checkpointは出力先への反映が成功した位置だけを保存します。読み取った直後にcheckpointを進めると、その後のDB更新が失敗したときに未反映データを読み飛ばします。

最小の差分レコード

差分同期の入力には、安定したID、変更順を判断できる値、削除状態を含めます。

{
  "id": "asset-0042",
  "updated_at": "2026-09-30T01:30:00Z",
  "version": 7,
  "deleted_at": null,
  "data": {
    "name": "Hydraulic unit A",
    "material": "SS400"
  }
}
  • id:同じ対象を再送時にも特定する文字列
  • updated_at:送信元で変更が確定した日時
  • version:対象ごとに増加する変更版
  • deleted_at:有効ならnull、削除済みなら削除日時
  • data:業務上のfieldを持つobject

日時はRFC 3339形式など、オフセットを含む交換形式に固定します。ただし、時計だけで完全な順序を保証できるとは限りません。送信元が変更sequence、revision、opaqueな同期tokenを提供するなら、その契約を優先します。

時刻だけのwatermarkでは欠落する

前回のcheckpointが2026-09-30T01:30:00Zだったとします。同じ時刻を持つレコードが複数あり、batchの途中で区切られた場合、次回を次の条件で始めると未処理レコードを除外します。

WHERE updated_at > :last_time

そこで、日時と一意なIDを組み合わせた複合watermarkを使います。取得順と再開条件を同じ組み合わせにすることが重要です。

SELECT id, name, material, version, updated_at, deleted_at
FROM source_assets
WHERE updated_at > :last_time
   OR (updated_at = :last_time AND id > :last_id)
ORDER BY updated_at ASC, id ASC
LIMIT :batch_size;

このときcheckpointはupdated_atだけでなくidも保存します。IDの大小関係は送信元DBの照合規則に依存するため、API・DB・実装で同じ並び順を使います。

ページcursorとcheckpointを分ける

ページネーションのcursorと差分同期のcheckpointは、どちらも「続き」を示しますが寿命が異なります。

値有効範囲更新する時点
page cursor一回の一覧走査や短期的な再開次ページを取得するとき
checkpoint日次・定期実行をまたぐ同期位置出力先への反映が確定した後

一回の同期では、checkpoint以降を複数ページに分けて取得できます。各ページを安全に反映してからcheckpointを進めるか、走査全体を一つの再実行単位にします。APIのcursorが期限切れになる仕様なら、永続checkpointとしてそのまま保存できるとは限りません。

動く境界を固定する

同期中にも送信元が更新され続けると、一覧の終端が動きます。実行開始時に上限となるhigh-water markを固定し、今回の範囲を次のように定義します。

前回checkpointより後、かつ今回のhigh-water mark以下

上限には、送信元が提供するsnapshot ID、変更sequence、整合性が保証された時刻などを使います。単にクライアントPCの現在時刻を使うと時計差が入るため、送信元の仕様に合わせます。

今回の上限を越えた変更は次回へ回します。これにより、実行中に追加されたデータを追い続けて処理が終わらない状態や、ページ間で境界が揺れる問題を避けられます。

反映とcheckpointを同じtransactionで確定する

差分処理は、次の責務に分けます。

段階主な処理失敗時の扱い
取得checkpoint以降を決定順で読むcheckpointを変更しない
検証ID、型、日時、version、削除状態を確認batchを隔離または中止
変換入力schemaを保存schemaへ正規化元データと変換版を記録
反映IDをkeyにUPSERTまたは削除反映transactionをROLLBACK
確定同じtransactionでcheckpointを更新COMMIT後だけ次回へ進む

SQLiteのtransactionは、BEGIN、COMMIT、ROLLBACKで変更単位を明示できます。対象データだけをCOMMITしてcheckpoint更新に失敗すると再送が起きますが、UPSERTが冪等なら同じ結果へ戻せます。逆にcheckpointだけが先へ進むと欠落するため、両方を同じtransactionへ含めます。

SQLiteの保存schema

出力データと同期位置を別tableにします。ここでは削除後もtombstoneを保持し、通常検索ではdeleted_at IS NULLを条件にします。

CREATE TABLE assets (
    id             TEXT PRIMARY KEY NOT NULL,
    name           TEXT,
    material       TEXT,
    source_version INTEGER NOT NULL,
    updated_at     TEXT NOT NULL,
    deleted_at     TEXT
);

CREATE TABLE sync_checkpoints (
    pipeline   TEXT PRIMARY KEY NOT NULL,
    updated_at TEXT NOT NULL,
    record_id  TEXT NOT NULL
);

source_versionは、遅れて届いた古い変更が新しい状態を巻き戻すのを防ぎます。削除済みレコードでもID、version、日時を残せば、古い更新の再送による復活を防止できます。

Pythonで再実行可能なbatch反映を作る

次の例は、同じbatchを再実行してもレコード数を増やさず、古いversionを上書きせず、データ反映とcheckpoint更新を一つのtransactionで確定します。

import sqlite3
from datetime import datetime


def parse_rfc3339(value: str) -> datetime:
    if not isinstance(value, str) or not value.endswith("Z"):
        raise ValueError("日時はUTCのRFC 3339文字列が必要です")
    return datetime.fromisoformat(value.replace("Z", "+00:00"))


def validate_row(row: dict) -> None:
    if not isinstance(row.get("id"), str) or not row["id"]:
        raise ValueError("idが不正です")
    if isinstance(row.get("version"), bool) or not isinstance(row.get("version"), int):
        raise ValueError("versionは整数が必要です")
    parse_rfc3339(row.get("updated_at"))
    if row.get("deleted_at") is not None:
        parse_rfc3339(row["deleted_at"])
    if not isinstance(row.get("data"), dict):
        raise ValueError("dataはobjectが必要です")


def apply_batch(
    conn: sqlite3.Connection,
    pipeline: str,
    rows: list[dict],
) -> None:
    ordered = sorted(rows, key=lambda row: (row["updated_at"], row["id"]))
    for row in ordered:
        validate_row(row)

    if not ordered:
        return

    with conn:
        for row in ordered:
            data = row["data"]
            conn.execute(
                """
                INSERT INTO assets (
                    id, name, material, source_version,
                    updated_at, deleted_at
                ) VALUES (?, ?, ?, ?, ?, ?)
                ON CONFLICT(id) DO UPDATE SET
                    name = excluded.name,
                    material = excluded.material,
                    source_version = excluded.source_version,
                    updated_at = excluded.updated_at,
                    deleted_at = excluded.deleted_at
                WHERE excluded.source_version > assets.source_version
                """,
                (
                    row["id"],
                    data.get("name"),
                    data.get("material"),
                    row["version"],
                    row["updated_at"],
                    row.get("deleted_at"),
                ),
            )

        last = ordered[-1]
        conn.execute(
            """
            INSERT INTO sync_checkpoints (
                pipeline, updated_at, record_id
            ) VALUES (?, ?, ?)
            ON CONFLICT(pipeline) DO UPDATE SET
                updated_at = excluded.updated_at,
                record_id = excluded.record_id
            WHERE excluded.updated_at > sync_checkpoints.updated_at
               OR (
                    excluded.updated_at = sync_checkpoints.updated_at
                    AND excluded.record_id > sync_checkpoints.record_id
               )
            """,
            (pipeline, last["updated_at"], last["id"]),
        )

SQLiteのUPSERTは、一意性制約との衝突時にUPDATEまたは何もしない処理へ切り替えられます。この例ではidが同じ場合、受信versionが新しいときだけ更新します。checkpointにも複合watermarkの比較条件を付け、古いbatchの手動再実行で同期位置が後退しないようにしています。

入力検証をtransaction前に行っていますが、実際の変換やDB制約違反がtransaction内で起きても、Pythonのsqlite3 connection context managerによりROLLBACKされます。再試行時は旧checkpointから同じ範囲を読み直します。

削除は通常の更新とは別に設計する

物理削除された行は、updated_atを条件にした一覧から消えます。存在しないデータを見ても、受信側は「削除された」のか「今回の検索範囲外」なのか判断できません。

代表的な方法は次の3つです。

  • soft delete:元recordへdeleted_atを設定し、通常の差分として送る
  • tombstone event:ID、削除日時、versionだけを持つ削除イベントを送る
  • reconciliation:定期的にID一覧やsnapshotを照合し、消失を検出する

Google AIP-164はsoft deleteの設計例として、物理削除せず削除状態を付け、delete_timeやpurge_timeを持つ方法を示しています。すべてのAPIに同じfieldを要求する規格ではありませんが、削除事実をデータとして保持する参考になります。

tombstoneを一定期間後に消す場合、その保持期間は想定する最大停止時間と再試行期間より長くします。受信側が停止中にtombstoneまで消えると、復旧後に削除を知る方法がなくなります。

overlapで遅延到着と精度差を吸収する

DB commitの遅延や日時精度の違いがあると、checkpoint直後の検索だけでは境界付近を取りこぼす場合があります。そこで、前回時刻より少し前から再読込します。

取得開始 = checkpoint時刻 - overlap幅

overlapで同じrecordが再取得されるため、IDとversionによるUPSERTが前提です。幅を広げるほど安全余裕は増えますが、再処理量も増えます。実測した遅延、送信元の時刻精度、障害復旧時間から決めます。

ただしoverlapは、安定した変更sequenceやsnapshotの代わりにはなりません。送信元が更新日時を後から変更できる、時計が逆行する、古い日時の行が新規追加される場合は、変更ログやCDCなど別の仕組みが必要です。

正常系と異常系をテストする

ケース期待結果
初回batch全件反映後にcheckpointが末尾へ進む
同一batchを再実行件数と最終状態が変わらない
同一時刻に複数ID(updated_at, id)順ですべて処理される
古いversionが遅れて到着最新recordを上書きしない
tombstoneを受信削除状態を保持し、通常検索から除外する
batch途中でDB error全変更をROLLBACKし、checkpointを進めない
未知field・型違反推測変換せず隔離または失敗にする

checkpointの値だけでなく、再試行後の件数、各version、削除状態まで検証します。障害試験では、データ更新後かつcheckpoint更新前に例外を発生させ、原子性を確認します。

セキュリティと監査

checkpointやcursorは同期位置であり、認証情報ではありません。送信元APIのtoken、DB credential、同期設定とは分離してsecret管理機能へ保存します。

差分payloadには、変更前後の値や削除済み情報が含まれる場合があります。ログには必要最小限のID、version、処理結果を残し、本文全体やcredentialを出力しません。また、入力件数、文字列長、許可field、日時範囲を検証し、SQLはparameter bindingで実行します。

監査用には、pipeline名、実行ID、開始checkpoint、終了checkpoint、high-water mark、取得件数、追加・更新・削除・拒否件数、変換versionを記録すると、欠落調査と再処理が容易になります。

CAD・BOM・設計データへ展開する

図面、部品、BOM、計算条件では、業務revisionと同期用versionを分けます。図面revisionがCでも、名称修正や公開範囲変更により同期versionは複数回進むことがあります。

親子関係を持つBOMでは、親より子が先に届く場合や、削除が参照先へ影響する場合があります。取得順だけに依存せず、一時tableへ受けてから参照整合性を確認する、削除済みIDをtombstoneとして保持する、依存関係の再構築を冪等にする、といった設計が必要です。

watermark、checkpoint、source version、tombstoneを共通データ契約にすれば、API、CSV、DB、CAD抽出、定期ETLで同じ再実行ルールを共有できます。

まとめ

差分同期の品質は、検索条件だけでなく「失敗後にどこから安全にやり直せるか」で決まります。

  • watermarkは日時だけでなく一意なIDと組み合わせる
  • 取得順と再開条件を同じ複合キーにする
  • page cursorと実行間checkpointを分ける
  • 出力反映とcheckpoint更新を同じtransactionで確定する
  • IDとsource versionによるUPSERTで再実行可能にする
  • soft deleteまたはtombstoneで削除を伝播する
  • overlapは遅延吸収に使い、変更ログの代替にしない
  • high-water markで今回の処理範囲を固定する
  • credential、同期位置、業務payloadを分離する

この契約を先に固めれば、処理途中で停止しても旧checkpointから再開でき、同じ差分を読み直しても同じ状態へ収束します。全件同期を減らしながら、更新・削除・再送を欠落なく扱える連携基盤になります。

参考情報

参考になったらシェアしてください
  • URLをコピーしました!
  • URLをコピーしました!

この記事を書いた人

機械設計・油圧・CAD・Python・AIなど、ものづくりに関わる技術を扱っています。工学知識を整理・構造化し、設計や自動化に再利用できる形へ変えていくことを目指しています。

目次