בניית מערכת משלוח Webhook: Postbacks אמינים בשנת 2026

מדריך מעשי לשנת 2026 לאספקת Webhook ו-Affiliate-Postback עמידה עם תורים, ניסיונות חוזרים מרצדים, מפסקי זרם, אידמפוטנטיות, חתימות וטלמטריה של השהייה לכל נקודת קצה.
מערכת משלוח Webhook - בניית מערכת משלוח Webhook: Postbacks אמינים בשנת 2026

עודכן לאחרונה ב -24 ביוני 2026 על ידי קיסר פיקסון

תשובה ישירה: שירות משלוח אמין של webhook או שירות משלוח של שותפים-פוסטבקס זקוק לתור עמיד, ניסיונות חוזרים מוגבלים עם ריצוד, שבירת מעגלים לכל נקודת קצה, הגנה מפני כפילויות, אימות חתימות וטלמטריה שמראה אם ​​המסירה מאטה לפני שהיא נכשלת.

יישום זה שומר על מודל הביצוע קטן במכוון: פונקציות Python ו-PostgreSQL עבור תור, היסטוריית מסירה ובריאות נקודות קצה. זוהי נקודת התחלה מעשית עבור B2B SaaS, פעולות שותפים של iGaming וכל מוצר שחייב לספק אירועי יציאה מבלי להתייחס לפסק זמן HTTP יחיד כהמרה אבודה.

מה חייב להיות מובטח בחוזה האספקה

שליטה למה זה משנה בדיקה תפעולית
רישום אירועים עמיד אירועים שורדים את ההפעלות מחדש של העובדים. לכל משלוח שאושר יש מזהה וסטטוס יציבים.
בקשה חתומה מקבלי המכשיר יכולים לאמת את השולח ולזהות שיבוש גוף. השתמש בחתימה עם חותמת זמן ודחה בקשות ישנות.
הגנה מפני כפילויות ניסיונות חוזרים עלולים ליצור המרות או עדכונים כפולים. שלח מזהה אירוע והפוך את המקבל לאידמפוטנטי.
ניסיונות חוזרים מוגבלים כשלים חולפים מתאוששים מבלי להעמיס על נקודת קצה פגומה. יש לסגת עם רעידות; לעצור לאחר מגבלת ניסיונות מתועדת.
טלמטריה של נקודות הקצה עומק התור לבדו מסתיר שותפים איטיים או כושלים. מעקב אחר שיעור הצלחה, האירוע הממתין הישן ביותר ואחוזוני השהייה לכל נקודת קצה.

בקרות אבטחה וכפילויות מגיעות לפני כוונון חוזר

התייחסו לגוף הבקשה היוצא כאל מידע תפעולי רגיש. יש לחתום על גוף הבקשה המדויק, לכלול חותמת זמן של מסירה ומזהה אירוע בלתי משתנה, לסובב סודות חתימה, ולוודא שהאפליקציה המקבלת מתעלמת בבטחה משחזור של אותו אירוע. תגובת HTTP מוצלחת אינה הוכחה לכך שאירוע עסקי הוחל פעם אחת בדיוק; הצד המקבל חייב לקבל החלטה זו.

עבור תהליך עבודה של החזרת חשבונות (postback) של שותפים או גיימינג דיגיטלי (iGaming), יש להרחיק מזהי קליקים, מזהי המרות ושדות רלוונטיים לתשלום מהיומנים, אלא אם כן בקרות הגישה וכללי השמירה מפורשים. ההפניה הקיימת של Scaleo להלן רלוונטית רק כדוגמה להשהיית החזרת חשבונות (postback) של שותפים; היא אינה תחליף לתיעוד חוזה האספקה ​​בין השירותים שלכם.

מקורות שימושיים ליישום: אימות חתימת webhook של Stripe, SELECT ו-SKIP של PostgreSQL נעולים, ו הנחיות AWS בנושא ניתוק אקספוננציאלי וריצוד.

מודל הנתונים

הכל מתחיל עם שתי טבלאות: אחת לתור המסירה ואחת למדידות ההשהיה.

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
);

שלושה דברים שכדאי לציין בנוגע לסכימה הזו.

ראשית, webhook_queue הטבלה משתמשת ב- next_attempt_at עמודה במקום מנגנון תזמון נפרד. העובד בוחן שורות שבהן status IN ('pending', 'failed') AND next_attempt_at <= NOW()זהו תור עיכובים של אדם עני והוא עובד מצוין עד כ-10,000 מסירות בדקה. מעבר לכך, החליפו למתווך הודעות תקין.

שנית, 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 LOCKED, שני עובדים ייחסמו באותה שורה. בעזרתו, כל עובד תופס קבוצה שונה של webhooks ממתינים. זה נותן לך קנה מידה אופקי פשוט על ידי הפעלת תהליכי עובד נוספים.

השמיים time.monotonic() להתקשר במקום time.time() הוא מכוון. time.time() יכול לקפוץ אחורה במהלך התאמות NTP. time.monotonic() אף פעם לא חוזר אחורה, וזה חשוב כשמודדים השהייה של פחות משנייה.

