본문으로 건너뛰기

증분 동기화

한 번 전체 복사한 뒤 바뀐 것만 반영하세요.

증분 동기화는 워크스페이스를 한 번 복사한 뒤, 연동이 돌아가는 동안 바뀐 것만 반영하는 방식입니다. 네 가지 규칙을 지켜야 정확하게 동작합니다. 하나라도 틀리면 동기화는 멀쩡해 보이지만 레코드를 조용히 잃어 갑니다.

전체 흐름

  1. 이벤트 워터마크(변경 피드에서 이어서 읽을 기준 sequence)를 기록합니다.
  2. 목록 엔드포인트를 페이지별로 읽어 전체를 복사합니다(백필).
  3. 그 뒤로는 워터마크부터 변경 피드를 읽습니다.
  4. updated_after는 복구할 때만 쓰고, 주 루프로는 쓰지 않습니다.

규칙 1: 워터마크는 백필 전에 기록하기

피드는 오래된 것부터 읽으므로, 워터마크는 볼 수 있는 가장 최근 이벤트의 sequence입니다. 피드를 끝까지 읽어 찾으세요.

콜드 스타트
import os

import requests

BASE = "https://tahoe.workonward.com/api/partner/v1"
session = requests.Session()
session.headers["Authorization"] = f"Bearer {os.environ['TAHOE_API_KEY']}"


def current_sequence():
    """The sequence of the newest event this key can see, or None if none."""
    after, last = None, None
    while True:
        params = {"limit": 100}
        if after is not None:
            params["after"] = after
        response = session.get(f"{BASE}/events", params=params, timeout=30)
        response.raise_for_status()
        body = response.json()
        if body["data"]:
            last = body["data"][-1]["sequence"]
        if not body["has_more"]:
            return last
        after = body["next_after"]


def cold_start(store):
    # 1. The watermark FIRST, held in memory for now.
    watermark = current_sequence()

    # 2. The backfill. Slow, and that is fine: it happens once.
    backfill_jobs(store)
    backfill_applicants(store)
    backfill_applications(store)

    # 3. Save the watermark only now. A crash during the backfill then
    #    restarts the whole cold start, instead of resuming from a point
    #    the copy never reached.
    store.save_sequence(watermark)

규칙 2: 상태는 sequence 번호입니다

진행 상황을 타임스탬프로 추적하지 마세요. 이벤트는 동시에 일어나는 변경들이 기록하므로, 시계 순서와 조금 다른 순서로 보일 수 있습니다. 그래서 “10:04 이후 전부”로 읽으면, 10:03으로 찍혔지만 10:05에 보이게 된 이벤트를 놓칠 수 있습니다. sequence 번호는 이벤트 순서를 정확히 나타내므로 이런 틈이 없습니다.

규칙 3: 워터마크는 작업 후에 저장하기

드레인 루프
def drain(store):
    after = store.load_sequence()

    while True:
        params = {"limit": 100}
        if after is not None:
            params["after"] = after
        response = session.get(f"{BASE}/events", params=params, timeout=30)
        response.raise_for_status()
        body = response.json()

        for event in body["data"]:
            # At-least-once delivery: the same event can arrive more than
            # once, so this check is not an optimization.
            if not store.already_processed(event["id"]):
                dispatch(store, event)
                store.mark_processed(event["id"])
            # AFTER the event is handled. Saving first turns a crash halfway
            # through a page into events that are skipped for good. Use the
            # event's own sequence: next_after is null on the last page.
            after = event["sequence"]
            store.save_sequence(after)

        if not body["has_more"]:
            return

먼저 저장하고 나중에 작업하면, 중간에 프로세스가 죽을 때 이벤트를 잃습니다. 먼저 작업하고 나중에 저장하면 이벤트가 다시 재생되는데, 핸들러가 두 번 실행해도 안전하다면 재생은 손해가 없습니다. 핸들러는 어차피 그렇게 만들어야 합니다. 장애가 없어도 전송은 최소 한 번(at least once) 방식이기 때문입니다.

규칙 4: 페이로드를 믿지 말고 리소스를 다시 읽기

이벤트 페이로드는 식별자와 간단한 값 몇 가지만 담습니다. 이벤트를 처리할 때 핸들로 리소스를 읽으세요.

