Keyboard shortcuts

Press ← or → to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

Batch job: chạy lại an toàn sau lỗi giữa chừng

Câu hỏi: worker đã ghi hiệu ứng nhưng chưa nhận được ACK; chạy lại có cộng tiền lần nữa không?

Cần biết trước: SQL transaction, MVCC, deadlock/retry. Chép fixture cách A của lab database vào thư mục trống. Lab dùng Python3.14.4, PostgreSQL18.6, dữ liệu đơn hàng giả và hai connection thật; không dùng queue/cloud.

Chọn identity và checkpoint

Một event có event_id ổn định cùng amount cent. Cùng ID/cùng payload là replay; cách xử lý idempotent giữ cùng kết quả khi replay input đã thành công. cùng ID/amount khác là lỗi identity, không âm thầm coi là đã xử lý. Mỗi effect có PK event_id và tham chiếu job; tổng tiền lấy từ effect, không cộng một biến tổng rồi quên ghi identity. Constraints.

Trạng tháiÝ nghĩa trong labAi nhìn thấy?
pendingCó input, chưa checkpoint thành côngConnection mới thấy
runningWorker giữ row lock, đang xử lý trong transactionWorker thấy; chưa commit nên reader khác còn thấy pending
completedEffect và trạng thái cùng commitReplay thấy và trả already
failedInput amount âm, lỗi vĩnh viễn được ghi nhậnReplay trả failed, cần sửa/duyệt input theo policy

Checkpoint là commit cả effect lẫn completed, không phải log “xong”. Fault trước commit rollback về pending. Running không lưu riêng, nên không cần lease recovery trong demo này. Với job lâu có checkpoint nhiều transaction, phải thiết kế ownership, lease/fencing và phục hồi; không bê nguyên transaction dài của lab sang production.

Mã lab tự chứa

Lưu file dưới đây. FOR UPDATE khóa job trước khi xem terminal state. Worker sau chờ worker trước commit và đọc row mới ở READ COMMITTED. Row lock.

from __future__ import annotations

import os
import subprocess
import sys
import time
from collections.abc import Callable
from pathlib import Path

LAB_DIR = os.environ["LAB_DIR"]
COMMAND = [
    "psql",
    "-X",
    "-q",
    "-w",
    "-h",
    LAB_DIR,
    "-p",
    "5432",
    "-U",
    "lab",
    "-d",
    "postgres",
    "-At",
    "-v",
    "ON_ERROR_STOP=1",
    "-v",
    "VERBOSITY=verbose",
]


def environment(name: str = "batch-monitor") -> dict[str, str]:
    env = dict(os.environ)
    for key in ("PGSERVICE", "PGHOSTADDR", "PGPASSWORD", "PGOPTIONS"):
        env.pop(key, None)
    env.update(
        PGAPPNAME=name,
        PGSERVICEFILE="/dev/null",
        PGPASSFILE=str(Path(LAB_DIR) / "no-pass"),
        PGOPTIONS="-c statement_timeout=7000 -c lock_timeout=5000 -c client_min_messages=warning",
    )
    return env


def sql(statement: str) -> str:
    return subprocess.run(
        COMMAND,
        input=statement,
        text=True,
        capture_output=True,
        check=True,
        timeout=10,
        env=environment(),
    ).stdout.strip()


def reset() -> None:
    sql(
        "TRUNCATE wiki_lab.batch_effects,wiki_lab.batch_jobs; INSERT INTO wiki_lab.batch_jobs VALUES (1,100,'pending'),(2,-1,'pending');"
    )


def state() -> str:
    return sql(
        "SELECT status,(SELECT count(*) FROM wiki_lab.batch_effects) FROM wiki_lab.batch_jobs WHERE event_id=1;"
    )


class AckLost(RuntimeError):
    pass


def retry(fault: str, always: bool = False) -> tuple[str, int]:
    for attempt in range(3):
        try:
            current = fault if always or attempt == 0 else "none"
            sql_fault = "none" if current == "after_commit" else current
            result = sql(
                f"BEGIN; SELECT wiki_lab.apply_event(1,100,'{sql_fault}',0); COMMIT;"
            )
            if current == "after_commit":
                raise AckLost("commit xong nhưng client mất ack giả lập")
            return result, attempt + 1
        except subprocess.CalledProcessError as error:
            if "40001" not in error.stderr:
                raise
        except AckLost:
            pass
        if attempt < 2:
            time.sleep(0.01 * 2**attempt)
    raise RuntimeError("retry exhausted: 3")