נסה שוב לוגיקה עם גיבוי אקספוננציאלי וריצוד

כאשר מסירה נכשלת, תזמון הניסיון החוזר קובע האם המערכת שלך מתאוששת בצורה חלקה או יוצרת עדר רועם שפוגע בנקודת קצה מתקשה.

פִּיתוֹן

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)

למה להשתמש ב-full jitter במקום dearcorrelated jitter או equal jitter? AWS פרסמה את הניתוח הסופי בנושא. full jitter (אקראי בין 0 לתקרת הערך האקספוננציאלית) מייצרת את זמן ההשלמה הכולל הנמוך ביותר בכל הלקוחות. equal jitter (אקראי בין חצי מהתקרת הערך לתקרת הערך המלאה) הוא שמרני יותר אך איטי יותר לניקוז צבר הניסיון החוזר. עבור אספקת webhook שבה יש לך נקודות קצה רבות ועצמאיות, full jitter הוא הבחירה הנכונה מכיוון שהניסיונות החוזרים של כל נקודת קצה הם בלתי תלויים - אינך מתאם ביניהם.

מפסק מעגל: תפסיקו להכות נקודות קצה שבורות

תבנית מפסק המעגל מונעת מהמערכת שלך לבזבז משאבים על נקודות קצה שנכשלות באופן עקבי. בלעדיה, נקודת קצה מתה צוברת מאות ניסיונות חוזרים ממתינים שכולם מגיעים ל-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,))

הסלמה של half_open → open עם זמן קירור כפול היא הפרט שרוב המימושים מפספסים. אם נקודת קצה נכשלת במהלך הבדיקה (מצב half_open), לא כדאי לנסות שוב בעוד 5 דקות. נקודת הקצה עדיין שבורה. הכפילו את זמן הקירור ל-10 דקות, לאחר מכן 20, עם הגבלה של שעה. זה מונע מהמפסק להפוך למנגנון פטיש מחזורי.

מעקב אחר אחוזוני השהייה

הממוצעים שקרים. נקודת קצה עם זמן תגובה ממוצע של 200ms עשויה להגיב תוך 50ms ב-95% מהמקרים וב-3,000ms ב-5% הנותרים. הממוצע נראה תקין. ה-P95 מגלה בעיה שמשפיעה על 1 מכל 20 מסירות.

פִּיתוֹן

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. עבור מסירת webhook עם חלון של 24 שעות, אחוזונים מדויקים על הנתונים הגולמיים מהירים מספיק עד כ-100,000 מדידות לכל נקודת קצה.

השמיים get_slow_endpoints הפונקציה הזו היא מה שאני מפעיל כבדיקה מתוזמנת כל 15 דקות. נקודות קצה עם P95 מעל 2 שניות מסומנות לבדיקה. נקודות קצה עם P95 מעל 5 שניות מקבלות סף מפסק המעגל שלהן מופחת - הן מורשות פחות כשלים רצופים לפני פתיחת המעגל, מכיוון שכל מסירה כושלת קושרת worker thread למשך הזמן הקצוב המלא.

ניטור תקינות מסירת Webhook

הנה שאילתת הניטור שאני מפעיל כל חמש דקות. היא מייצרת סיכום תקינות של כל צינור האספקה:

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 ערך הוא המדד החשוב ביותר בשאילתה זו. אם הוא ישן יותר מחלון הניסיון החוזר המרבי שלך (סכום כל עיכובי ה-backoff), משהו לא בסדר מבחינה מבנית - או שהעובד תקוע, נקודת הקצה נמצאת בחור שחור, או שהתור גדל מהר יותר ממה שאתה יכול לרוקן אותו.

אני מתריע בשלושה תנאים: מספר האותיות המתות עולה (נקודות הקצה כושלות לצמיתות ואף אחד לא חוקר), גיל התור הממתין עולה על 30 דקות (המשלוח מפגר), ו השהיית P95 לכל נקודת קצה חורגת מספי מעקב המצביעים על ירידה באמינות המסירה השלישי הוא אות האזהרה המוקדם - ההשהיה עולה לפני שתקלות קורות. נקודת קצה שהגיבה תוך 200 מילישניות ומתחילה להגיב תוך 3 שניות עומדת להתחיל פסק זמן.

תור האותיות המתות אינו רק אחסון

רוב הצוותים מיישמים תור אותיות מוחלטות כטבלה שאליה עוברים webhooks שנכשלו כדי למות. הם בודקים זאת מדי פעם במהלך תגובה לאירועים. זה בזבוז.

תור האותיות המתות הוא מערך הנתונים החשוב ביותר שלך לניפוי שגיאות. כל שורה מייצגת מסירה שהמערכת שלך ניסתה מספר פעמים וויתרה עליה. דפוס האותיות המתות אומר לך דברים שמדד ההצלחה לעולם לא יאמר.

פִּיתוֹן

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()

כשאני סוקר אותיות מתות, אני מחפש שלוש דפוסים.

כשלים באשכול: 50 אותיות חסרות עבור אותה נקודת קצה באותה שעה פירושן שנקודת הקצה נפלה ולא התאוששה בתוך חלון הניסיונות החוזרים. פעולה: הארכת חלון הניסיונות החוזרים או יישום שינוי תור ידני.

