diff --git a/native/csrc/catalog/lease_coordinator.cpp b/native/csrc/catalog/lease_coordinator.cpp index 05633ce4d..b54f49f56 100644 --- a/native/csrc/catalog/lease_coordinator.cpp +++ b/native/csrc/catalog/lease_coordinator.cpp @@ -119,6 +119,16 @@ PublisherLease LeaseCoordinator::claim_with_rival( std::set owners; for (const Row& row : rows) owners.insert(row[0]); if (owners == std::set{lease_id}) { + // rows[0] is this claim's own row, unordered and ungrouped though the + // read is: every row at this term carries our lease_id, and we wrote one + // (the client repeats an INSERT only when its connection was never made, + // so nothing reached the server, and CatalogWriter quarantines after one + // of unknown outcome rather than retrying). The expiry rule is + // the minimum expires_at_ns under (term, lease_id), which differs from + // any one row only once a release's tombstone shares the key, and none + // can here: a release writes at the holder's own term and a claim always + // goes to head + 1, so no claim lands on a term that already holds its + // own tombstone. lease_ = PublisherLease{ term, lease_id, holder, parse_u64_field(rows[0][1], "lease acquisition"), diff --git a/native/csrc/catalog/storage_service.cpp b/native/csrc/catalog/storage_service.cpp index 941a3a038..cc1254c96 100644 --- a/native/csrc/catalog/storage_service.cpp +++ b/native/csrc/catalog/storage_service.cpp @@ -632,8 +632,24 @@ void CaptureStorageService::keep_lease() { } void CaptureStorageService::renew_lease_if_due() { - // Renew once a third of the TTL has passed without a publish, which leaves - // two more tries before a rival could claim it. + // Renew once a third of the TTL has passed since last_renew_ns_. The lease + // thread wakes every ttl/6, so with no index() in the way the renewal + // fires within about a tick of falling due, leaving at least roughly half + // the TTL for it to land before the row expires. An index() can leave far + // less, or none. It holds lease_mutex_ throughout, so this thread cannot + // renew until it returns, and index_bounded() then stamps last_renew_ns_ + // whenever it indexed or skipped a pack, as though the row had just been + // renewed. After a publish that stamp trails the publish's own last + // renewal by the watermark INSERT, its read-backs and commit_packs; and + // it is taken even when every pack was already committed, so index() + // published nothing and renewed nothing. Whatever slack is left covers a + // renewal that runs late, not one that fails. A failed renewal costs the + // lease at once whatever the cause. A refusal drops it in the coordinator, + // and any ClickHouse error (transport, timeout, or a server error) takes + // renew_for_publish()'s catch, which quarantines the writer on the first + // error that survives the client's retries (a write is repeated only when + // its connection was never made; one that may have reached the server + // never is). const uint64_t ttl = config_.writer.lease_ttl_ns; if (ttl == 0 || steady_ns() - last_renew_ns_ < ttl / 3) return; writer_.renew_lease(); diff --git a/native/csrc/catalog/storage_service.h b/native/csrc/catalog/storage_service.h index 13b77dd87..3c798e04a 100644 --- a/native/csrc/catalog/storage_service.h +++ b/native/csrc/catalog/storage_service.h @@ -105,7 +105,21 @@ struct StorageServiceConfig { // spool. Recover() deletes every .open file this object does not own, so it // is only safe while no writer is live: start() must run before the sink // opens the spool. It runs after the lease is taken, so a start refused - // the catalog never touches the spool. + // the catalog never touches the spool. That keeps a second process off a + // live spool only usually: a holder that stops renewing for a TTL + // (quarantined, or stalled) lets its row lapse, and a second process can + // take the lease and sweep while the first is still writing; one on + // another (database, table_prefix) never meets the lease at all. Its + // Recover() also lists the first's sealed packs, which the first may + // upload too. A cycle checks the lease once, before its upload batch, and + // UploadPending does not stop when the lease is lost: a holder that is + // quarantined, or refused a renewal or publish, while a batch is in + // flight finishes that batch (which can outlast the TTL), and only its + // later cycles upload nothing while it holds no lease. A stalled one + // still holds the lease locally and starts new batches until a renewal + // or publish is refused. Two on different (database, table_prefix) pairs + // each hold a lease and upload freely. The spool itself is not locked; + // one process per spool is the caller's job. bool sweep_spool_on_start = true; bool reconcile_on_start = true; }; diff --git a/src/dmi/storage/native_capture.py b/src/dmi/storage/native_capture.py index d8e8022e1..aed635898 100644 --- a/src/dmi/storage/native_capture.py +++ b/src/dmi/storage/native_capture.py @@ -269,11 +269,15 @@ class NativeCaptureStorageConfig: clock_skew_s: float = 0.0 # How long start() waits for a predecessor's lease to expire before # failing with it held. None waits lease_ttl_s + publish_timeout_s + - # clock_skew_s, enough to outlast a crashed predecessor that ran with the - # same knobs; 0 fails at once. A predecessor with a longer TTL (the - # native default is 30 s, which processes predating these knobs used) - # can outlast it: set this explicitly for the first restart after such - # a process. + # clock_skew_s; 0 fails at once. That is guaranteed to outlast a crashed + # predecessor only when its TTL is at most lease_ttl_s + + # publish_timeout_s: clock_skew_s cancels, since the wait adds it and a + # lagging replica sees the row live that much longer -- which assumes + # clock_skew_s bounds the offset between the replica that stamped the + # predecessor's row and the one serving the read. Set this explicitly + # for the first restart after a predecessor above that threshold -- 20 s + # on these defaults, so one on the native 30 s default (which processes + # predating these knobs used) outlasts the default wait by 10 s. start_lease_wait_s: Optional[float] = None def __post_init__(self) -> None: