# 청구 저장소의 javadoc
25: *
Claiming is {@code FOR UPDATE SKIP LOCKED} inside the select, for the same reason the delivery
26: * queue is: two workers polling the same table must not both take the same job.
# 청구 문장과 그것을 실행하는 메서드
44: SELECT id, attempt_id, provider_profile_id, next_check_at, attempts, last_result,
45: created_at, updated_at
46: FROM notification_reconciliation_job
47: WHERE next_check_at <= ?
48: ORDER BY next_check_at, id
49: LIMIT ?
50: FOR UPDATE SKIP LOCKED
51: """;
52:
113: public List claimDue(int limit, Instant now) {
114- Objects.requireNonNull(now, "now");
115- if (limit < 1) {
116- throw new IllegalArgumentException("limit");
117- }
118- return jdbc.query(CLAIM_DUE, MAPPER, Timestamp.from(now), limit);
119- }
120-
# 그 클래스의 @Transactional / TransactionPort: 0
# 같은 패키지의 전달 큐는 같은 잠금 구문에 UPDATE 를 붙인다
27: "WITH claimable AS ("
28: + " SELECT id FROM notification_recipient_delivery"
29: + " WHERE next_dispatch_at <= :now"
30: + " AND delivery_state IN ('PENDING', 'READY_TO_DISPATCH', 'RETRY_WAITING')"
31: + " AND (expires_at IS NULL OR expires_at > :now)"
32: + " AND (lease_until IS NULL OR lease_until < :now)"
33: + " ORDER BY next_dispatch_at, id"
34: + " FOR UPDATE SKIP LOCKED"
35: + " LIMIT :batchSize"
36: + ") "
37: + "UPDATE notification_recipient_delivery AS d "
38: + " SET lease_owner = :owner,"
39: + " lease_until = :leaseUntil,"
40: + " lease_fence = d.lease_fence + 1,"
41: + " delivery_state = 'DISPATCHING',"
42: // The JPA @Version column, advanced by the native write. Without this a managed entity
43: // loaded before the claim still holds the old version, so its flush succeeds and
44: // overwrites the lease and the state with pre-claim values — optimistic locking that
45: // reports no conflict because the native statement never told it there was one.
46: + " version = d.version + 1,"
47: + " updated_at = :now"
48: + " FROM claimable"
49: + " WHERE d.id = claimable.id"
50: + " RETURNING d.id, d.lease_fence";
# 작업자는 조회와 정산 사이에 제공자를 부른다
80: public int reconcileOnce() {
81- List due = jobs.claimDue(batchSize, clock.instant());
82- int settled = 0;
83- for (ReconciliationJob job : due) {
84- // One job's failure is not the pass's: a provider that is refusing connections would
85- // otherwise stop every other provider's jobs behind it.
86- try {
87- if (handle(job)) {
88- settled++;
89- }
90- } catch (RuntimeException failure) {
91- jobs.reschedule(
92- job, failure.getClass().getSimpleName(), clock.instant().plus(retryBackoff));
93- }
94- }
95- return settled;
96- }
# 그 클래스의 @Transactional / TransactionPort: 0
# 조정 작업 표의 열
68:CREATE TABLE notification_reconciliation_job (
69- id uuid PRIMARY KEY,
70- attempt_id uuid NOT NULL REFERENCES notification_delivery_attempt(id),
71- provider_profile_id varchar(120) NOT NULL,
72- next_check_at timestamptz NOT NULL,
73- attempts integer NOT NULL DEFAULT 0,
74- last_result varchar(40),
75- created_at timestamptz NOT NULL,
76- updated_at timestamptz NOT NULL,
77- CONSTRAINT uk_notification_reconciliation_attempt UNIQUE (attempt_id),
78- CONSTRAINT ck_notification_reconciliation_attempts CHECK (attempts >= 0)
79-);
80-
# 그 표 정의 안에서 owner/lease/status 를 포함한 줄: 0
# 이 표를 건드리는 마이그레이션 전부
V3__notification_platform_inbox_admin.sql:68:CREATE TABLE notification_reconciliation_job (
V3__notification_platform_inbox_admin.sql:82: ON notification_reconciliation_job (next_check_at, id);
V4__notification_platform_alignment_and_activation.sql:178: ('notification_reconciliation_job'),
# 도달 조건
486: enabled: ${APP_NOTIFICATION_PLATFORM_ENABLED:false}
489: mode: ${APP_NOTIFICATION_PLATFORM_MODE:SERVING}
61: Executors.newScheduledThreadPool(
62- 1,
63- runnable -> {
64- Thread thread = new Thread(runnable, "notification-background");
78: schedule("reconciliation", reconciliation::reconcileOnce);