דפוסי קוד סטטוס: עלייה חדה במספר אותיות ריקות (401/403) פירושה שנקודת הקצה סיבבה את פרטי ההרשאה ואף אחד לא עדכן את תצורת ה-webhook. עלייה חדה ב-429 (Too Many Requests) פירושה שאתם חורגים ממגבלת הקצב שלהם וצריכים לווסת.

הצטברות הדרגתית: 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, ה-worker מסיים את האצווה הנוכחית שלו, מבצע את הטרנזקציה ויוצא. בלעדיו, שורות נתקעות ב- in_flight סטטוס ללא עובד שמעבד אותם.

מה שהמערכת הזו לא עושה (וכאשר אתם צריכים יותר)

יישום זה מטפל בעד כ-10,000 מסירות לדקה על מופע PostgreSQL יחיד עם 2-3 תהליכי עבודה. מעבר לכך, נדרשים שלושה שינויים.

ראשית, החלף את תור PostgreSQL ב-Redis Streams או RabbitMQ. SELECT FOR UPDATE SKIP LOCKED תבנית זו יוצרת מחלוקת כתיבה בטבלת התור בתפוקה גבוהה. מתווך הודעות ייעודי מבטל זאת.

שנית, הוסיפו הגבלת קצב לכל נקודת קצה. לחלק מנקודות הקצה המקבלות יש מגבלות קצב (100 בקשות לדקה, 1,000 לשעה). ללא הגבלת קצב בצד הלקוח, תבזבזו את המכסה שלהם ותקבלו 429 בקשות. הטמיעו דלי אסימונים לכל נקודת קצה.

שלישית, הוסף חתימת בקשה. חתימות HMAC-SHA256 על המטען מאפשרות לנקודת הקצה המקבלת לאמת שה-webhook הגיע מהמערכת שלך ולא טופלה בו במהלך ההעברה. זהו טבלת סטייקים עבור כל מערכת webhook ששולחת נתונים פיננסיים.

המערכת במאמר זה היא הבסיס. היא מטפלת בבעיות הקשות - לוגיקת ניסיונות חוזרים, שבירת מעגלים, מדידת השהייה, ניתוח אותיות חסומות - שכל מערכת אספקת webhook צריכה ללא קשר לקנה המידה. הרכיבים הספציפיים שתוסיפו (מתווך הודעות, מגביל קצב, חתימת בקשות) תלויים בדרישות התפוקה והאבטחה שלכם.

החלק החשוב ביותר הוא החלק שרוב הצוותים מדלגים עליו: מדידת מערכת האספקה ​​עצמה. אם אינכם יכולים לענות על "מהי זמן השהיית האספקה ​​של P95 לנקודת הקצה X ב-24 השעות האחרונות", אתם פועלים בעיוורון. בנו קודם את המכשור. כל השאר יבוא לאחר מכן.

שאלות נפוצות לגבי משלוח Webhook

האם שולח ב-webhook צריך להבטיח מסירה חד פעמית בדיוק?

בדרך כלל לא. על השולח להפוך ניסיונות חוזרים לגלויים ולספק מזהה אירוע יציב; על המקבל להפוך את העיבוד לאידמפוטנטי כך שניתן יהיה לספק אירוע יותר מפעם אחת מבלי לשכפל את תוצאת העסק.

אילו כשלים יש לנסות שוב?

נסה שוב רק כשלים שהחוזה שלך מסווג כחולפים, כגון שגיאות רשת, פסקי זמן ותגובות שרת נבחרות. אין לנסות שוב ושוב בקשות בעלות מבנה שגוי, כשלי אימות או שגיאות קבועות אחרות ללא נתיב תיקון מפורש.

מתי על צוות לעבור מעבר לתור מגובה על ידי מסד נתונים?

העברה כאשר מדדי מתח, גיל צבר הזמנות, תפוקה או צורכי התאוששות תפעולית מראים שתור מסד הנתונים כבר לא עומד בחוזה האספקה. החלטות קיבולת צריכות להתבסס על עומס העבודה שנצפה, ולא על טענה כללית של קצב בקשות.

לכתבה קודמת

תוכנת מעקב השותפים הטובה ביותר ל-iGaming בשנת 2026

הכתבה הבאה

בינה מלאכותית במשחקי iGaming: מקרי שימוש ב-ChatGPT עבור בתי קזינו

קיסר פיקסון
מְחַבֵּר:

קיסר פיקסון

אני אנליסט נתונים של משחקים מקוונים המתמחה בבחינה ופירוש נתונים הקשורים לפלטפורמות משחקים מקוונות ופעילויות הימורים וכן למגמות שוק. אני מנתח התנהגות שחקנים, ביצועי משחקים ומגמות הכנסות כדי לייעל את חוויות המשחק ואת האסטרטגיות העסקיות.

בקש הדגמה
שלב 1 OF 3
תודה - אתה בתור.
מהנדס פתרונות של NowG ייצור איתך קשר תוך יום עסקים אחד כדי לתאם את הסיור שלך.
מדד