증분 동기화
한 번 전체 복사한 뒤 바뀐 것만 반영하세요.
증분 동기화는 워크스페이스를 한 번 복사한 뒤, 연동이 돌아가는 동안 바뀐 것만 반영하는 방식입니다. 네 가지 규칙을 지켜야 정확하게 동작합니다. 하나라도 틀리면 동기화는 멀쩡해 보이지만 레코드를 조용히 잃어 갑니다.
전체 흐름
- 이벤트 워터마크(변경 피드에서 이어서 읽을 기준 sequence)를 기록합니다.
- 목록 엔드포인트를 페이지별로 읽어 전체를 복사합니다(백필).
- 그 뒤로는 워터마크부터 변경 피드를 읽습니다.
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헤더 값을 내 실행 기록과 함께 남겨 두면, 지원팀이 정확한 요청을 찾을 수 있습니다.