Execution labs: kiểm chứng ở ranh giới lỗi
Không dừng ở “message đã được gửi”. Hãy làm hệ thống hỏng đúng thời điểm, giữ lại dấu vết transport và dùng dữ liệu nghiệp vụ để chứng minh điều gì thực sự xảy ra.
0. Chuẩn bị một thí nghiệm có thể lặp lại
Đọc phần chuẩn bị một lần, sau đó chạy A → B → C → D. A, B và D thực hiện trên cả RabbitMQ và Kafka; C tập trung vào Kafka. D tích hợp thêm degradation/capacity, không tạo lab thứ năm.
Phạm vi code: cấu hình broker, SQL, class Failpoint và replay.py có nội dung cụ thể để sao chép. Các đoạn ghi integration sketch là phần người học phải nối vào consumer/producer/relay của mình; trang này không giả định có sẵn một JAR, endpoint hay project ứng dụng trong ZIP.
| Lớp | Hợp đồng của sandbox | Phải ghi lại trước khi chạy |
|---|---|---|
| Runtime | Bash trên Linux/macOS hoặc WSL2; Docker + Compose v2; JDK 21 cho hook; Python 3.11+ cho replay; curl, jq. | OS, CPU, RAM cấp cho Docker/JVM, disk type/free space, timezone, clock skew, Docker/Compose/JDK/client version. |
| Database | PostgreSQL; một DB exec83a. Inbox và business effect nằm trong cùng transaction của consumer. | Schema/migration, isolation level, pool size, transaction timeout, constraints và application commit SHA. |
| RabbitMQ | 3 node; vhost exec83a; durable direct exchanges; quorum queues với 3 members đã xác nhận; manual ack. | Queue membership/leader, effective policies, prefetch, persistence, confirm mode, retry/DLQ route và permissions. |
| Kafka | 3 combined KRaft broker/controller cho sandbox; topic 3 partitions, RF=3, min.insync.replicas=2; producer acks=all. | ISR/leader từng partition; client assignor/group protocol; batching, linger, compression; auto commit tắt. |
| Cô lập | Chỉ dữ liệu giả; port bind loopback; không public Internet, không dùng credential production. | TLS/SASL/auth/ACL đang bật hay tắt. Kết quả plaintext local không đại diện deployment production. |
Mở bộ cấu hình sandbox và lệnh khởi tạo
Lưu khối sau thành compose.yaml trong một thư mục sandbox mới. Đây là topology local 7 container. Cấu hình Kafka dùng listeners nội bộ cho CLI trong container và listeners loopback cho ứng dụng chạy trên host; không dùng localhost của host từ một app container khác.
name: exec83a
x-rabbit: &rabbit
image: rabbitmq:4.1.0-management
environment: &rabbit-env
RABBITMQ_DEFAULT_USER: lab
RABBITMQ_DEFAULT_PASS: lab_local_only
RABBITMQ_ERLANG_COOKIE: exec83a_local_cookie
healthcheck:
test: ["CMD", "rabbitmq-diagnostics", "-q", "ping"]
interval: 5s
timeout: 5s
retries: 30
x-kafka: &kafka
image: apache/kafka:4.0.0
environment: &kafka-env
CLUSTER_ID: 4L6g3nShT-eMCtK--X86sw
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka1:9093,2@kafka2:9093,3@kafka3:9093
KAFKA_LISTENERS: INTERNAL://:9092,CONTROLLER://:9093,EXTERNAL://:19092
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: INTERNAL:PLAINTEXT,CONTROLLER:PLAINTEXT,EXTERNAL:PLAINTEXT
KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL
KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 3
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 3
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 2
KAFKA_DEFAULT_REPLICATION_FACTOR: 3
KAFKA_MIN_INSYNC_REPLICAS: 2
KAFKA_AUTO_CREATE_TOPICS_ENABLE: "false"
KAFKA_LOG_DIRS: /tmp/kraft-combined-logs
services:
pg:
image: postgres:17.4
environment:
POSTGRES_DB: exec83a
POSTGRES_USER: lab
POSTGRES_PASSWORD: lab_local_only
ports: ["127.0.0.1:5432:5432"]
volumes: ["pg-data:/var/lib/postgresql/data"]
healthcheck:
test: ["CMD-SHELL", "pg_isready -U lab -d exec83a"]
interval: 5s
timeout: 5s
retries: 30
rabbit1:
<<: *rabbit
hostname: rabbit1
ports: ["127.0.0.1:5672:5672", "127.0.0.1:15672:15672"]
volumes: ["rabbit1-data:/var/lib/rabbitmq"]
rabbit2:
<<: *rabbit
hostname: rabbit2
ports: ["127.0.0.1:5673:5672", "127.0.0.1:15673:15672"]
volumes: ["rabbit2-data:/var/lib/rabbitmq"]
rabbit3:
<<: *rabbit
hostname: rabbit3
ports: ["127.0.0.1:5674:5672", "127.0.0.1:15674:15672"]
volumes: ["rabbit3-data:/var/lib/rabbitmq"]
kafka1:
<<: *kafka
hostname: kafka1
ports: ["127.0.0.1:19092:19092"]
environment:
<<: *kafka-env
KAFKA_NODE_ID: 1
KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka1:9092,EXTERNAL://localhost:19092
kafka2:
<<: *kafka
hostname: kafka2
ports: ["127.0.0.1:29092:19092"]
environment:
<<: *kafka-env
KAFKA_NODE_ID: 2
KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka2:9092,EXTERNAL://localhost:29092
kafka3:
<<: *kafka
hostname: kafka3
ports: ["127.0.0.1:39092:19092"]
environment:
<<: *kafka-env
KAFKA_NODE_ID: 3
KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka3:9092,EXTERNAL://localhost:39092
volumes:
pg-data:
rabbit1-data:
rabbit2-data:
rabbit3-data:Kafka lưu log trong filesystem container để ví dụ không phụ thuộc mount permissions. kill/start giữ filesystem; down, xóa hoặc recreate container có thể xóa log. Không dùng các thao tác đó trong phép thử phục hồi dữ liệu. PostgreSQL/RabbitMQ dùng named volumes; không xóa volumes trước khi chụp evidence.
mkdir -p evidence/{A,B,C,D}
docker compose config --quiet
docker compose up -d
# Wait for PostgreSQL and every RabbitMQ node to become ready.
docker compose exec -T pg pg_isready -U lab -d exec83a
for node in rabbit1 rabbit2 rabbit3; do
until docker compose exec -T "$node" rabbitmq-diagnostics -q ping; do
sleep 2
done
done
# FIRST BOOT ONLY: rabbit2/rabbit3 must be empty disposable lab nodes.
# Do not run reset on an existing cluster or on a node containing evidence.
for node in rabbit2 rabbit3; do
docker compose exec -T "$node" rabbitmqctl stop_app
docker compose exec -T "$node" rabbitmqctl reset
docker compose exec -T "$node" rabbitmqctl join_cluster rabbit@rabbit1
docker compose exec -T "$node" rabbitmqctl start_app
done
docker compose exec -T rabbit1 rabbitmqctl cluster_status
docker compose exec -T rabbit1 rabbitmqctl add_vhost exec83a
docker compose exec -T rabbit1 rabbitmqctl set_permissions -p exec83a lab '.*' '.*' '.*'
# Reuse these Bash functions in later sections.
kcli() { docker compose exec -T kafka1 /opt/kafka/bin/"$@"; }
db() { docker compose exec -T pg psql -X -v ON_ERROR_STOP=1 -U lab -d exec83a "$@"; }
until kcli kafka-topics.sh --bootstrap-server kafka1:9092 --list; do
sleep 2
done
for topic in lab.a.events lab.b.events lab.b.retry lab.b.dlq lab.c.events lab.d.events; do
kcli kafka-topics.sh --bootstrap-server kafka1:9092 \
--create --if-not-exists --topic "$topic" \
--partitions 3 --replication-factor 3 \
--config min.insync.replicas=2 \
--config cleanup.policy=delete \
--config retention.ms=604800000
done
export RMQ_API=http://127.0.0.1:15672
export RMQ_AUTH=lab:lab_local_only
export RMQ_URL=amqp://lab:lab_local_only@127.0.0.1:5672/exec83a
export KAFKA_BOOTSTRAP=localhost:19092,localhost:29092,localhost:39092
# Snapshot effective config, not just the desired YAML.
docker version > evidence/docker-version.txt
docker compose version > evidence/compose-version.txt
docker compose images > evidence/images.txt
docker image inspect rabbitmq:4.1.0-management apache/kafka:4.0.0 postgres:17.4 \
--format '{{json .RepoDigests}}' > evidence/image-digests.txt
docker compose exec -T rabbit1 rabbitmqctl version > evidence/rabbit-version.txt
kcli kafka-topics.sh --version > evidence/kafka-version.txt
db -c 'SELECT version();' > evidence/postgres-version.txt
kcli kafka-topics.sh --bootstrap-server kafka1:9092 --describe \
> evidence/kafka-topology.txtĐợi các topic có đủ ISR=3 rồi mới làm negative control. Nếu startup lỗi, lưu log và giải quyết trước; không diễn giải lỗi bootstrap thành kết quả crash lab. set_permissions .* là quyền admin rộng trong vhost lab; không sao chép sang production.
# Create separate RabbitMQ routes for Labs A, B and D.
for lab in a b d; do
prefix="lab.$lab"
for suffix in events dlx retry; do
curl --fail-with-body -sS -u "$RMQ_AUTH" \
-H 'content-type: application/json' -X PUT \
"$RMQ_API/api/exchanges/exec83a/$prefix.$suffix" \
-d '{"type":"direct","durable":true,"auto_delete":false,"internal":false,"arguments":{}}'
done
for suffix in main retry.5s dlq; do
curl --fail-with-body -sS -u "$RMQ_AUTH" \
-H 'content-type: application/json' -X PUT \
"$RMQ_API/api/queues/exec83a/$prefix.$suffix" \
-d '{"durable":true,"auto_delete":false,"arguments":{"x-queue-type":"quorum","x-quorum-initial-group-size":3}}'
done
for route in 'events main work' 'retry retry.5s retry' 'dlx dlq dead'; do
read -r exchange queue key <<< "$route"
curl --fail-with-body -sS -u "$RMQ_AUTH" \
-H 'content-type: application/json' -X POST \
"$RMQ_API/api/bindings/exec83a/e/$prefix.$exchange/q/$prefix.$queue" \
-d "{\"routing_key\":\"$key\",\"arguments\":{}}"
done
main_policy=$(jq -nc --arg dlx "$prefix.dlx" \
'{"dead-letter-exchange":$dlx,"dead-letter-routing-key":"dead","dead-letter-strategy":"at-least-once","overflow":"reject-publish","delivery-limit":20}')
retry_policy=$(jq -nc --arg events "$prefix.events" \
'{"message-ttl":5000,"dead-letter-exchange":$events,"dead-letter-routing-key":"work","dead-letter-strategy":"at-least-once","overflow":"reject-publish"}')
docker compose exec -T rabbit1 rabbitmqctl set_policy -p exec83a \
"exec83a-$lab-main" "^lab[.]$lab[.]main$" "$main_policy" \
--apply-to quorum_queues --priority 10
docker compose exec -T rabbit1 rabbitmqctl set_policy -p exec83a \
"exec83a-$lab-delay" "^lab[.]$lab[.]retry[.]5s$" "$retry_policy" \
--apply-to quorum_queues --priority 10
docker compose exec -T rabbit1 rabbitmq-queues quorum_status \
--vhost exec83a "$prefix.main"
done
curl --fail-with-body -sS -u "$RMQ_AUTH" \
"$RMQ_API/api/queues/exec83a" | jq > evidence/rabbit-queues.json
docker compose exec -T rabbit1 rabbitmqctl list_policies -p exec83a \
> evidence/rabbit-policies.txtQuorum membership không tự được chứng minh bởi số node trong cluster. Kiểm tra thật bằng quorum_status. At-least-once dead lettering cần cấu hình tương ứng; TTL 5 giây là khoảng chờ tối thiểu theo topology này, không phải cam kết message sẽ được xử lý sau đúng 5 giây. [R2], [R3], [R4], [R9], [R10]
| Lab | RabbitMQ route | Kafka topic |
|---|---|---|
| A | lab.a.events → key work → lab.a.main | lab.a.events |
| B | lab.b.events → lab.b.main; retry/DLQ cùng prefix lab.b | lab.b.events, lab.b.retry, lab.b.dlq |
| C | Không dùng RabbitMQ trong lab này. | lab.c.events |
| D | lab.d.events → key work → lab.d.main | lab.d.events |
Đối chiếu bản lỗi/bản sửa bằng DB consumer namespace riêng. Với group Kafka mới, ghi mốc bắt đầu từng partition trước khi publish fixture, hoặc tạo topic case mới cùng config; không dùng earliest rồi bỏ qua việc group sẽ đọc cả record của case cũ. Giữa các case, lưu offsets/log và xác nhận không còn công việc ngoài fixture đang chạy.
0.1. Định danh, log và oracle
Tạo bộ fixture trước khi publish và lưu vào expected_event. Mỗi case dùng ID mới, nhưng redelivery, relay retry và replay của cùng sự kiện phải giữ event_id, business_key và origin_run_id. run_id của lần thực thi replay là metadata riêng, không được làm thay đổi danh tính nghiệp vụ.
{
"event_id": "11111111-1111-4111-8111-111111111111",
"origin_run_id": "A-fixed-after-commit",
"business_key": "tenant-demo:charge:order-001",
"aggregate_id": "order-001",
"version": 1,
"schema_version": 1,
"amount": 10,
"occurred_at": "2026-09-12T00:00:00Z"
}business_key định danh một thao tác nghiệp vụ, không chỉ một aggregate: hai lần thanh toán hợp lệ phải có key khác nhau. Phạm vi dedup phải chứa tenant và tên consumer/projection. Hash lưu trong Inbox/effect là fingerprint của nội dung nghiệp vụ bất biến theo canonicalization đã công bố (ví dụ business key, aggregate, amount, schema/version); không đưa event_id, origin_run_id, attempt, timestamp giao nhận hay replay headers vào fingerprint này. Cùng ID mà khác fingerprint phải bị báo lỗi, không được im lặng coi là duplicate.
Transport evidence
JSONL append-only: run_id, event_id, stage, wall_time_utc, monotonic_ns, process_id, broker. RabbitMQ thêm connection/channel ID, delivery tag, redelivered, confirm/nack/return; Kafka thêm topic, partition, offset, group/member, assignment và committed next offset. Log phải flush trước crash hook.
Business correctness
Fixture kỳ vọng độc lập với consumer; query số effect, amount, business key, version và payload hash sau recovery. Phải kiểm cả missing lẫn duplicate: bảng chỉ có một dòng không chứng minh một message khác chưa bị mất.
Mở schema dùng chung — lưu thành schema.sql
CREATE TABLE expected_event (
consumer_name text NOT NULL,
event_id uuid NOT NULL,
origin_run_id text NOT NULL,
business_key text NOT NULL,
amount bigint NOT NULL,
PRIMARY KEY (consumer_name, event_id)
);
CREATE TABLE inbox (
consumer_name text NOT NULL,
event_id uuid NOT NULL,
payload_hash text NOT NULL,
created_at timestamptz NOT NULL DEFAULT clock_timestamp(),
PRIMARY KEY (consumer_name, event_id)
);
CREATE TABLE business_effect (
effect_id bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
consumer_name text NOT NULL,
event_id uuid NOT NULL,
origin_run_id text NOT NULL,
business_key text NOT NULL,
amount bigint NOT NULL,
payload_hash text NOT NULL,
created_at timestamptz NOT NULL DEFAULT clock_timestamp()
);
-- Keep broken-run evidence separate from the fixed implementation.
CREATE TABLE business_effect_broken (LIKE business_effect INCLUDING ALL);
-- Lab A adds the unique business index on business_effect only.
CREATE TABLE projection_state (
projection_name text NOT NULL,
aggregate_id text NOT NULL,
version bigint NOT NULL,
state jsonb NOT NULL,
updated_at timestamptz NOT NULL DEFAULT clock_timestamp(),
PRIMARY KEY (projection_name, aggregate_id)
);
CREATE TABLE business_order (
order_id text PRIMARY KEY,
state text NOT NULL,
version bigint NOT NULL
);
CREATE TABLE outbox (
event_id uuid PRIMARY KEY,
origin_run_id text NOT NULL,
aggregate_id text NOT NULL,
payload jsonb NOT NULL,
created_at timestamptz NOT NULL DEFAULT clock_timestamp(),
sent_at timestamptz,
lease_owner text,
lease_until timestamptz,
lease_epoch bigint NOT NULL DEFAULT 0,
attempts integer NOT NULL DEFAULT 0
);
CREATE INDEX outbox_pending ON outbox (created_at, event_id)
WHERE sent_at IS NULL;db < schema.sql0.2. Ghi ngưỡng trước, đo kết quả sau
Trước mỗi case, ghi seed, số message N, payload distribution, producer rate, concurrency, thời gian warm-up/measurement và giới hạn tài nguyên. Chọn trước SLO xử lý, giới hạn retry, backlog, recovery window, replay rate và ngưỡng dừng vì disk/RAM. Các số như delay 5 giây, 3 attempts, lease 30 giây trong trang chỉ là tham số thí nghiệm, không phải benchmark hay khuyến nghị phổ quát.
Gate về correctness là bất biến; gate hiệu năng do người học chốt trước khi chạy. Lặp ít nhất ba lần mỗi case ổn định, giữ raw samples thay vì chỉ ảnh dashboard. Sau negative control, archive dữ liệu lỗi bằng consumer/run namespace riêng; không “sửa” kết quả bằng cách xóa hàng trùng rồi tuyên bố pass.
Lab A · Crash trước/sau ack
Goal: phân biệt “được giao lại” với “bị thực thi hai lần”, và chứng minh một thao tác nghiệp vụ không bị mất hoặc bị áp dụng lặp khi consumer chết ở từng ranh giới.
- Crash trước DB commit, sau DB commit nhưng trước ack/offset commit, và sau ack.
- Quan sát redelivery/loss semantics.
- Thêm Inbox/unique business key cùng transaction với effect.
- Assert duplicate deliveries chỉ tạo một effect.
A1. Failure-first: tạo lỗi trước khi thêm Inbox
Dùng business_effect_broken, append effect mỗi lần nhận message, chưa dedup. Gửi một event amount=10; commit DB rồi crash trước ack/offset commit. Khởi động lại cùng RabbitMQ queue hoặc cùng Kafka group. Quan sát số lần giao và số effect. Thêm negative control thứ hai: ack/commit offset trước DB transaction, chờ bằng chứng broker đã ghi nhận, rồi crash trước business write.
FAILPOINT, nếu không consumer sẽ crash lại vô hạn.public final class Failpoint {
private Failpoint() {}
public static void hit(String stage, String eventId) {
String active = System.getenv("FAILPOINT");
String target = System.getenv("FAIL_EVENT_ID");
if (!stage.equals(active) || !eventId.equals(target)) return;
System.err.printf(
"{\"stage\":\"%s\",\"event_id\":\"%s\",\"pid\":%d,\"nano\":%d}%n",
stage, eventId, ProcessHandle.current().pid(), System.nanoTime());
System.err.flush();
// Deliberate crash: no shutdown hook, finally block or graceful ack.
Runtime.getRuntime().halt(137);
}
}javac Failpoint.java
# Run the application's real start command with these environment variables:
export FAILPOINT=AFTER_DB_COMMIT
export FAIL_EVENT_ID=11111111-1111-4111-8111-111111111111
# Start YOUR consumer here, then publish the matching fixture event.
# After the deliberate exit:
unset FAILPOINT FAIL_EVENT_ID
# Restart the SAME consumer group/queue with your usual application command.Hook dùng Runtime.halt để không chạy graceful shutdown; chỉ bật trong sandbox. [R17] Hook BEFORE_DB_COMMIT đặt sau business SQL nhưng trước commit thật. Hook AFTER_DB_COMMIT đặt ở caller sau khi transaction đã hoàn tất. Nếu dùng Spring proxy, self-invocation có thể không đi qua transaction interceptor; không gắn nhãn “committed” chỉ vì đã chạy hết thân method. [R11]
A2. Crash matrix và quan sát cần ghi
| Case | Tiêm lỗi | Expected transport | Expected business / cách kết luận |
|---|---|---|---|
| A−1 · không dedup | DB commit → crash → restart. | RabbitMQ có thể redeliver; Kafka đọc lại từ committed next offset cũ. | Negative control phải cho thấy cùng operation có >1 effect. Nếu không tái hiện được delivery lặp, case chưa đủ bằng chứng. |
| A−2 · ack quá sớm | Broker đã nhận ack/offset commit → crash trước DB write. | Message không còn được giao lại tự nhiên trong queue/group hiện tại; Kafka log vẫn có thể dùng để replay thủ công. | Oracle phát hiện missing effect. Đây là mất xử lý ở application, không nhất thiết broker đã mất record. |
| A1 · trước DB commit | Halt sau SQL nhưng trước commit. | Chưa ack/commit offset; chờ broker phát hiện mất consumer rồi giao lại. | Transaction cũ rollback. Sau recovery có đúng một effect. |
| A2 · sau DB commit | Commit đã thành công → halt trước ack/offset commit. | Redelivery/re-read có thể xảy ra; cùng logical event ID. | Inbox/effect đã tồn tại. Delivery lặp phải kết thúc bằng dedup, không cộng amount lần hai. |
| A3 · sau ack call | Ack/commit sau DB commit → halt. | Kafka: commitSync thành công là mốc đã xác nhận. RabbitMQ: ACK_SENT chỉ là client call, ack có thể chưa được broker xử lý. | Một effect dù có redelivery hay không. Không được suy luận “sau basicAck thì chắc chắn không giao lại”. |
| A4 · commit outcome unknown | Ngắt kết nối DB đúng lúc commit; giữ log trước/sau và query bằng connection mới. | Không ack khi chưa xác minh được outcome; replay cùng ID. | DB có thể đã commit hoặc chưa. Retry với Inbox phải hội tụ về một effect, không tự kết luận timeout là rollback. |
Publisher confirm và consumer ack là hai protocol độc lập. Ghi confirm/nack/timeout cùng event ID; với RabbitMQ dùng mandatory=true và xử lý returned/unroutable message. Với Kafka, acks=all và idempotent producer không làm transaction DB trở thành atomic với offset commit. [R1], [R5]
Transport evidence
Lưu trước/sau crash: RabbitMQ ready/unacked, delivery tag + redelivered và ack sent; Kafka record offset n, group committed next offset, read lại n, commit n+1. Một lần publish lại với cùng event ID có thể mang RabbitMQ redelivered=false; đừng dùng cờ này làm khóa dedup.
Business correctness
Với operation amount=10, count=1 và sum=10 sau recovery. Case ack sớm phải làm query missing trả ra ID tương ứng. Dòng log “processed” hoặc “committed” của ứng dụng không thay thế query bằng một connection DB độc lập.
A3. Sửa atomic boundary: Inbox + unique business key
@Transactional
void consume(Event event) {
if (!inbox.tryInsert(event.id())) return;
applyBusinessEffect(event);
} // ack/offset commit xảy ra sau local transaction theo adapter designThêm Inbox và unique business key cùng local transaction với effect; chỉ ack/commit offset khi transaction thành công. Nếu event đã tồn tại nhưng payload fingerprint khác, báo conflict và đưa vào luồng điều tra. Giữ bảng lỗi business_effect_broken để đối chiếu, không xóa evidence.
Mở SQL sửa idempotency và integration sketch
CREATE UNIQUE INDEX one_business_effect
ON business_effect (consumer_name, business_key);
CREATE OR REPLACE FUNCTION apply_charge(
p_consumer text, p_event uuid, p_run text,
p_key text, p_amount bigint, p_hash text
) RETURNS text LANGUAGE plpgsql AS $$
DECLARE
v_claim uuid;
v_effect bigint;
v_hash text;
v_amount bigint;
BEGIN
INSERT INTO inbox(consumer_name, event_id, payload_hash)
VALUES (p_consumer, p_event, p_hash)
ON CONFLICT (consumer_name, event_id) DO NOTHING
RETURNING event_id INTO v_claim;
IF v_claim IS NULL THEN
SELECT payload_hash INTO v_hash FROM inbox
WHERE consumer_name=p_consumer AND event_id=p_event;
IF v_hash IS DISTINCT FROM p_hash THEN
RAISE EXCEPTION 'EVENT_ID_PAYLOAD_CONFLICT';
END IF;
RETURN 'DUPLICATE_EVENT';
END IF;
INSERT INTO business_effect
(consumer_name,event_id,origin_run_id,business_key,amount,payload_hash)
VALUES (p_consumer,p_event,p_run,p_key,p_amount,p_hash)
ON CONFLICT (consumer_name,business_key) DO NOTHING
RETURNING effect_id INTO v_effect;
IF v_effect IS NULL THEN
SELECT payload_hash, amount INTO v_hash, v_amount
FROM business_effect
WHERE consumer_name=p_consumer AND business_key=p_key;
IF v_hash IS DISTINCT FROM p_hash OR v_amount IS DISTINCT FROM p_amount THEN
RAISE EXCEPTION 'BUSINESS_KEY_PAYLOAD_CONFLICT';
END IF;
RETURN 'DUPLICATE_BUSINESS_OPERATION';
END IF;
RETURN 'APPLIED';
END $$;// Integration sketch: use a separate transactional service/proxy.
// decodeAndValidate and audit are application code, not supplied APIs.
void onDelivery(Delivery d) {
Event e = decodeAndValidate(d);
audit("RECEIVED", e, d);
service.applyInDatabase(e); // Returns AFTER the proxy commits.
audit("DB_COMMITTED", e, d);
Failpoint.hit("AFTER_DB_COMMIT", e.id().toString());
if (d.isRabbit()) {
channel.basicAck(d.deliveryTag(), false); // Same receiving channel.
audit("ACK_SENT", e, d); // Does NOT mean broker-confirmed ack.
} else {
consumer.commitSync(Map.of(
new TopicPartition(d.topic(), d.partition()),
new OffsetAndMetadata(d.offset() + 1)));
audit("OFFSET_COMMIT_CONFIRMED", e, d);
}
Failpoint.hit("AFTER_ACK_CALL", e.id().toString());
}
@Transactional
public void applyInDatabase(Event e) {
// SELECT apply_charge(...); uses THIS transaction/connection.
repository.applyCharge(e);
Failpoint.hit("BEFORE_DB_COMMIT", e.id().toString());
}ON CONFLICT dựa trên unique constraint để phân xử các insert cạnh tranh. Chạy thêm hai consumers nhận cùng event đồng thời, và hai event ID khác nhau nhưng cùng business key/fingerprint. Ở isolation khác hoặc khi gặp deadlock/serialization error, retry toàn transaction với cùng ID; không ack một transaction thất bại. [R8]
A4. Assert một effect và không thiếu effect
Oracle gom theo business_key để hai event ID hợp lệ cùng biểu diễn một thao tác không bị đếm kỳ vọng hai lần. Fixture conflict/permanent-invalid được kiểm ở tập negative riêng, không trộn vào tập expected-success.
-- psql variables: substitute the fixed consumer/run being evaluated.
\set c 'charge-fixed'
\set r 'A-fixed-after-commit'
-- MUST return zero rows after all expected valid events have recovered.
WITH expected_operations AS (
SELECT DISTINCT consumer_name,business_key,amount
FROM expected_event
WHERE consumer_name=:'c' AND origin_run_id=:'r'
)
SELECT e.business_key, count(b.effect_id) AS effects,
coalesce(sum(b.amount),0) AS actual_amount, e.amount AS expected_amount
FROM expected_operations e
LEFT JOIN business_effect b
ON b.consumer_name=e.consumer_name AND b.business_key=e.business_key
GROUP BY e.business_key,e.amount
HAVING count(b.effect_id) <> 1 OR coalesce(sum(b.amount),0) <> e.amount;
-- MUST return zero rows: no effect outside the selected run's fixture.
SELECT b.event_id,b.business_key
FROM business_effect b
WHERE b.consumer_name=:'c' AND b.origin_run_id=:'r'
AND NOT EXISTS (
SELECT 1 FROM expected_event e
WHERE e.consumer_name=b.consumer_name
AND e.origin_run_id=b.origin_run_id AND e.business_key=b.business_key
);
-- MUST return zero rows in the fixed namespace.
SELECT consumer_name,business_key,count(*)
FROM business_effect
WHERE consumer_name=:'c'
GROUP BY consumer_name,business_key HAVING count(*) > 1;
-- Negative-control evidence: this query SHOULD expose duplicates.
SELECT consumer_name,business_key,count(*),sum(amount)
FROM business_effect_broken
GROUP BY consumer_name,business_key HAVING count(*) > 1;Acceptance gate A
- Hai negative control tái hiện được duplicate effect và missing effect; có raw evidence, không chỉ mô tả.
- Cả ba cửa sổ nguồn được chạy trên cả hai broker. Bản sửa hội tụ về một effect/operation; query missing, unexpected và duplicate đều trả 0 hàng sau recovery.
- Chạy duplicate đồng thời và conflict fingerprint; conflict không bị coi là thành công. Ack/commit offset chỉ diễn ra sau transaction hoặc quyết định quarantine đã bền vững.
- Matrix phân biệt ACK_SENT, broker-observed ack, committed offset và DB commit; những outcome chưa rõ được đánh dấu UNKNOWN, không đổi thành PASS.
Artifact A — 5 tệp: A/crash-matrix.csv, A/transport.jsonl, A/assertions.sql, A/assertions.txt, A/fix.diff. Matrix chứa broker, case, run, event ID, stage, delivery count, effect count, amount, kết luận và đường dẫn log.
Lab B · Poison, retry và DLQ
Goal: một poison message không tạo vòng lặp vô hạn, không làm mất sự kiện tốt và chỉ được replay khi có owner, điều kiện sửa lỗi cùng dấu vết audit.
- Gửi payload malformed và business-permanent failure.
- So immediate retry, delayed retry và blocking partition/consumer.
- Bound attempts, record reason/original ID/schema/attempts vào DLQ.
- Xây replay command có filter, dry-run, rate limit và audit.
B1. Failure-first: tạo poison ở cả decode và business layer
Tạo fixture M1 có JSON thiếu dấu đóng nhưng có event ID trong message property/header; P1 parse được nhưng vi phạm business rule vĩnh viễn; T1 có dependency giả lập trả lỗi transient ở hai lần gọi đầu rồi thành công. Chèn G1, G2 hợp lệ sau poison: một record cùng Kafka partition và một record ở partition khác. Ghi offsets thực tế, không giả định hai key khác nhau chắc chắn vào hai partition.
Bản lỗi bắt mọi exception rồi requeue ngay RabbitMQ, hoặc seek về cùng Kafka offset mãi. Chạy trong cửa sổ ngắn có watchdog và giới hạn tổng số lần gọi dependency; lưu tốc độ retry, CPU, lag, tuổi event tốt. Đây là negative control có chủ đích, không phải chính sách production.
B2. So ba chiến lược trên cùng fixture
| Chiến lược | RabbitMQ | Kafka | Phép đo / trade-off |
|---|---|---|---|
| Immediate retry | Requeue có giới hạn hoặc retry ngay trong handler; không dùng nack(requeue=true) làm policy vô hạn. | Seek/retry có giới hạn; không commit qua record chưa xử lý. | So invocation rate, CPU, duplicate count, tuổi G1/G2. Lỗi permanent không có lý do để retry như transient. |
| Delayed retry | Publish vào lab.b.retry/retry, chờ confirm rồi ack nguồn. TTL queue + DLX trả về main; giữ ID/attempt metadata. | Đưa vào retry topic và processor/scheduler có not_before; Kafka không tự delay record chỉ vì có header này. | Có thêm handoff crash window; good messages tiếp tục nhưng retry có thể đổi thứ tự nghiệp vụ. Fixed TTL bucket không hứa thời điểm delivery chính xác. |
| Blocking / giữ thứ tự | Giữ message chưa ack; với prefetch=1, consumer/channel đó dừng tiến độ. Consumer khác có thể vẫn chạy. | Pause partition chứa poison; tiếp tục poll cho membership và xử lý partition khác. Không commit vượt lỗ hổng của partition bị chặn. | Bảo toàn thứ tự cục bộ đổi lấy head-of-line blocking. Sleep toàn poll thread lâu có thể gây mất ownership/rebalance. |
Với RabbitMQ, x-death mô tả các lần dead-lettering, không phải bộ đếm chính xác mọi lần callback được gọi. Quorum delivery-limit là lớp chặn bảo vệ bổ sung, không thay application retry budget. Với Kafka, polling/commit phải tuân theo client và group protocol đã chọn. [R2], [R3], [R6]
B3. Bound attempts và đóng handoff window
Đặt tổng budget minh họa là 3 attempts = lần đầu + 2 retries. M1/P1 vào quarantine ngay sau phân loại; T1 chỉ được gọi tối đa ba lần theo ledger bền vững. Lưu attempt trước invocation để crash không reset budget. Phân biệt delivery count, business invocation count và replay generation. Việc nhận lại một terminal event không được tự mở thêm budget.
# Integration sketch: durable retry ledger and adapters are application code.
receive raw delivery
-> recover original ID / raw digest before JSON decoding
-> load terminal state and reserve bounded attempt in durable ledger
-> decode + validate schema + invoke business transaction
-> success: settle source only after DB commit
-> permanent / exhausted:
persist reason + original bytes + immutable metadata
publish DLQ envelope; await transport confirmation
settle source
-> transient with budget remaining:
persist incremented attempt and not_before
publish delayed retry; await transport confirmation
settle source
# Any handoff timeout:
# - keep source unsettled; stop/throttle the failing pipeline
# - reconcile by original ID; do not manufacture a fresh event ID
# - if terminal state already exists, hand off again without invoking business| Crash / fault | Expected observation | Acceptance |
|---|---|---|
| Publish retry/DLQ thành công → crash trước ack nguồn | Nguồn được nhận lại; đích có thể có hơn một bản cùng original ID. | Không gọi business quá budget. Duplicate DLQ entry/retry delivery được nhận diện; business effect không lặp. |
| DLQ target không sẵn sàng | Handoff không hoàn tất, nguồn hoặc pending ledger còn giữ trách nhiệm; tuổi backlog tăng. | Không ack rồi bỏ message. Có cảnh báo, circuit breaker/backpressure và recovery evidence. |
| Native DLX bật nhưng không kiểm safety | At-most-once dead-letter strategy có thể làm mất handoff khi target lỗi; không suy từ “có DLX” ra an toàn. | Dùng policy at-least-once đã kiểm ở setup hoặc application republish-confirm-then-ack. Log actual policy và queue membership. |
| Kafka DLQ produce → crash trước source offset commit | Không dùng transaction: DLQ có thể trùng; source record quay lại. | Giữ original topic/partition/offset + ID để dedup. Nếu dùng Kafka transaction, gửi DLQ và source offsets trong cùng transaction, downstream dùng read_committed; DB vẫn ở boundary riêng. |
| Replay chạy lại chính manifest | Transport có thể publish lại, kể cả khi producer idempotence bật. | Original ID/key/fingerprint không đổi; effect đã đúng phải giữ nguyên. Audit chỉ phản ánh publish, không tự đánh dấu business-resolved. |
Không tự gán application reason vào một message bằng cách chỉ basicNack: nack không mang payload/header cập nhật. Muốn DLQ có business reason, ghi failure ledger hoặc republish một envelope đã enrich rồi confirm trước ack. Native DLX vẫn hữu ích như safety net cho delivery-limit và expiration. [R1], [R3]
B4. Replay command: filter, dry-run, rate limit, audit
Đóng băng một manifest export không làm mất DLQ nguồn; ghi hash và review trước khi chạy. Với Kafka có thể dùng group export riêng; với RabbitMQ dùng failure ledger/mirror đã lưu hoặc một quy trình peek có kiểm soát. Không “drain DLQ để xem” rồi coi đó là backup. Replay tool dưới đây chỉ đọc manifest và publish, không xóa/ack DLQ.
Field bắt buộc: dlq_id, original_id, original_key, origin_run_id, reason, schema_version, attempts, first_failed_at, payload_b64, payload_sha256, replay_approved. Lưu thêm original queue/exchange hoặc topic/partition/offset, exception class, first/last failure, correlation ID và application version trong bản export đầy đủ. Hash raw bytes để phát hiện sửa manifest khác với semantic fingerprint dùng cho Inbox.
python3 -m venv .venv
. .venv/bin/activate
python -m pip install 'pika==1.3.2' 'confluent-kafka==2.10.0'
# Save the full Python block below as replay.py.
# The time window below is a sample: replace it with YOUR evidence timestamps.
# DRY RUN (default): no broker connection, no publish, no DLQ deletion.
python replay.py --manifest evidence/B/dlq-manifest.jsonl --broker kafka --reason SCHEMA_UNKNOWN --schema 2 --from-utc 2026-09-12T00:00:00Z --to-utc 2026-09-13T00:00:00Z --run-id B-plan-01 --owner learner --limit 20 --rate 2 --audit evidence/B/replay-plan.jsonl
# EXECUTE only after consumer schema support and approval are recorded.
python replay.py --manifest evidence/B/dlq-manifest.jsonl --broker kafka --reason SCHEMA_UNKNOWN --schema 2 --from-utc 2026-09-12T00:00:00Z --to-utc 2026-09-13T00:00:00Z --run-id B-replay-01 --owner learner --approval LAB-B-REVIEW-01 --limit 20 --rate 2 --execute --audit evidence/B/replay-audit.jsonl
# RabbitMQ run: use --broker rabbit and a NEW audit path/run-id.
# RMQ_URL / KAFKA_BOOTSTRAP were exported in the shared setup.Mở replay.py — công cụ replay tối thiểu, không phải production replayer
#!/usr/bin/env python3
"""Finite, reviewed DLQ manifest replay for a LOCAL LAB.
Default is dry-run. Transport confirmation is NOT business success.
Requires pika==1.3.2 and confluent-kafka==2.10.0 only for --execute.
"""
import argparse
import base64
import hashlib
import json
import os
import sys
import signal
from contextlib import contextmanager
import time
from datetime import datetime, timezone
from pathlib import Path
@contextmanager
def deadline(seconds):
# Main-thread POSIX watchdog; native Windows must use WSL for this lab.
def expired(_signum, _frame):
raise TimeoutError("Broker operation deadline exceeded")
previous = signal.signal(signal.SIGALRM, expired)
signal.setitimer(signal.ITIMER_REAL, seconds)
try:
yield
finally:
signal.setitimer(signal.ITIMER_REAL, 0)
signal.signal(signal.SIGALRM, previous)
def utc(value):
dt = datetime.fromisoformat(value.replace("Z", "+00:00"))
if dt.tzinfo is None:
raise ValueError("Timestamp must contain a timezone")
return dt.astimezone(timezone.utc)
def audit_write(f, record):
record = {"at": datetime.now(timezone.utc).isoformat(), **record}
f.write(json.dumps(record, ensure_ascii=False) + "\n")
f.flush()
os.fsync(f.fileno())
def parse_args():
p = argparse.ArgumentParser(description=__doc__)
p.add_argument("--manifest", required=True)
p.add_argument("--broker", choices=("rabbit", "kafka"), required=True)
p.add_argument("--reason", required=True)
p.add_argument("--schema", type=int, required=True)
p.add_argument("--from-utc", required=True)
p.add_argument("--to-utc", required=True)
p.add_argument("--run-id", required=True)
p.add_argument("--owner", required=True)
p.add_argument("--approval")
p.add_argument("--limit", type=int, default=100)
p.add_argument("--rate", type=float, default=1.0)
p.add_argument("--audit", required=True)
p.add_argument("--execute", action="store_true")
return p.parse_args()
def main():
a = parse_args()
if a.limit < 1 or not 0 < a.rate <= 1000:
raise ValueError("Require limit >= 1 and 0 < rate <= 1000")
if a.execute and not a.approval:
raise ValueError("--execute requires --approval")
start, end = utc(a.from_utc), utc(a.to_utc)
if start >= end:
raise ValueError("Require from-utc < to-utc; end is exclusive")
manifest_bytes = Path(a.manifest).read_bytes()
manifest_hash = hashlib.sha256(manifest_bytes).hexdigest()
items = [json.loads(line) for line in manifest_bytes.splitlines() if line.strip()]
candidates = []
seen_dlq = set()
for item in items:
if item["dlq_id"] in seen_dlq:
raise ValueError("Duplicate dlq_id in the frozen manifest")
seen_dlq.add(item["dlq_id"])
ts = utc(item["first_failed_at"])
if not (item.get("replay_approved") is True
and item["reason"] == a.reason
and item["schema_version"] == a.schema
and start <= ts < end):
continue
if not item.get("original_id") or not item.get("original_key"):
raise ValueError("Approved replay requires original ID and key")
if not isinstance(item["attempts"], int) or item["attempts"] < 0:
raise ValueError("Invalid original attempt count")
payload = base64.b64decode(item["payload_b64"], validate=True)
if hashlib.sha256(payload).hexdigest() != item["payload_sha256"]:
raise ValueError("Raw payload hash mismatch")
candidates.append((item, payload))
candidates.sort(key=lambda pair: (
utc(pair[0]["first_failed_at"]), pair[0]["dlq_id"]))
selected = candidates[:a.limit]
# Exclusive create prevents accidental reuse/overwrite of an audit run.
with open(a.audit, "x", encoding="utf-8") as log:
audit_write(log, {
"stage": "PLAN", "run_id": a.run_id, "owner": a.owner,
"approval": a.approval, "execute": a.execute,
"manifest_sha256": manifest_hash, "broker": a.broker,
"reason": a.reason, "schema": a.schema,
"from_utc": a.from_utc, "to_utc": a.to_utc,
"rate_per_second": a.rate, "limit": a.limit,
"selected_count": len(selected)})
if not a.execute:
for item, _ in selected:
audit_write(log, {
"stage": "WOULD_PUBLISH", "dlq_id": item["dlq_id"],
"event_id": item["original_id"], "key": item["original_key"]})
print(json.dumps({"mode": "dry-run", "selected": len(selected),
"manifest_sha256": manifest_hash}))
return 0
if not selected:
return 0
connection = producer = channel = None
try:
if a.broker == "rabbit":
import pika
params = pika.URLParameters(os.environ["RMQ_URL"])
params.socket_timeout = 10
params.blocked_connection_timeout = 10
params.heartbeat = 30
connection = pika.BlockingConnection(params)
channel = connection.channel()
channel.confirm_delivery()
else:
from confluent_kafka import Producer
producer = Producer({
"bootstrap.servers": os.environ["KAFKA_BOOTSTRAP"],
"client.id": "replay-" + a.run_id,
"enable.idempotence": True,
"acks": "all",
"delivery.timeout.ms": 10000})
next_due = time.monotonic()
for item, payload in selected:
delay = next_due - time.monotonic()
if delay > 0:
if connection is not None:
connection.sleep(delay) # Keep Rabbit heartbeats alive.
else:
time.sleep(delay)
key = item["original_key"]
metadata = {
"x-original-id": item["original_id"],
"x-origin-run-id": item["origin_run_id"],
"x-replay-run-id": a.run_id,
"x-dlq-id": item["dlq_id"],
"x-attempts": str(item["attempts"]),
"x-schema-version": str(item["schema_version"]),
"x-original-key": key}
audit_write(log, {
"stage": "PUBLISH_PLANNED", "run_id": a.run_id,
"event_id": item["original_id"], "dlq_id": item["dlq_id"],
"payload_sha256": item["payload_sha256"]})
broker_position = {}
try:
if channel is not None:
import pika
with deadline(15):
channel.basic_publish(
exchange="lab.b.events", routing_key="work",
body=payload, mandatory=True,
properties=pika.BasicProperties(
message_id=item["original_id"],
content_type="application/json",
delivery_mode=2, headers=metadata))
broker_position = {"exchange": "lab.b.events",
"routing_key": "work"}
else:
outcomes = []
def delivered(error, message):
outcomes.append((error, message))
producer.produce(
"lab.b.events", key=key.encode(), value=payload,
headers=[(k, v.encode()) for k, v in metadata.items()],
on_delivery=delivered)
remaining = producer.flush(15)
if remaining or not outcomes or outcomes[0][0] is not None:
raise RuntimeError("Kafka delivery not confirmed")
m = outcomes[0][1]
broker_position = {"topic": m.topic(),
"partition": m.partition(),
"offset": m.offset()}
audit_write(log, {
"stage": "TRANSPORT_CONFIRMED", "run_id": a.run_id,
"event_id": item["original_id"],
"dlq_id": item["dlq_id"], **broker_position})
except Exception as exc:
# Conservatively stop. Some errors are definitive, others
# mean the publish succeeded but its response was lost.
audit_write(log, {
"stage": "UNKNOWN_OR_FAILED_STOP", "run_id": a.run_id,
"event_id": item["original_id"], "dlq_id": item["dlq_id"],
"error_type": type(exc).__name__})
raise
# Sequential confirms + spacing; no burst after a long stall.
next_due = time.monotonic() + 1.0 / a.rate
finally:
if connection is not None and connection.is_open:
try:
with deadline(5):
connection.close()
except Exception as exc:
print("Cleanup warning: " + type(exc).__name__,
file=sys.stderr)
return 0
if __name__ == "__main__":
try:
sys.exit(main())
except Exception as exc:
print(type(exc).__name__ + ": " + str(exc), file=sys.stderr)
sys.exit(2)Để kiểm thử schema/replay, thêm một event hợp lệ theo schema v2 nhưng consumer v1 chưa hỗ trợ; đưa vào DLQ với SCHEMA_UNKNOWN. Sau khi deploy decoder v2 tương thích, approve đúng subset và replay nguyên bytes. M1 malformed không được sửa bytes rồi giữ cùng ID/fingerprint; khi cần sửa nội dung nghiệp vụ, tạo correction event theo policy riêng và link tới bản gốc.
Tool dùng POSIX watchdog 15 giây để chặn thời gian chờ RabbitMQ confirm (Windows chạy qua WSL), và dừng khi publish lỗi hoặc outcome không rõ. Một PUBLISH_PLANNED không có TRANSPORT_CONFIRMED sau crash phải được reconcile, không được mặc định “chưa gửi”. Preserve attempts để app không tự cấp budget mới; việc cấp thêm budget cho replay cần approval/ledger riêng. Tool không triển khai quyền nhiều người duyệt, distributed rate limit, resume ledger hay distributed locking; chạy một process duy nhất cho manifest thí nghiệm. [R12], [R15]
Transport evidence
Strategy matrix có actual attempts, first/last failure và tuổi event tốt. DLQ giữ original ID/schema/reason/attempts. Audit ghi manifest hash, filter, owner/approval, publish planned, confirm hoặc unknown. So tốc độ phát thực tế với rate limit; chứng minh dry-run không làm broker counter tăng.
Business correctness
M1/P1 không tạo effect; T1 tạo đúng một effect sau lần thành công; G1/G2 không mất. Sau hỗ trợ schema v2, event đã duyệt tạo đúng một effect; replay manifest lần hai không đổi count/amount. Unapproved subset phải còn nguyên và không được publish.
Acceptance gate B
- So đủ immediate, delayed và blocking bằng cùng fixture; có giới hạn chạy negative control, không retry vô hạn.
- Budget được giữ qua restart; malformed/permanent được phân loại; lỗi parse không bỏ mất raw payload. DLQ metadata đủ reason, original ID, schema, attempts.
- Tiêm lỗi sau handoff confirm/trước source settle và khi target DLQ unavailable: không mất trách nhiệm xử lý; duplicate transport không tạo duplicate effect.
- Dry-run có selected manifest/hash nhưng không publish. Execute có filter, rate, approval và audit; replay lần hai vẫn idempotent; business resolution dựa vào assertion riêng.
Artifact B — 6 tệp: B/strategy-matrix.csv, B/dlq-manifest.jsonl, B/replay-plan.jsonl, B/replay-audit.jsonl, B/assertions.txt, B/replay.py. Lưu các lượt chạy phụ dưới tên có run ID để không ghi đè audit.
Lab C · Ordering và rebalance
Goal: không nhầm append order, completion order và business version; không commit offset của công việc chưa hoàn thành khi partition đổi chủ.
- Publish aggregate events cùng key và khác key; lưu partition/offset/version.
- Scale consumer group khi message đang xử lý.
- Gửi v3 rồi v2; state-transition guard bỏ stale event.
- Đo rebalance duration, consumer lag và processing latency.
C1. Failure-first: cố tình phá completion order
Dùng topic ba partitions từ setup, một consumer, một executor chạy xử lý song song không phân chia theo partition. Làm record đầu của một partition chậm hơn record sau; bản lỗi commit offset lớn nhất vừa xong và update state theo “message tới sau cùng”. Lưu event_id, key, partition, offset, version, start, finish, member_id. Oracle phải chỉ ra record chưa xử lý bị vượt qua hoặc state đi lùi.
Sau khi chụp bản lỗi, sửa thành xử lý tuần tự theo partition/key và commit theo completion watermark; dùng group mới lab-c-fixed với namespace DB mới để không mang offset/effect lỗi sang bản sửa. Việc đổi group chỉ để phân tách hai phiên bản thí nghiệm; khi crash/restart bản sửa phải giữ nguyên group.
bootstrap.servers=localhost:19092,localhost:29092,localhost:39092
group.id=lab-c-fixed
group.protocol=classic
enable.auto.commit=false
auto.offset.reset=earliest
partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignor
max.poll.records=1
max.poll.interval.ms=10000
session.timeout.ms=10000
heartbeat.interval.ms=3000
key.deserializer=org.apache.kafka.common.serialization.StringDeserializer
value.deserializer=org.apache.kafka.common.serialization.StringDeserializer
# 10s poll interval is an experiment parameter, not a production default.
# Give every process a distinct client.id; keep the group.id the same.# Publish snapshots in THIS order: A v1 -> A v3 -> A v2, plus another key.
# Pipe separator avoids conflicts with JSON colons.
printf '%s\n' \
'order-A|{"event_id":"aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaa1","aggregate_id":"order-A","version":1,"schema_version":1,"state":"NEW"}' \
'order-B|{"event_id":"bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbb1","aggregate_id":"order-B","version":1,"schema_version":1,"state":"NEW"}' \
'order-A|{"event_id":"aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaa3","aggregate_id":"order-A","version":3,"schema_version":1,"state":"PAID"}' \
'order-A|{"event_id":"aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaa2","aggregate_id":"order-A","version":2,"schema_version":1,"state":"RESERVED"}' |
kcli kafka-console-producer.sh --bootstrap-server kafka1:9092 \
--topic lab.c.events --property parse.key=true --property 'key.separator=|' \
--producer-property acks=all --producer-property enable.idempotence=true
# Capture actual group/offset state. The application consumer must be running.
kcli kafka-consumer-groups.sh --bootstrap-server kafka1:9092 \
--describe --group lab-c-fixed > evidence/C/offsets-before.txt
kcli kafka-consumer-groups.sh --bootstrap-server kafka1:9092 \
--describe --group lab-c-fixed --members --verbose \
> evidence/C/members-before.txt
kcli kafka-consumer-groups.sh --bootstrap-server kafka1:9092 \
--describe --group lab-c-fixed --state > evidence/C/group-before.txtVới partitioner và partition count cố định, key giúp duy trì phạm vi phân vùng; khác key có thể vẫn cùng partition. Kafka offset phản ánh vị trí trong partition, không phải global sequence. V3 được publish trước v2 trong fixture trên nên offset v3 nhỏ hơn v2 là kết quả đúng của log, không phải broker tự đảo message. [R5], [R7]
C2. Scale khi record chưa hoàn tất
Tạo thêm một batch có seed và ID riêng đủ để còn record in-flight. Đặt barrier ngay trước DB commit của một event được chọn. Tăng từ một lên hai rồi ba consumer cùng group; bỏ barrier để record hoàn tất. Lặp lại với scale-down và kill consumer đang sở hữu partition. Chụp assignments và committed offsets trước, trong và sau mỗi lần.
Log callback onPartitionsRevoked, onPartitionsAssigned, onPartitionsLost; mỗi task mang assignment epoch/ownership token. Tách một lần “handler đứng quá max.poll.interval.ms” để quan sát group behavior, rồi sửa bằng bounded processing và tiếp tục poll thay vì chỉ tăng timeout. Cooperative rebalance có thể có nhiều vòng chuyển giao; không lấy một callback đơn lẻ làm toàn bộ thời gian rebalance.
# Integration sketch for an asynchronous Java consumer.
# KafkaConsumer calls remain on ONE poll/owner thread.
on records polled:
attach current assignment epoch to each record
enqueue in bounded FIFO for its partition
never let later work bypass an unfinished record for the same key
on task completion:
validate result's assignment epoch
record DB outcome (APPLIED / DUPLICATE / durable quarantine)
advance only the processed-prefix watermark of delivered records
commit {partition -> next offset after that safe prefix}
on revoke:
stop scheduling revoked partitions
drain only within a bounded deadline
commit only the completed safe prefix while ownership is still valid
abandon/replay unfinished work; do not invent a "completed" offset
on lost:
do not commit for lost partitions
invalidate old assignment epoch
allow redelivery; DB Inbox and version guard protect late effects| Fault | Bằng chứng phải thấy | Không được làm |
|---|---|---|
| Offset 10 chậm; 11 hoàn tất trước | Committed next offset chưa vượt phần việc ở 10. Nếu safe prefix kết thúc ở 9, next offset vẫn là 10. | Commit 12 chỉ vì offset 11 đã xong. Offset commit không xác nhận riêng lẻ một record. |
| DB xong, partition bị revoke trước commit | New owner có thể nhận lại event; Inbox trả duplicate; state không đổi sai. | Commit bằng old owner sau mất ownership hoặc coi CommitFailedException là business rollback. |
| Worker cũ hoàn tất sau khi partition bị giao lại | Hai worker có thể tạo DB contention nhưng effect vẫn idempotent và version không lùi. | Giả định group membership tự fence mọi side effect trong DB. |
| Poll bị chặn quá lâu | Membership/assignment thay đổi hoặc client báo vi phạm poll interval theo protocol đã chọn; đo thực tế. | Sleep toàn poll thread để chờ retry rồi gọi đó là “partition pause”. |
| Thêm partitions giữa test | Có thể thay ánh xạ key với partitioner; dữ liệu cũ vẫn ở partitions cũ. | Thêm partition để tăng throughput rồi vẫn khẳng định thứ tự per-key xuyên migration chưa kiểm. |
Watermark là prefix đã hoàn tất của các record được giao theo partition; không giả định mọi số offset liên tiếp đều có một business record nhìn thấy được (ví dụ filtering/transactional records). KafkaConsumer không thread-safe; xử lý worker phải trả kết quả về thread sở hữu consumer để commit/pause/resume. [R7]
C3. V3 rồi v2: state-transition guard đúng loại event
-- Execute inside the SAME transaction as this projection's Inbox insert.
-- Parameters below demonstrate a snapshot update, not a delta application.
INSERT INTO projection_state
(projection_name,aggregate_id,version,state)
VALUES ('order-view','order-A',3,'{"state":"PAID"}'::jsonb)
ON CONFLICT (projection_name,aggregate_id) DO UPDATE
SET version=EXCLUDED.version,
state=EXCLUDED.state,
updated_at=clock_timestamp()
WHERE projection_state.version < EXCLUDED.version
RETURNING aggregate_id,version,state;
-- Now run the same statement with version=2 and state={"state":"RESERVED"}.
-- It must return zero updated rows; final state must remain version=3.
-- If zero rows: read current row and distinguish:
-- lower version -> STALE
-- equal version + identical content -> DUPLICATE
-- equal version + different content -> CONFLICT, not a silent success
SELECT aggregate_id,version,state
FROM projection_state
WHERE projection_name='order-view' AND aggregate_id='order-A';Snapshot: v3 chứa toàn bộ trạng thái cần thiết, nên v2 đến sau là stale và có thể bỏ qua có audit. Guard version nằm trong cùng transaction với Inbox; update kiểu read-then-write không khóa có thể race. Equal version nhưng payload khác phải bị coi là conflict.
Delta: nếu v2 là “+10” và v3 là “−5”, bỏ v2 chỉ vì đã thấy v3 sẽ làm sai kết quả. Chỉ áp dụng khi incoming_version = current_version + 1; version lớn hơn tạo GAP, cần buffer/quarantine bền vững và phục hồi sequence. Đừng đánh dấu Inbox “đã áp dụng” cho gap chưa xử lý; ack chỉ khi đã hoàn tất handoff bền vững. Chạy thêm fixture delta để chứng minh lý do hai policy khác nhau.
C4. Đo ba lớp thời gian, không chỉ nhìn lag
| Metric | Định nghĩa để báo cáo | Lưu ý |
|---|---|---|
| Rebalance duration | Ghi riêng: thời gian chuyển ownership quan sát ở client (revoke → assignment ổn định); và recovery interruption (fault/membership change → business completion đầu tiên ở owner mới). | Có timestamp/callback cho từng process; include scale command time. Không gộp hai đại lượng thành một số mà không nêu cách đo. |
| Consumer lag | Per-partition log-end offset trừ committed next offset tại thời điểm snapshot; ghi cả paused partitions. | Lag nhỏ không chứng minh DB đúng; commit quá sớm hoặc chuyển hết sang retry topic có thể che pending business work. |
| Processing latency | Thời gian trong một process từ handler start tới DB commit hoặc durable terminal decision; đo bằng monotonic clock. | Tách queue wait, DB time, commit-offset time và retry wait; không chỉ đo duration của poll(). |
| End-to-end age | Thời điểm business completion trừ occurred_at của event, theo đồng hồ được đồng bộ/ghi sai lệch. | Báo p50/p95/p99, max, số samples, lỗi và timeouts. Với sample nhỏ, không diễn giải percentile như production tail latency. |
Transport evidence
Ordering CSV chứa key/partition/offset/version/start/finish/assignment epoch. Group snapshots trước và sau scaling; callback logs và commit failures. Chứng minh không commit qua record chưa hoàn tất ở cả normal và revoke paths.
Business correctness
Snapshot order-A kết thúc ở v3/PAID dù v2 được nhận sau; equal-version conflict được chặn. Với delta có gap, trạng thái không bị cập nhật thiếu bước. Inbox bảo vệ duplicate khi old/new owner chồng lấn; fixture hợp lệ không thiếu effect.
Acceptance gate C
- Có cùng-key và khác-key fixtures, actual partition/offset/version, cùng bằng chứng completion order của bản lỗi.
- Scale up/down và kill khi record in-flight không làm bản sửa commit qua unfinished work. Task của owner cũ không tạo thêm effect hoặc làm version lùi.
- V3 rồi v2 đi qua guard và ra đúng state; báo rõ snapshot/delta contract, stale/duplicate/conflict/gap, không chỉ check số message.
- Có rebalance duration, per-partition lag và processing latency với định nghĩa, sample size và môi trường đo. Không khẳng định một ngưỡng cố định đúng cho mọi máy.
Artifact C — 5 tệp: C/ordering.csv, C/rebalance.jsonl, C/offset-snapshots.txt, C/latency.csv, C/assertions.txt. Ghi ordering scope, partitioner, partition count, group protocol và assignor trong kết luận.
Lab D · Outbox relay
Goal: business commit luôn để lại một nhiệm vụ publish bền vững; relay có thể chết, bị thay thế hoặc gặp broker suy giảm mà không làm mất event hay nhân đôi business effect.
- Business update + Outbox trong một transaction.
- Relay publish rồi crash trước mark-sent; chứng minh duplicate.
- Claim batch có lease; inject worker chết để worker khác reclaim.
- Monitor oldest unpublished age và cleanup retention.
D1. Failure-first: dual-write không có atomic boundary
Bản lỗi cập nhật business_order, commit DB rồi publish trực tiếp. Crash sau commit, trước publish: state nguồn đã đổi nhưng consumer không bao giờ thấy event và không có pending task để retry. Ghi cả source state và missing event. Sau đó thay bằng business update + Outbox insert trong một transaction, giữ nguyên fixture contract.
Chạy thêm crash trước commit: cả source update và Outbox insert phải rollback. Sau commit nhưng trước relay: source và Outbox đều tồn tại, sent_at còn null. Chỉ đọc Outbox đã commit; không gọi broker trong local business transaction để “thử tạo atomicity”.
-- Fresh fixture per experiment. This event ID remains stable during relay retry.
BEGIN;
INSERT INTO business_order(order_id,state,version)
VALUES ('order-D-001','NEW',0)
ON CONFLICT (order_id) DO NOTHING;
WITH changed AS (
UPDATE business_order SET state='PAID',version=1
WHERE order_id='order-D-001' AND version=0
RETURNING order_id,version
)
INSERT INTO outbox(event_id,origin_run_id,aggregate_id,payload)
SELECT
'dddddddd-dddd-4ddd-8ddd-ddddddddddd1'::uuid,
'D-after-publish', order_id,
jsonb_build_object(
'event_id','dddddddd-dddd-4ddd-8ddd-ddddddddddd1',
'origin_run_id','D-after-publish',
'business_key','tenant-demo:paid:order-D-001',
'aggregate_id',order_id, 'version',version,
'schema_version',1, 'amount',10,
'occurred_at',clock_timestamp())
FROM changed
RETURNING event_id,payload;
-- Application hook BEFORE_DB_COMMIT goes before this actual COMMIT.
-- Require exactly one outbox row for a new fixture; zero means no new transition.
COMMIT;D2. Publish rồi crash trước mark-sent
# Integration sketch: two independent relay workers.
claim batch and COMMIT claim transaction
for each claimed event:
renew/check lease if needed
publish original payload with the SAME event ID and aggregate key
RabbitMQ: persistent + mandatory + wait for confirm, inspect return/nack
Kafka: acks=all + idempotent producer; wait for send result
log TRANSPORT_CONFIRMED with event ID and returned metadata
Failpoint.hit("AFTER_PUBLISH_BEFORE_MARK_SENT", eventId)
mark sent with worker + epoch + still-valid lease
if mark count == 0: record LOST_CLAIM; reconcile, do not force updateTiêm AFTER_PUBLISH_BEFORE_MARK_SENT cho đúng một event. Sau khi lease hết hạn, worker khác publish lại cùng event. Lưu hai lần delivery/record positions cùng ID và chứng minh consumer có một effect. Idempotent Kafka producer có thể ngăn một số retry trong phiên producer, nhưng một publish chủ động mới sau relay restart không phải lời hứa dedup theo business event ID.
Khi timeout/connection mất trước confirm, đừng đánh dấu sent chỉ vì lời gọi publish đã được gửi. Outcome có thể chưa biết: giữ event pending/reconcile. Với RabbitMQ, một publish unroutable phải được coi là handoff thất bại kể cả có publisher confirm; kiểm mandatory return.
[R1], [R5]D3. Lease, reclaim và fencing token
Chạy hai workers với owner khác nhau. Claim bằng FOR UPDATE SKIP LOCKED, tăng lease_epoch mỗi lần claim; commit ngay rồi mới gọi broker. Lease 30 giây là tham số minh họa. Batch size, publish timeout và lịch renew phải bảo đảm event ở cuối batch không bị lease hết hạn trước khi bắt đầu publish.
Mở SQL claim, mark-sent và lease renewal
-- psql demo variables. Application workers use parameterized statements.
\set worker 'relay-1'
\set batch 10
BEGIN;
WITH picked AS (
SELECT event_id FROM outbox
WHERE sent_at IS NULL
AND (lease_until IS NULL OR lease_until < clock_timestamp())
ORDER BY created_at,event_id
LIMIT :batch
FOR UPDATE SKIP LOCKED
)
UPDATE outbox o
SET lease_owner=:'worker',
lease_until=clock_timestamp() + interval '30 seconds',
lease_epoch=o.lease_epoch+1,
attempts=o.attempts+1
FROM picked p
WHERE o.event_id=p.event_id
RETURNING o.event_id,o.payload,o.lease_owner,o.lease_epoch,o.lease_until;
COMMIT;
-- Publish outside the DB transaction, not while holding row locks.-- Values must come from THIS worker's returned claim, not a fresh SELECT.
\set worker 'relay-1'
\set event 'dddddddd-dddd-4ddd-8ddd-ddddddddddd1'
\set epoch 1
UPDATE outbox
SET sent_at=clock_timestamp(),lease_owner=NULL,lease_until=NULL
WHERE event_id=:'event'::uuid
AND sent_at IS NULL
AND lease_owner=:'worker'
AND lease_epoch=:epoch
AND lease_until > clock_timestamp()
RETURNING event_id,sent_at;
-- Exactly 1 row = successful mark; 0 = expired/lost/already marked claim.
-- A publish may still have succeeded when this UPDATE returns zero.-- Optional heartbeat for a bounded, still-owned claim.
UPDATE outbox
SET lease_until=clock_timestamp() + interval '30 seconds'
WHERE event_id=:'event'::uuid AND sent_at IS NULL
AND lease_owner=:'worker' AND lease_epoch=:epoch
AND lease_until > clock_timestamp()
RETURNING lease_until;
-- If zero rows, STOP using the old claim. Do not silently adopt a new epoch.SKIP LOCKED phù hợp cho kiểu nhiều worker lấy việc từ bảng queue, nhưng không tạo một ảnh chụp nhất quán để báo cáo nghiệp vụ. Token tăng đơn điệu ở đây chỉ fence việc cập nhật Outbox; nó không thu hồi một request publish đã ra mạng và không ngăn được mọi duplicate ở broker. [R13]
| Case | Cách tiêm lỗi | Expected observations / gate |
|---|---|---|
| D-claim-death | Worker 1 chết ngay sau claim commit, trước publish. | Hàng còn unsent, lease_owner=worker1. Trước expiry worker2 không claim; sau expiry worker2 có epoch lớn hơn, publish rồi mark-sent. Không pending bị kẹt vô hạn. |
| D-publish-death | Worker 1 chết sau confirm nhưng trước mark-sent. | Sau reclaim có duplicate transport cùng event ID. Downstream Inbox + unique business key giữ đúng một effect. |
| D-stale-worker | Tạm dừng worker1 sau claim; để lease expire và worker2 reclaim; cho worker1 tiếp tục với token cũ. | UPDATE mark-sent của worker1 trả 0 hàng. Một publish muộn vẫn có thể xảy ra; không được tuyên bố fencing DB làm broker exactly-once. |
| D-slow-batch | Publish chậm hơn ngân sách lease của batch. | Có renew hợp lệ hoặc claim được bỏ/reclaim. Không tiếp tục dùng epoch cũ sau khi mất lease; có metrics expired claims. |
| D-route/permission | Xóa binding trong sandbox hoặc dùng principal không có write permission. | Return/authorization error được phân loại, không mark-sent; pending/oldest age tăng có cảnh báo. Không retry auth failure như transient vô hạn. |
| D-order | Hai relay workers publish events của cùng aggregate từ các batch khác nhau. | Không mặc định created_at order = delivery order. Dùng partition/key, serial per-aggregate relay hoặc version/gap guard của C theo yêu cầu nghiệp vụ. |
D4. Broker degradation và capacity — giữ correctness dưới tải
Chạy workload steady, burst, một node chết và recovery catch-up với cùng payload/producer config. Trong lúc fault, relay có thể không publish được; source transaction vẫn tạo Outbox cho tới ngưỡng backpressure đã chọn. Theo dõi dung lượng DB/broker và dừng ingress ở ngưỡng an toàn thay vì để host đầy disk.
# Run only in the isolated Compose sandbox. Keep producers/consumers observed.
# RabbitMQ: verify membership is 3 before killing one member.
docker compose exec -T rabbit1 rabbitmq-queues quorum_status --vhost exec83a lab.d.main
docker compose kill -s SIGKILL rabbit3
# Measure confirms, ready/unacked, pending outbox, age, client reconnect.
# Restore and wait until membership is healthy before the next case:
docker compose start rabbit3
docker compose exec -T rabbit1 rabbitmq-queues quorum_status --vhost exec83a lab.d.main
# RabbitMQ majority loss: with all 3 healthy first, stop two.
docker compose kill -s SIGKILL rabbit2 rabbit3
# Keep rabbit1/producer alive; capture lack of successful durable progress.
docker compose start rabbit2 rabbit3
# Wait for quorum recovery, then prove backlog drains without duplicate effects.
# Kafka: one broker loss at the baseline RF=3 / minISR=2.
kcli kafka-topics.sh --bootstrap-server kafka1:9092 --describe --topic lab.d.events
docker compose kill -s SIGKILL kafka3
kcli kafka-topics.sh --bootstrap-server kafka1:9092 --describe --topic lab.d.events
# Publish new fixtures with acks=all and record actual result/latency.
docker compose start kafka3
# Wait for ISR=3, not just "container is running".
# Isolate insufficient-ISR behavior without losing the controller majority.
# Temporary TEST PROFILE: raise the topic requirement from 2 to 3.
kcli kafka-configs.sh --bootstrap-server kafka1:9092 \
--entity-type topics --entity-name lab.d.events --alter \
--add-config min.insync.replicas=3
docker compose kill -s SIGKILL kafka3
# Now ISR=2 < required 3 while 2/3 controllers remain available.
# Record producer errors/timeouts and prove relay does not mark sent.
docker compose start kafka3
# Wait for full ISR recovery, then RESTORE the documented baseline:
kcli kafka-configs.sh --bootstrap-server kafka1:9092 \
--entity-type topics --entity-name lab.d.events --alter \
--add-config min.insync.replicas=2
kcli kafka-configs.sh --bootstrap-server kafka1:9092 \
--entity-type topics --entity-name lab.d.events --describeRabbitMQ quorum queue ba members cần majority để tiếp tục hoạt động; mất hai members không phải chỉ “chậm hơn”. Với Kafka, phải kiểm actual ISR và yêu cầu min.insync.replicas cùng acks=all. Ví dụ dùng ba combined broker/controller: kill hai node còn làm mất controller quorum, nên không được quy mọi lỗi chỉ cho ISR. Phép thử tạm nâng minISR lên 3 và kill một node phía trên giúp tách điều kiện ISR mà vẫn giữ controller majority. [R2], [R5], [R14]
| Profile | Biến thay đổi có kiểm soát | Bằng chứng cần so |
|---|---|---|
| Healthy baseline | 3/3 nodes, workload/limits cố định; đo warm-up riêng. | Publish-confirm latency, useful business completions/s, error rate, pending count/age, DB pool, disk/network/CPU. |
| One-node loss | Giữ ingress tương đương; không hạ durability để làm đẹp kết quả. | Actual quorum/ISR, election/reconnect interval, backlog slope và latency tails. Không cam kết outage hay throughput bằng một con số cố định. |
| No quorum / insufficient ISR | Mất RabbitMQ majority hoặc ISR thấp hơn yêu cầu test profile; source/relay vẫn được quan sát. | Không mark sent khi chưa có bằng chứng broker; Outbox backlog tăng có giới hạn. Timeouts không được biến thành success. |
| Recovery + replay | Restore nodes; đợi replicas catch up; thêm replay có rate limit được ghi lại. | Thời gian drain, oldest age, duplicate deliveries, DB idempotency và ảnh hưởng workload mới. Replay không được chiếm hết downstream capacity. |
| Capacity sweep | Tăng arrival rate theo các bậc do người học chọn; giữ TLS/auth/payload/client config giống nhau. | Xác định vùng backlog tăng không hồi phục hoặc SLO vi phạm. Điểm gãy là kết quả của môi trường, không phải năng lực phổ quát của broker. |
profile,run_id,seed,payload_bytes,arrival_per_s,useful_complete_per_s,publish_p50_ms,publish_p95_ms,publish_p99_ms,publish_errors,pending_peak,oldest_age_peak_s,recovery_s,duplicate_deliveries,duplicate_effects,missing_effects,cpu_limit,ram_limit,tls_auth,notesƯớc lượng drain time có điều kiện: backlog / (μ − λ), trong đó μ là useful completion rate đo được trong recovery, λ là arrival rate mới và μ > λ. Nếu μ ≤ λ thì backlog không thể drain trong mô hình steady đó. Đây là phép ước lượng, không tính sẵn election, skew, retry hay downstream contention; đối chiếu với đường cong age thực tế.
Không so RabbitMQ với Kafka bằng các config durability khác nhau rồi kết luận một broker “nhanh hơn”. Ghi rõ plaintext sandbox hiện tại; một capacity claim cho production phải đo lại với TLS/auth, replication, quotas, payload distribution và giới hạn tài nguyên tương ứng.
D5. Oldest unpublished age và cleanup retention
SELECT count(*) AS pending,
count(*) FILTER (WHERE lease_until > clock_timestamp()) AS leased,
count(*) FILTER (WHERE lease_until <= clock_timestamp()) AS expired_leases,
coalesce(extract(epoch FROM
(clock_timestamp()-min(created_at))),0) AS oldest_unpublished_seconds
FROM outbox
WHERE sent_at IS NULL;
SELECT event_id,lease_owner,lease_epoch,lease_until,attempts,sent_at
FROM outbox
ORDER BY created_at,event_id;
-- The source transition and its outbox record must agree.
SELECT b.order_id,b.state,b.version,o.event_id,o.sent_at
FROM business_order b
LEFT JOIN outbox o ON o.aggregate_id=b.order_id
WHERE b.order_id='order-D-001';Sample pending, oldest_unpublished_seconds, publish latency/errors, attempts, expired leases và useful completion rate theo thời gian. Đặt cảnh báo theo age/SLO; count nhỏ vẫn có thể chứa một event bị kẹt rất lâu. Không loại hàng leased khỏi oldest-age query vì sẽ che worker treo.
Giữ Outbox đã gửi theo retention được phê duyệt và audit nhu cầu replay/điều tra. Inbox/business unique-key retention là quyết định riêng: tối thiểu phải bao phủ toàn bộ khoảng mà một event cũ còn có thể quay lại từ retry, broker retention, DLQ, backup hoặc replay. Xóa dedup record quá sớm có thể tái tạo effect dù relay hoàn toàn đúng.
Mở cleanup có preview và giới hạn batch
-- Review a retention cutoff before running. Example ONLY; choose your own.
\set cutoff '2026-09-01T00:00:00Z'
-- Preview: only already-sent rows can be eligible.
SELECT count(*) AS eligible FROM outbox
WHERE sent_at < :'cutoff'::timestamptz
AND lease_owner IS NULL;
SELECT count(*) AS pending_before FROM outbox WHERE sent_at IS NULL;
-- Bounded cleanup after archiving required evidence.
WITH eligible AS (
SELECT event_id FROM outbox
WHERE sent_at < :'cutoff'::timestamptz
AND lease_owner IS NULL
ORDER BY sent_at,event_id
LIMIT 1000
FOR UPDATE SKIP LOCKED
)
DELETE FROM outbox o USING eligible e
WHERE o.event_id=e.event_id
RETURNING o.event_id;
SELECT count(*) AS pending_after FROM outbox WHERE sent_at IS NULL;
-- pending_before == pending_after when relay/producer are paused for this test.
-- Inbox/unique business key retention is a separate policy; do not delete it here.Transport evidence
Timeline nối source commit → claim/epoch → publish confirm → crash → reclaim → publish lại → mark-sent. Có quorum/ISR snapshots, producer failure classification, pending/oldest-age curve và cleanup preview. Số outbox rows marked sent phải dựa trên confirm + valid claim.
Business correctness
Trước source commit: không state change và không Outbox. Sau commit: source state và Outbox cùng tồn tại. Relay publish lặp không làm business_effect tăng; mọi event hợp lệ cuối cùng có effect hoặc trạng thái pending/terminal được giải trình. Sau recovery và hết pending, oracle không có missing/duplicate.
Acceptance gate D
- Negative control dual-write tạo missing event; bản Outbox loại được cửa sổ thiếu task sau business commit, và rollback không để lại orphan event.
- Đã chứng minh duplicate sau publish-before-mark-sent trên cả RabbitMQ và Kafka, trong khi DB effect đúng một lần.
- Hai workers reclaim được lease của worker chết; late worker không mark-sent bằng token cũ; publish duplicate vẫn được xử lý idempotent.
- One-node loss, no-quorum/insufficient-ISR và recovery/capacity sweep có raw measurements, config, ngưỡng dừng và không mark-sent khi publish chưa xác nhận.
- Oldest unpublished age bao gồm leased rows; cleanup chỉ xóa sent rows đủ tuổi. Retention của dedup được giải thích theo replay horizon, không chỉ theo kích thước bảng.
Artifact D — 7 tệp: D/relay-timeline.jsonl, D/lease-snapshots.csv, D/outbox-age.csv, D/degradation-capacity.csv, D/cleanup.sql, D/assertions.txt, D/runbook.md. Runbook nêu pause ingress, khôi phục quorum/ISR, kiểm claim, drain backlog, replay giới hạn và xác nhận DB.
5. Deliverables và cách chấm bằng chứng
Giữ đủ bốn nhóm nguồn dưới đây. Bộ nộp tối thiểu được chuẩn hóa thành 26 tệp: 3 tệp chung + A (5) + B (6) + C (5) + D (7). Có thể bổ sung raw logs/snapshots theo từng run; không ghi đè lượt chạy cũ. Các tệp này là evidence người học tạo ra, không phải kết quả đã được chạy sẵn trong tài liệu.
| Deliverable nguồn — giữ nguyên | Evidence tối thiểu | Review gate |
|---|---|---|
| Broker/config/version và topology. | environment.md, configs.txt; config mong muốn + effective config, image digest, topology/ISR/quorum, app/client version. | Không có version/topology thực dùng → không thể tái hiện, chưa đạt. |
| Crash-boundary matrix với raw logs/offsets/acks. | A/crash-matrix.csv và transport logs; bổ sung B handoff, C offsets/rebalance, D claim/publish timeline. | Chỉ ảnh “queue empty” hay log “success” → chưa đủ. |
| DB assertions chứng minh một effect. | A/assertions.sql, A/B/C/D assertions.txt; expected ledger, duplicate/missing, amount và state/version. | Một duplicate effect hoặc missing effect chưa giải trình → fail correctness, dù throughput tốt. |
| DLQ/replay evidence và ordering scope. | B/dlq-manifest.jsonl, replay plan/audit và C/ordering.csv; schema evolution, original ID/key, ordering scope. | Replay không filter/audit hoặc gọi global ordering khi chỉ kiểm per-partition → chưa đạt. |
evidence/
environment.md
configs.txt
acceptance-summary.md
A/
crash-matrix.csv
transport.jsonl
assertions.sql
assertions.txt
fix.diff
B/
strategy-matrix.csv
dlq-manifest.jsonl
replay-plan.jsonl
replay-audit.jsonl
assertions.txt
replay.py
C/
ordering.csv
rebalance.jsonl
offset-snapshots.txt
latency.csv
assertions.txt
D/
relay-timeline.jsonl
lease-snapshots.csv
outbox-age.csv
degradation-capacity.csv
cleanup.sql
assertions.txt
runbook.mdacceptance-summary.md ghi từng gate là PASS / FAIL / NOT RUN / UNKNOWN, link về case/run/log/query tương ứng, nêu bản lỗi và bản sửa khác nhau ở đâu. Gate không được chạy thì không được gộp vào PASS. configs.txt tổng hợp phiên bản, image digest, topology, policy, client config và SHA ứng dụng; loại bỏ secret trước khi nộp.
Cleanup sandbox: xuất evidence, phục hồi broker/ISR/quorum và policy về baseline; dừng app/relay/replay rồi docker compose stop khi không dùng. Chỉ down/xóa volumes sau khi chắc chắn không còn cần dữ liệu lab. Không đưa lệnh phá hủy vào runbook production.
6. Tài liệu đối chiếu semantics
Các experiments và acceptance gate là thiết kế của trang này; các liên kết dưới đây dùng để đối chiếu protocol/API. Tài liệu version-pinned ưu tiên khớp baseline, không hàm ý đó là bản mới nhất. Spring và Confluent có thể hiển thị phiên bản hiện hành; ghi lại dependency thực dùng và kiểm tra API tương ứng. Cấu hình Docker cần đối chiếu listeners và runtime contract của Apache Kafka [R16].
- R1. RabbitMQ 4.1 — Consumer Acknowledgements and Publisher Confirms
- R2. RabbitMQ 4.1 — Quorum Queues
- R3. RabbitMQ 4.1 — Dead Letter Exchanges
- R4. RabbitMQ 4.1 — Time-To-Live and Expiration
- R5. Apache Kafka 4.0 — Producer Configs
- R6. Apache Kafka 4.0 — Consumer Configs
- R7. Apache Kafka 4.0 — KafkaConsumer API
- R8. PostgreSQL 17 — INSERT / ON CONFLICT
- R9. RabbitMQ 4.1 — Clustering Guide
- R10. RabbitMQ 4.1 — HTTP API Reference
- R11. Spring Framework — Using @Transactional
- R12. Confluent — Python client API / delivery callbacks / flush
- R13. PostgreSQL 17 — SELECT / FOR UPDATE / SKIP LOCKED
- R14. Apache Kafka 4.0 — Broker Configs / min.insync.replicas
- R15. Pika 1.3.2 — Blocking Connection Adapter
- R16. Apache Kafka 4.0.0 — Docker combined-cluster example
- R17. Java SE 21 — Runtime.halt