def wait_for(predicate: Callable[[], bool]) -> None:
    deadline = time.monotonic() + 1.5
    while time.monotonic() < deadline:
        if predicate():
            return
        time.sleep(0.01)
    raise TimeoutError("không quan sát được barrier")


def race() -> None:
    reset()
    workers: list[subprocess.Popen[str]] = []
    try:
        workers.append(
            subprocess.Popen(
                [*COMMAND, "-c", "SELECT wiki_lab.apply_event(1,100,'none',3);"],
                stdout=subprocess.PIPE,
                stderr=subprocess.PIPE,
                text=True,
                env=environment("batch-a"),
            )
        )
        wait_for(
            lambda: (
                sql(
                    "SELECT count(*) FROM pg_stat_activity WHERE application_name='batch-a' AND wait_event='PgSleep';"
                )
                == "1"
            )
        )
        workers.append(
            subprocess.Popen(
                [*COMMAND, "-c", "SELECT wiki_lab.apply_event(1,100,'none',0);"],
                stdout=subprocess.PIPE,
                stderr=subprocess.PIPE,
                text=True,
                env=environment("batch-b"),
            )
        )
        wait_for(
            lambda: (
                sql(
                    "SELECT count(*) FROM pg_stat_activity b WHERE b.application_name='batch-b' AND EXISTS(SELECT 1 FROM pg_stat_activity a WHERE a.application_name='batch-a' AND a.pid=ANY(pg_blocking_pids(b.pid)));"
                )
                == "1"
            )
        )
        outputs = [worker.communicate(timeout=10)[0].strip() for worker in workers]
        assert all(worker.returncode == 0 for worker in workers)
        assert outputs == ["applied", "already"] and state() == "completed|1"
        assert sql("SELECT sum(amount) FROM wiki_lab.batch_effects;") == "100"
        sys.stdout.write(
            "race: blocked_B_by_A=observed applied/already effect_count=1 sum=100\n"
        )
    finally:
        for worker in workers:
            if worker.poll() is None:
                worker.terminate()
                try:
                    worker.communicate(timeout=2)
                except subprocess.TimeoutExpired:
                    worker.kill()
                    worker.communicate(timeout=2)


def main() -> None:
    assert Path(LAB_DIR, ".wiki-lab").is_file()
    assert sql("SHOW unix_socket_directories;") == LAB_DIR
    assert sql("SHOW server_version_num;") == "180006"
    sql("""
CREATE TABLE wiki_lab.batch_jobs (
 event_id integer PRIMARY KEY, amount integer NOT NULL,
 status text NOT NULL CHECK(status IN ('pending','running','completed','failed'))
);
CREATE TABLE wiki_lab.batch_effects (
 event_id integer PRIMARY KEY REFERENCES wiki_lab.batch_jobs(event_id),
 amount integer NOT NULL CHECK(amount>=0)
);
CREATE FUNCTION wiki_lab.apply_event(p_id integer,p_amount integer,p_fault text,p_delay double precision)
RETURNS text LANGUAGE plpgsql AS $body$
DECLARE item wiki_lab.batch_jobs%ROWTYPE;
BEGIN
 SELECT * INTO STRICT item FROM wiki_lab.batch_jobs WHERE event_id=p_id FOR UPDATE;
 IF item.amount<>p_amount THEN RAISE EXCEPTION 'identity payload mismatch' USING ERRCODE='22000'; END IF;
 IF item.status='completed' THEN RETURN 'already'; END IF;
 IF item.status='failed' THEN RETURN 'failed'; END IF;
 IF p_amount<0 THEN UPDATE wiki_lab.batch_jobs SET status='failed' WHERE event_id=p_id; RETURN 'failed'; END IF;
 UPDATE wiki_lab.batch_jobs SET status='running' WHERE event_id=p_id;
 PERFORM pg_sleep(p_delay);
 IF p_fault='before_write' THEN RAISE EXCEPTION 'fault before' USING ERRCODE='40001'; END IF;
 INSERT INTO wiki_lab.batch_effects VALUES(p_id,p_amount);
 IF p_fault='after_write' THEN RAISE EXCEPTION 'fault after' USING ERRCODE='40001'; END IF;
 UPDATE wiki_lab.batch_jobs SET status='completed' WHERE event_id=p_id;
 RETURN 'applied';
END;
$body$;
""")
    for fault in ("before_write", "after_write", "after_commit"):
        reset()
        if fault != "after_commit":
            try:
                sql(f"SELECT wiki_lab.apply_event(1,100,'{fault}',0);")
            except subprocess.CalledProcessError as error:
                assert "40001" in error.stderr and state() == "pending|0"
            else:
                raise AssertionError("fault không xảy ra")
        result, attempts = retry(fault)
        assert attempts == 2 and state() == "completed|1"
        assert result == ("already" if fault == "after_commit" else "applied")
        assert sql("SELECT wiki_lab.apply_event(1,100,'none',0);") == "already"
        sys.stdout.write(f"{fault}: attempts=2 replay=already effect_count=1\n")
    reset()
    try:
        retry("before_write", always=True)
    except RuntimeError as error:
        assert str(error) == "retry exhausted: 3" and state() == "pending|0"
    else:
        raise AssertionError("retry không dừng")
    assert sql("SELECT wiki_lab.apply_event(2,-1,'none',0);") == "failed"
    assert sql("SELECT wiki_lab.apply_event(2,-1,'none',0);") == "failed"
    try:
        sql("SELECT wiki_lab.apply_event(1,999,'none',0);")
    except subprocess.CalledProcessError as error:
        assert "22000" in error.stderr and state() == "pending|0"
    else:
        raise AssertionError("payload conflict không bị từ chối")
    sys.stdout.write(
        "exhausted=3 pending|0 permanent_failed=yes payload_conflict=22000\n"
    )
    race()
    race()
    sql(
        "DROP FUNCTION wiki_lab.apply_event(integer,integer,text,double precision); DROP TABLE wiki_lab.batch_effects,wiki_lab.batch_jobs;"
    )