사용하는 것결과
이벤트 페이로드를 데이터로 사용이벤트가 기록될 때 스코프가 허용한 내용을 계속 갖고 있게 됩니다. 나중에 유료 잠금 상태가 바뀌거나 삭제(소거)가 일어나도 오래된 사본이 그대로 남습니다.
핸들로 새로 읽기지금 스코프가 허용하는 내용을 받고, 삭제된 리소스는 404로 응답합니다. 이것이 올바른 응답입니다.

페이로드에서 바로 쓸 만한 부분이 하나 있습니다. 일부 이벤트에는 data.object.changed에 변경된 필드 목록이 담깁니다. 그중 내가 복사하는 필드가 하나도 없다면 다시 읽지 않아도 됩니다.

필요 없는 읽기 건너뛰기
# This system records outcomes only, not pipeline stages.
MIRRORED_APPLICATION_FIELDS = {"status"}


def on_application_event(store, event):
    changed = event["data"]["object"].get("changed")
    # When the event lists what changed and you copy none of it, skip the
    # read. When there is no list, read again to be safe.
    if changed and not set(changed) & MIRRORED_APPLICATION_FIELDS:
        return

    handle = event["data"]["object"]["id"]
    response = session.get(f"{BASE}/applications/{handle}", timeout=30)
    if response.status_code == 404:
        store.delete_application(handle)  # gone since the event was written
        return
    response.raise_for_status()
    store.upsert_application(response.json())

updated_after는 언제 쓰나요

채용 공고, 지원서, 지원자, 소싱한 프로필, 공유 인재풀 목록은 ?updated_after=를 받습니다. 이 파라미터는 평소 동기화가 아니라 복구에 쓰는 것입니다.

  • 데이터가 어긋났다고 의심될 때 특정 기간을 확인하려는 경우.
  • 내 스키마에 필드를 추가해서 값을 채워 넣어야 하는 경우.
  • 안전망이 필요할 때. 하루 한 번 지난 하루나 이틀을 확인하면, 피드로 받지 못한 것을 잡을 수 있습니다.

리더가 뒤처졌을 때

이벤트는 30일 동안 보관됩니다. 며칠 멈췄던 리더는 워터마크부터 피드를 읽어 따라잡으면 됩니다. 30일 넘게 멈췄던 리더는 다시는 받을 수 없는 이벤트를 잃은 것입니다. 먼저 삭제 통지를 따라잡은 다음 콜드 스타트를 다시 실행하고, 새 백필에서 돌아오지 않은 레코드는 모두 지우세요. 그 레코드의 삭제 이벤트가 놓친 이벤트 중에 있었을 수 있습니다.

웹훅, 그리고 그래도 피드를 유지해야 하는 이유

웹훅은 같은 이벤트를 같은 봉투에 담아 보내므로, 핸들러 하나로 둘 다 처리할 수 있습니다. 웹훅은 더 빨리 알기 위한 수단이지, 피드를 대신하는 것이 아닙니다.

전송은 실패할 수 있고, 실패가 반복되면 엔드포인트가 비활성화되며, 시도 횟수를 다 쓴 전송은 포기됩니다. 순서가 보장된 피드가 기준 데이터이자 유일한 복구 수단입니다. 드레인 루프를 남겨 두고 일정에 따라 안전망으로 실행하며, 워터마크도 계속 그 루프에서 저장하세요.

속도 조절

  • limit=100을 쓰세요. 요청 네 번이 한 번으로 줄어듭니다.
  • 키 하나에 워커 하나. 병렬 워커는 같은 한도를 나눠 쓰므로, 429 응답이 무작위로 나오는 것처럼 보이기 시작합니다.
  • Retry-After를 따르세요. 개인 데이터 한도의 경우 몇 시간이 될 수도 있으며, 믿어야 할 값은 이것입니다.
  • personal_data_reads_used_today를 지켜보세요. 백필하는 동안 GET /me에서 확인하고, 한도에 닿기 전에 계획적으로 멈추세요. 중간에 한도에 걸리면 무엇이 저장되었는지 추측해야 합니다.
  • 요청 ID를 기록하세요. Tahoe-Request-Id 헤더 값을 내 실행 기록과 함께 남겨 두면, 지원팀이 정확한 요청을 찾을 수 있습니다.

관련 문서