最終更新日:24年2026月XNUMX日 シーザー・フィクソン
直接的な答え: 信頼性の高いウェブフックまたはアフィリエイトポストバック配信サービスには、耐久性のあるキュー、ジッター付きの制限付き再試行、エンドポイントごとのサーキットブレーカー、重複防止、署名検証、および配信が失敗する前に速度が低下しているかどうかを示すテレメトリが必要です。
この実装では、実行モデルを意図的に小規模に抑えています。Pythonワーカーと、キュー、配信履歴、エンドポイントの状態管理にPostgreSQLを使用します。B2B SaaS、iGamingアフィリエイト事業、およびHTTPタイムアウトをコンバージョン損失とみなさずにアウトバウンドイベントを配信する必要のあるあらゆる製品にとって、実用的な出発点となります。
配送契約で保証しなければならないこと
| 管理 | 重要性 | 動作確認 |
|---|---|---|
| 耐久性のあるイベント記録 | イベントはワーカーの再起動後も継続します。 | 受理されたすべての配送には、安定したIDとステータスが付与されます。 |
| 署名済みの依頼書 | 受信者は送信者を確認し、本体の改ざんを検出できる。 | タイムスタンプ付きの署名を使用し、期限切れのリクエストを拒否してください。 |
| 重複防止機能 | 再試行を行うと、重複した変換や更新が発生する可能性があります。 | イベントIDを送信し、受信側を冪等にする。 |
| 制限付き再試行 | 一時的な障害は、劣化したエンドポイントに過負荷をかけることなく復旧する。 | ジッターを徐々に減らし、規定の試行回数制限に達したら停止する。 |
| エンドポイントテレメトリ | キューの深さだけでは、処理速度の遅いパートナーや、接続に失敗するパートナーを見逃してしまう。 | エンドポイントごとに、成功率、最も古い保留中のイベント、およびレイテンシのパーセンタイルを追跡します。 |
セキュリティと重複制御は、再試行チューニングよりも優先されます。
送信本文は機密性の高い運用データとして扱ってください。リクエスト本文に署名し、配信タイムスタンプと変更不可能なイベントIDを含め、署名シークレットをローテーションし、受信アプリケーションが同じイベントの再送信を安全に無視するようにしてください。HTTPレスポンスが成功したからといって、ビジネスイベントが正確に一度だけ適用されたという証拠にはなりません。受信側でその判断を下す必要があります。
アフィリエイトまたはiGamingのポストバックワークフローでは、アクセス制御と保持ルールが明示的に定められていない限り、クリックID、コンバージョンID、および支払い関連フィールドをログに記録しないようにしてください。下記のScaleoの既存のリファレンスは、アフィリエイトのポストバック遅延の例としてのみ有効であり、自社サービス間の配信契約を文書化することの代わりとなるものではありません。
役立つ実装例: Stripeウェブフック署名検証, PostgreSQLのSELECTとSKIPがロックされています, AWSの指数バックオフとジッターに関するガイダンス.
データモデル
すべては2つのテーブルから始まります。1つは配信キュー用、もう1つはレイテンシ測定用です。
SQL
-- Pending and in-flight webhook deliveries
CREATE TABLE webhook_queue (
id SERIAL PRIMARY KEY,
endpoint_id INTEGER NOT NULL,
endpoint_url TEXT NOT NULL,
payload JSONB NOT NULL,
-- Delivery state
status VARCHAR(20) NOT NULL DEFAULT 'pending',
-- pending, in_flight, delivered, failed, dead_letter
attempt_count INTEGER NOT NULL DEFAULT 0,
max_attempts INTEGER NOT NULL DEFAULT 5,
-- Timing
created_at TIMESTAMP NOT NULL DEFAULT NOW(),
next_attempt_at TIMESTAMP NOT NULL DEFAULT NOW(),
delivered_at TIMESTAMP,
last_error TEXT,
-- Response tracking
last_status_code INTEGER,
last_response_ms INTEGER, -- response time in milliseconds
INDEX idx_queue_next (status, next_attempt_at)
WHERE status IN ('pending', 'failed')
);
-- Per-endpoint latency measurements (ring buffer)
CREATE TABLE endpoint_latency (
id SERIAL PRIMARY KEY,
endpoint_id INTEGER NOT NULL,
response_ms INTEGER NOT NULL,
status_code INTEGER,
measured_at TIMESTAMP NOT NULL DEFAULT NOW(),
INDEX idx_latency_endpoint (endpoint_id, measured_at)
);
-- Endpoint health state (circuit breaker)
CREATE TABLE endpoint_health (
endpoint_id INTEGER PRIMARY KEY,
endpoint_url TEXT NOT NULL,
state VARCHAR(20) NOT NULL DEFAULT 'closed',
-- closed (healthy), open (broken), half_open (testing)
consecutive_failures INTEGER NOT NULL DEFAULT 0,
failure_threshold INTEGER NOT NULL DEFAULT 10,
last_failure_at TIMESTAMP,
last_success_at TIMESTAMP,
opened_at TIMESTAMP, -- when circuit opened
cooldown_seconds INTEGER NOT NULL DEFAULT 300, -- 5 min before half_open
-- Latency stats (updated periodically)
p50_ms INTEGER,
p95_ms INTEGER,
p99_ms INTEGER,
sample_count INTEGER NOT NULL DEFAULT 0
);
この図式について注意すべき点が3つあります。
まず、 webhook_queue テーブルは next_attempt_at 個別のスケジューリングメカニズムの代わりに列を使用します。ワーカーは、次の行をポーリングします。 status IN ('pending', 'failed') AND next_attempt_at <= NOW()これは簡易的な遅延キューであり、1分間に約1万件の配信までは問題なく動作します。それ以上の配信量の場合は、適切なメッセージブローカーに切り替えてください。
第二に、 endpoint_latency テーブルはリングバッファとして機能します。24 時間以上前の行は定期的に削除します。レイテンシのパーセンタイルは endpoint_health これらはこの移動平均線から計算されるものであり、過去の平均値ではなく、最近の動向を表しています。
第三に、 endpoint_health このテーブルは、回路遮断器の状態遷移図を実装しています。詳細については後述します。
配達員
コアとなるワーカーループは意図的にシンプルに設計されている。複雑な処理はリトライロジックとサーキットブレーカーに任せるべきであり、配信パス自体には任せるべきではない。
パイソン
import requests
import time
import psycopg2
from psycopg2.extras import RealDictCursor
from datetime import datetime, timedelta
DB_DSN = "postgresql://user:pass@localhost/webhooks"
def get_connection():
return psycopg2.connect(DB_DSN)
def deliver_webhooks(batch_size=50):
"""
Fetch pending webhooks and attempt delivery.
Uses SELECT FOR UPDATE SKIP LOCKED for safe concurrent workers.
"""
conn = get_connection()
cur = conn.cursor(cursor_factory=RealDictCursor)
try:
cur.execute("""
SELECT id, endpoint_id, endpoint_url, payload,
attempt_count, max_attempts
FROM webhook_queue
WHERE status IN ('pending', 'failed')
AND next_attempt_at <= NOW()
ORDER BY next_attempt_at ASC
LIMIT %s
FOR UPDATE SKIP LOCKED
""", (batch_size,))
rows = cur.fetchall()
for row in rows:
# Check circuit breaker before attempting
if is_circuit_open(cur, row['endpoint_id']):
# Don't attempt delivery — reschedule for after cooldown
reschedule_for_cooldown(cur, row['id'], row['endpoint_id'])
continue
# Attempt delivery and measure latency
result = attempt_delivery(
row['endpoint_url'],
row['payload']
)
# Record latency measurement regardless of success/failure
record_latency(
cur,
row['endpoint_id'],
result['response_ms'],
result['status_code']
)
if result['success']:
mark_delivered(cur, row['id'], result)
record_success(cur, row['endpoint_id'])
else:
handle_failure(
cur, row['id'], row['endpoint_id'],
row['attempt_count'], row['max_attempts'],
result
)
conn.commit()
except Exception as e:
conn.rollback()
raise
finally:
cur.close()
conn.close()
def attempt_delivery(url, payload):
"""
Fire the webhook and measure response time.
Returns dict with success, status_code, response_ms, error.
"""
start = time.monotonic()
try:
response = requests.post(
url,
json=payload,
timeout=15, # 15 second hard timeout
headers={
'Content-Type': 'application/json',
'User-Agent': 'WebhookDelivery/1.0',
'X-Delivery-Timestamp': str(int(time.time()))
}
)
elapsed_ms = int((time.monotonic() - start) * 1000)
return {
'success': 200 <= response.status_code < 300,
'status_code': response.status_code,
'response_ms': elapsed_ms,
'error': None if response.ok else f"HTTP {response.status_code}"
}
except requests.Timeout:
elapsed_ms = int((time.monotonic() - start) * 1000)
return {
'success': False,
'status_code': None,
'response_ms': elapsed_ms,
'error': 'timeout_15s'
}
except requests.ConnectionError as e:
elapsed_ms = int((time.monotonic() - start) * 1000)
return {
'success': False,
'status_code': None,
'response_ms': elapsed_ms,
'error': f'connection_error: {str(e)[:200]}'
}
その SELECT FOR UPDATE SKIP LOCKED 句は、複数のワーカー インスタンスを実行するために重要です。 SKIP LOCKED2つのワーカーが同じ行でブロックしてしまう可能性があります。これを使用すると、各ワーカーが異なるバッチの保留中のWebhookを取得します。これにより、ワーカープロセスを増やすだけで水平スケーリングが可能になります。
その time.monotonic() 代わりに電話する time.time() 意図的なものです。 time.time() NTP調整中に前の状態に戻ることがあります。 time.monotonic() 決して後退しない。これは、1秒未満の遅延を測定する際には重要な点である。
指数バックオフとジッターを使用した再試行ロジック
配信が失敗した場合、再試行のタイミングによって、システムが適切に復旧するか、あるいは処理に苦戦しているエンドポイントに大量の負荷をかけるかが決まります。
パイソン
import random
def calculate_next_attempt(attempt_count, base_delay=30, max_delay=3600):
"""
Exponential backoff with full jitter.
Attempt 1: 0-30s
Attempt 2: 0-60s
Attempt 3: 0-120s
Attempt 4: 0-240s
Attempt 5: 0-480s (capped at max_delay)
Full jitter prevents thundering herd when an endpoint
recovers and hundreds of retries fire simultaneously.
"""
exponential_delay = base_delay * (2 ** attempt_count)
capped_delay = min(exponential_delay, max_delay)
jittered_delay = random.uniform(0, capped_delay)
return datetime.utcnow() + timedelta(seconds=jittered_delay)
def handle_failure(cur, webhook_id, endpoint_id,
attempt_count, max_attempts, result):
"""
Handle a failed delivery attempt.
Either retry with backoff or move to dead letter queue.
"""
new_attempt_count = attempt_count + 1
if new_attempt_count >= max_attempts:
# Exhausted retries — dead letter
cur.execute("""
UPDATE webhook_queue
SET status = 'dead_letter',
attempt_count = %s,
last_error = %s,
last_status_code = %s,
last_response_ms = %s
WHERE id = %s
""", (
new_attempt_count, result['error'],
result['status_code'], result['response_ms'],
webhook_id
))
else:
# Schedule retry with backoff
next_attempt = calculate_next_attempt(new_attempt_count)
cur.execute("""
UPDATE webhook_queue
SET status = 'failed',
attempt_count = %s,
next_attempt_at = %s,
last_error = %s,
last_status_code = %s,
last_response_ms = %s
WHERE id = %s
""", (
new_attempt_count, next_attempt,
result['error'], result['status_code'],
result['response_ms'], webhook_id
))
# Update circuit breaker
record_failure(cur, endpoint_id)
非相関ジッターや等ジッターではなく、フルジッターを使用する理由は?AWSがこの点について決定的な分析を公開しています。フルジッター(0から指数上限までの範囲でランダム化)は、すべてのクライアントで合計完了時間が最も短くなります。等ジッター(上限の半分から上限までの範囲でランダム化)はより保守的ですが、再試行バックログの処理に時間がかかります。多数の独立したエンドポイントがあるWebhook配信の場合、各エンドポイントの再試行は独立しているため、エンドポイント間で調整する必要がないことから、フルジッターが適切な選択肢となります。
回路遮断器:壊れたエンドポイントを叩くのをやめましょう
サーキットブレーカーパターンは、システムが継続的に障害を起こすエンドポイントにリソースを浪費するのを防ぎます。これがないと、障害が発生したエンドポイントには数百件もの再試行が蓄積され、それぞれ15秒でタイムアウトしてしまうため、決して成功しない配信にワーカーのリソースが浪費されてしまいます。
パイソン
def is_circuit_open(cur, endpoint_id):
"""
Check if the circuit breaker is open (endpoint is broken).
If open and cooldown has passed, transition to half_open.
"""
cur.execute("""
SELECT state, opened_at, cooldown_seconds
FROM endpoint_health
WHERE endpoint_id = %s
""", (endpoint_id,))
row = cur.fetchone()
if not row:
return False # no health record = assume healthy
if row['state'] == 'closed':
return False
if row['state'] == 'open':
# Check if cooldown period has elapsed
if row['opened_at'] and row['cooldown_seconds']:
elapsed = (datetime.utcnow() - row['opened_at']).total_seconds()
if elapsed >= row['cooldown_seconds']:
# Transition to half_open — allow one probe
cur.execute("""
UPDATE endpoint_health
SET state = 'half_open'
WHERE endpoint_id = %s
""", (endpoint_id,))
return False # allow the probe delivery
return True # still in cooldown
if row['state'] == 'half_open':
return False # allow probe delivery
return False
def record_failure(cur, endpoint_id):
"""
Record a delivery failure. Open circuit if threshold reached.
"""
cur.execute("""
UPDATE endpoint_health
SET consecutive_failures = consecutive_failures + 1,
last_failure_at = NOW()
WHERE endpoint_id = %s
RETURNING consecutive_failures, failure_threshold, state
""", (endpoint_id,))
row = cur.fetchone()
if not row:
# Create health record on first failure
cur.execute("""
INSERT INTO endpoint_health (endpoint_id, endpoint_url,
consecutive_failures, last_failure_at)
VALUES (%s, '', 1, NOW())
ON CONFLICT (endpoint_id) DO UPDATE
SET consecutive_failures = endpoint_health.consecutive_failures + 1,
last_failure_at = NOW()
""", (endpoint_id,))
return
if row['state'] == 'half_open':
# Probe failed — re-open circuit with longer cooldown
cur.execute("""
UPDATE endpoint_health
SET state = 'open',
opened_at = NOW(),
cooldown_seconds = LEAST(cooldown_seconds * 2, 3600)
WHERE endpoint_id = %s
""", (endpoint_id,))
elif (row['state'] == 'closed' and
row['consecutive_failures'] >= row['failure_threshold']):
# Threshold reached — open circuit
cur.execute("""
UPDATE endpoint_health
SET state = 'open',
opened_at = NOW(),
cooldown_seconds = 300 -- reset to 5 minutes
WHERE endpoint_id = %s
""", (endpoint_id,))
def record_success(cur, endpoint_id):
"""
Record a delivery success. Close circuit if half_open.
"""
cur.execute("""
UPDATE endpoint_health
SET consecutive_failures = 0,
last_success_at = NOW(),
state = 'closed'
WHERE endpoint_id = %s
""", (endpoint_id,))
半開状態から開状態へのエスカレーションと、クールダウン時間の倍増は、多くの実装で見落とされている重要な点です。プローブ中にエンドポイントが失敗した場合(半開状態)、さらに5分後に再試行するのは適切ではありません。エンドポイントは依然として機能していない状態です。クールダウン時間を10分、次に20分と倍増させ、最大1時間まで延長します。これにより、サーキットブレーカーが周期的にハンマーで叩きつけるような動作をするのを防ぎます。
レイテンシーパーセンタイル追跡
平均値は当てにならない。平均応答時間が200msのエンドポイントでも、95%の確率で50msで応答し、残りの5%では3,000msかかる可能性がある。平均値は問題ないように見えるかもしれないが、P95(95%信頼区間)を見ると、20回に1回の割合で問題が発生することが明らかになる。
パイソン
def record_latency(cur, endpoint_id, response_ms, status_code):
"""
Record a latency measurement and update percentile stats.
"""
cur.execute("""
INSERT INTO endpoint_latency
(endpoint_id, response_ms, status_code)
VALUES (%s, %s, %s)
""", (endpoint_id, response_ms, status_code))
def update_latency_percentiles(cur, endpoint_id, window_hours=24):
"""
Calculate P50, P95, P99 from the rolling window.
Uses PostgreSQL's percentile_cont for exact percentiles.
"""
cur.execute("""
SELECT
COUNT(*) as sample_count,
percentile_cont(0.50) WITHIN GROUP
(ORDER BY response_ms) AS p50,
percentile_cont(0.95) WITHIN GROUP
(ORDER BY response_ms) AS p95,
percentile_cont(0.99) WITHIN GROUP
(ORDER BY response_ms) AS p99
FROM endpoint_latency
WHERE endpoint_id = %s
AND measured_at >= NOW() - INTERVAL '%s hours'
""", (endpoint_id, window_hours))
row = cur.fetchone()
if row and row['sample_count'] > 0:
cur.execute("""
UPDATE endpoint_health
SET p50_ms = %s,
p95_ms = %s,
p99_ms = %s,
sample_count = %s
WHERE endpoint_id = %s
""", (
int(row['p50']), int(row['p95']),
int(row['p99']), row['sample_count'],
endpoint_id
))
return row
def get_slow_endpoints(cur, p95_threshold_ms=2000):
"""
Find endpoints whose P95 latency exceeds the threshold.
These are candidates for investigation or circuit opening.
"""
cur.execute("""
SELECT endpoint_id, endpoint_url,
p50_ms, p95_ms, p99_ms, sample_count,
state, consecutive_failures
FROM endpoint_health
WHERE p95_ms > %s
AND sample_count >= 20 -- need sufficient samples
ORDER BY p95_ms DESC
""", (p95_threshold_ms,))
return cur.fetchall()
PostgreSQLの percentile_cont は、正確なパーセンタイルを計算する順序付き集合集計関数です。大規模なデータセットの場合は、 に切り替えます。 percentile_disc (補間ではなく実際の観測値を返す)またはt-digest近似を使用します。24時間ウィンドウのWebhook配信の場合、エンドポイントあたり約100,000万回の測定値までは、生データに対する正確なパーセンタイルで十分高速です。
その get_slow_endpoints この関数は、15分ごとにスケジュールされたチェックとして実行されます。P95が2秒を超えるエンドポイントは調査対象としてフラグ付けされます。P95が5秒を超えるエンドポイントは、サーキットブレーカーのしきい値が引き下げられます。これは、配信失敗ごとにワーカー スレッドがタイムアウト期間全体にわたって占有されるため、サーキットが開くまでに許容される連続失敗回数が少なくなるためです。
Webhook配信状況の監視
以下は、私が5分ごとに実行している監視クエリです。このクエリは、配信パイプライン全体の健全性に関する概要を1行で出力します。
SQL
SELECT
-- Queue depth
COUNT(*) FILTER (WHERE status = 'pending') AS pending,
COUNT(*) FILTER (WHERE status = 'failed') AS awaiting_retry,
COUNT(*) FILTER (WHERE status = 'in_flight') AS in_flight,
COUNT(*) FILTER (WHERE status = 'dead_letter') AS dead_letter,
-- Delivery rate (last hour)
COUNT(*) FILTER (
WHERE status = 'delivered'
AND delivered_at >= NOW() - INTERVAL '1 hour'
) AS delivered_last_hour,
-- Failure rate (last hour)
COUNT(*) FILTER (
WHERE status IN ('failed', 'dead_letter')
AND created_at >= NOW() - INTERVAL '1 hour'
) AS failed_last_hour,
-- Oldest undelivered
MIN(created_at) FILTER (
WHERE status IN ('pending', 'failed')
) AS oldest_pending,
-- Average delivery latency (last hour, successful only)
AVG(last_response_ms) FILTER (
WHERE status = 'delivered'
AND delivered_at >= NOW() - INTERVAL '1 hour'
) AS avg_delivery_ms_last_hour
FROM webhook_queue;
その oldest_pending このクエリにおいて最も重要な指標は値です。値が最大再試行期間(すべてのバックオフ遅延の合計)よりも古い場合、何らかの構造的な問題があります。ワーカーが停止しているか、エンドポイントがブラックホール化されているか、キューが処理速度よりも速く増大しているかのいずれかです。
3 つの条件でアラートを発します: デッドレター数の増加 (エンドポイントが恒久的に故障しており、誰も調査していない)、保留キューの経過時間が 30 分を超えている (配信が遅れている)、 エンドポイントごとのP95レイテンシが、配信信頼性の低下を示す閾値を超えている。3つ目は早期警告信号です。障害が発生する前にレイテンシが増加します。200ミリ秒で応答していたエンドポイントが3秒かかるようになった場合、タイムアウトが発生し始める兆候です。
デッドレターキューは単なる保管場所ではない
ほとんどのチームは、失敗したWebhookが送られるテーブルとしてデッドレターキューを実装しています。そして、インシデント対応時に時折それをチェックします。これは無駄です。
デッドレターキューは、最も貴重なデバッグ用データセットです。各行は、システムが複数回試行したものの、最終的に配信を断念したメールを表しています。デッドレターのパターンからは、成功指標では決して得られない情報が得られます。
パイソン
def analyze_dead_letters(cur, hours=24):
"""
Analyze recent dead letter entries for patterns.
Returns per-endpoint failure analysis.
"""
cur.execute("""
SELECT
endpoint_id,
endpoint_url,
COUNT(*) AS dead_count,
-- Most common error
MODE() WITHIN GROUP (ORDER BY last_error) AS primary_error,
-- Most common status code
MODE() WITHIN GROUP (ORDER BY last_status_code)
AS primary_status_code,
-- Timing
MIN(created_at) AS first_dead,
MAX(created_at) AS last_dead,
-- Average attempts before giving up
AVG(attempt_count)::INTEGER AS avg_attempts
FROM webhook_queue
WHERE status = 'dead_letter'
AND created_at >= NOW() - INTERVAL '%s hours'
GROUP BY endpoint_id, endpoint_url
ORDER BY dead_count DESC
LIMIT 20
""", (hours,))
return cur.fetchall()
宛先不明の郵便物を確認する際、私は3つのパターンを探しています。
クラスター障害: 同じエンドポイントに対して同じ時間に50件のデッドレターが発生した場合、エンドポイントがダウンし、再試行ウィンドウ内に復旧しなかったことを意味します。対処方法:再試行ウィンドウを延長するか、手動でキューを再作成してください。
ステータスコードのパターン: 401/403エラー(デッドレター)が急増している場合は、エンドポイントの認証情報がローテーションされたにもかかわらず、誰もWebhookの設定を更新していないことを意味します。429エラー(リクエスト過多)が急増している場合は、レート制限を超えているため、スロットリングを行う必要があります。
徐々に蓄積していく: 単一のエンドポイントに対して、1日に2~3通のデッドレターが均等に分散して発生している。これは最も厄介なパターンで、エンドポイントは大部分は正常に動作しているものの、断続的な障害が発生し、時間の経過とともに再試行回数を使い果たしてしまう。修正は通常、増加傾向にある。 max_attempts その特定のエンドポイントに対して、またはタイムアウト時間を短縮する。
ワーカーの実行
すべてを繋ぎ合わせる主要なループ:
パイソン
import signal
import sys
running = True
def shutdown_handler(signum, frame):
global running
running = False
print(f"Received signal {signum}, shutting down gracefully...")
signal.signal(signal.SIGTERM, shutdown_handler)
signal.signal(signal.SIGINT, shutdown_handler)
def main():
print("Webhook delivery worker starting...")
while running:
try:
deliver_webhooks(batch_size=50)
except Exception as e:
print(f"Worker error: {e}")
time.sleep(5) # back off on errors
continue
# Update latency stats every 100 iterations
# (cheap operation, doesn't need to run every loop)
if int(time.time()) % 100 == 0:
conn = get_connection()
cur = conn.cursor(cursor_factory=RealDictCursor)
try:
cur.execute(
"SELECT DISTINCT endpoint_id FROM endpoint_health"
)
for row in cur.fetchall():
update_latency_percentiles(cur, row['endpoint_id'])
conn.commit()
finally:
cur.close()
conn.close()
# Poll interval — 500ms keeps latency low without
# hammering the database
time.sleep(0.5)
print("Worker shut down cleanly.")
if __name__ == '__main__':
main()
その SIGTERM ハンドラは、コンテナ化された環境でのクリーンシャットダウンに不可欠です。Kubernetes が SIGTERM を送信すると、ワーカーは現在のバッチを完了し、トランザクションをコミットして終了します。これがないと、行がスタックしてしまいます。 in_flight 処理するワーカーがいない状態です。
このシステムではできないこと(そして、より多くの機能が必要な場合)
この実装では、単一のPostgreSQLインスタンスと2~3個のワーカープロセスで、1分あたり約1万件の配送を処理できます。それ以上の処理能力が必要な場合は、3つの変更が必要です。
まず、PostgreSQLキューをRedis StreamsまたはRabbitMQに置き換えます。 SELECT FOR UPDATE SKIP LOCKED このパターンは、高スループット時にキューテーブルへの書き込み競合を引き起こします。専用のメッセージブローカーを使用することで、この問題を解消できます。
次に、エンドポイントごとのレート制限を追加します。受信エンドポイントの中には、レート制限(1分あたり100リクエスト、1時間あたり1,000リクエストなど)が設定されているものがあります。クライアント側のレート制限がないと、その制限を超えてしまい、429エラーが発生します。エンドポイントごとにトークンバケットを実装してください。
第三に、リクエスト署名を追加します。ペイロードにHMAC-SHA256署名を追加することで、受信側のエンドポイントは、Webhookが送信元のシステムから送信されたものであり、送信中に改ざんされていないことを検証できます。これは、金融データを送信するWebhookシステムにとって必須の要件です。
この記事で紹介するシステムは、その基盤となるものです。規模に関わらず、あらゆるWebhook配信システムに必要な、再試行ロジック、サーキットブレーカー、レイテンシー測定、デッドレター分析といった難題を処理します。追加する具体的なコンポーネント(メッセージブローカー、レートリミッター、リクエスト署名など)は、スループットとセキュリティ要件によって異なります。
最も重要な部分は、多くのチームが見落としている部分です。それは、配信システム自体を測定することです。「過去24時間におけるエンドポイントXへのP95配信遅延時間はどれくらいか」という質問に答えられないのであれば、手探りで運用しているようなものです。まず計測システムを構築しましょう。そうすれば、他のすべては自然とついてきます。
Webhook配信に関するよくある質問
ウェブフック送信者は、正確に1回限りの配信を約束すべきでしょうか?
通常はそうではありません。送信側は再試行を可視化し、安定したイベントIDを提供する必要があります。受信側は処理を冪等にすることで、イベントが複数回配信されてもビジネス上の結果が重複しないようにする必要があります。
どの失敗を再試行すべきか?
契約で一時的なエラーとして分類されているもの(ネットワークエラー、タイムアウト、特定のサーバー応答など)のみを再試行してください。不正なリクエスト、認証失敗、その他の永続的なエラーについては、明確な修復手順がない限り、繰り返し再試行しないでください。
チームはいつ、データベースを基盤としたキューから脱却すべきでしょうか?
計測された競合状況、バックログの経過時間、スループット、または運用復旧ニーズから、データベースキューが配信契約を満たさなくなったと判断された場合に、移行を実施してください。容量に関する決定は、一般的なリクエストレートの主張ではなく、観測されたワークロードに基づいて行うべきです。