if __name__ == "__main__":
    main()

Dựng, chạy và kết quả

set -euo pipefail
. ./lab-local.sh
lab_up
lab_seed
lab_whoami
lab_versions
set -euo pipefail
. ./lab.env
python3 batch.py
before_write: attempts=2 replay=already effect_count=1
after_write: attempts=2 replay=already effect_count=1
after_commit: attempts=2 replay=already effect_count=1
exhausted=3 pending|0 permanent_failed=yes payload_conflict=22000
race: blocked_B_by_A=observed applied/already effect_count=1 sum=100

Fault40001 được tiêm để kiểm retry path, chưa là một SSI conflict thật. Mất ACK là exception phía client sau commit, chưa cắt mạng. Hai worker là psql thật, barrier xác nhận A đã giữ lock và B bị A chặn, không suy từ sleep rằng đã có race. Race chạy hai lượt từ dữ liệu reset; unique effect và tổng100 kiểm hậu điều kiện.

Cleanup cả khi thất bại

set -euo pipefail
. ./lab-local.sh
if [ -f lab.env ]; then
  . ./lab.env
  old_dir=$LAB_DIR
  lab_clean
  [ ! -e "$old_dir" ]
fi
lab_clean
[ ! -e lab.env ]
echo 'cleanup ok'

Retry phải hữu hạn và đúng lỗi

Lab retry40001 hoặc ACK giả lập, tối đa3lượt, backoff0.01/0.02s. Lỗi22000 do payload đổi không retry; amount âm thành failed để kiểm/xử lý input. SQLSTATE giúp phân loại, nhưng production phải có policy cho từng lỗi thực tế, deadline tổng và jitter phù hợp. Exhausted giữ pending; một scheduler khác có thể thử lại theo ngân sách đã duyệt.

“Đã thấy job tồn tại” rồi ghi effect ngoài transaction là check-then-act, hai worker vẫn có thể cùng làm. Trong lab, khóa và effect/checkpoint cùng transaction mới giữ invariant. Unique constraint là lớp bảo vệ thêm, không thay thiết kế identity.

Hiệu ứng ngoài database (gửi email, charge API, file/cloud) không nằm trong commit này. Cần idempotency key của bên nhận hoặc outbox và consumer dedupe; không tuyên bố exactly-once cho toàn pipeline. Checkpoint từng item cũng khác “cả file atomic”.

Loại việc này hay gặp ở dịch vụ thu thập dữ liệu định kỳ từ nhiều nguồn: chạy lại sau lỗi không được tạo bản ghi trùng. JobCollect, sản phẩm khác của tác giả, tổng hợp tin tuyển dụng từ nhiều nền tảng theo mô tả trên trang công khai; bài này không mô tả cách JobCollect được xây dựng và lab không dùng nó.

Học tiếp: bulk import, transaction isolation, ADR.