feat(log-store): add the enqueued acknowledgement mode to the object store WAL (#9358)

* feat(log-store): add the enqueued acknowledgement mode to the object store WAL

Add `ack_mode` (`durable` by default, or `enqueued`) and the backlog
thresholds `max_unpersisted_bytes` and `max_unpersisted_age` to the
object store WAL config, validated by the datanode and the store.

In the `enqueued` mode `append_batch` returns on admission with the
entry ids assigned and the object is created in the background. At a
backlog threshold the next append is held back until an upload
completes. A transient create failure is repeated under the same
sequence with the same bytes; any conflicting object poisons the
store. `stop` uploads the backlog, or returns the error that dropped
it once stop began. `obsolete` clamps the watermark to the durable
entry id, and an id the store handed out needs no sequence floor.

Add `LogStore::wait_durable` with a default that returns at once. The
object store WAL answers it once the region is durable and indexed
through the entry id, and fails it after a backlog was lost.

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* fix(log-store): poison an enqueued store on permanent create failures

Repeat a failed create in the enqueued mode only when the storage error
is retryable; any other storage error poisons the store. A conflicting
object poisons an enqueued store without reading its epoch, so a failed
header read cannot turn the conflict into a retry.

A durability wait for an id above the highest id the store handed out
now waits for the handed-out ids of the region below it instead of
returning at once.

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* fix(log-store): poison on permanent create failures after stop begins

A create that fails with a storage error that is not retryable poisons
an enqueued store even after stop began; only a transient failure drops
the backlog without poisoning. Durability waiters whose callers stopped
waiting are pruned before a new waiter is queued. The backlog age test
no longer depends on a follow-up append finishing within the threshold.

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* fix(log-store): answer a durability wait once no earlier entry is pending

A durability wait now returns once the region holds no entry at or
below the target that is handed out but not durable, instead of waiting
for the largest id handed out to the region. A later object of the
region that is still being created no longer holds back a wait whose
target it does not cover.

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* test(log-store): order the pending durability wait check after the actor

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* docs(store-api): state that the default wait_durable keeps each store's guarantee

The default `LogStore::wait_durable` returns at once, which keeps each
log store's own acknowledgement guarantee; Raft Engine with
`sync_write = false` acknowledges before its periodic sync, so the
documentation no longer claims that every entry id a caller holds is
durable. The object store WAL configuration test now also serializes
the new options and reads them back.

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* docs(log-store): limit the acknowledgement guarantees to the durable mode

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* fix(log-store): repeat enqueued creates that a retry layer marks persistent

An object store wrapped in the OpenDAL retry layer reports a temporary
error that outlasted its retries as persistent rather than temporary.
The enqueued mode now repeats a create after any storage error that is
not permanent, so a transient outage behind the retry layer no longer
poisons the store and drops the acknowledged backlog; after stop began
such a failure still drops the backlog without poisoning.

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* test(log-store): cover a persistent create failure after stop begins

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

---------

Signed-off-by: jeremyhi <fengjiachun@gmail.com>
This commit is contained in:
jeremyhi
2026-09-28 09:17:45 +00:00
committed by GitHub
parent a310ca2bcf
commit 03823a9a01
9 changed files with 1140 additions and 71 deletions
+8 -2
View File
@@ -125,7 +125,10 @@
| `wal.overwrite_entry_start_id` | Bool | `false` | Ignore missing entries during read WAL.<br/>**It's only used when the provider is `kafka`**.<br/><br/>This option ensures that when Kafka messages are deleted, the system<br/>can still successfully replay memtable data without throwing an<br/>out-of-range error.<br/>However, enabling this option might lead to unexpected data loss,<br/>as the system will skip over missing entries instead of treating<br/>them as critical errors. |
| `wal.storage_provider` | String | `""` | The name of the storage provider that holds the WAL objects, an empty name selects the default object store.<br/>**It's only used when the provider is `experimental_object_store`**. |
| `wal.prefix` | String | `wal` | The path prefix of the WAL objects inside the storage provider.<br/>The objects are written under `<prefix>/datanodes/<node_id>/epochs/<generation>`, which is derived from this prefix.<br/>**It's only used when the provider is `experimental_object_store`**. |
| `wal.flush_interval` | String | `100ms` | The interval of flushing buffered entries to the object store, at least `10ms`, defaults to `100ms`.<br/>Each non-empty timer-triggered flush creates one object, and a batch that reaches `max_batch_bytes` is flushed immediately, so under sustained load the batch seals on size and the interval no longer matters.<br/>When writes are sparse, a shorter interval lowers the acknowledgement latency of appends and raises the number of object requests: timer-triggered sealing creates at most one object per interval per node.<br/>**It's only used when the provider is `experimental_object_store`**. |
| `wal.flush_interval` | String | `100ms` | The interval of flushing buffered entries to the object store, at least `10ms`, defaults to `100ms`.<br/>Each non-empty timer-triggered flush creates one object, and a batch that reaches `max_batch_bytes` is flushed immediately, so under sustained load the batch seals on size and the interval no longer matters.<br/>When writes are sparse, a shorter interval lowers the acknowledgement latency of `durable` appends and raises the number of object requests: timer-triggered sealing creates at most one object per interval per node, and in `enqueued` mode the backlog thresholds can seal earlier. An `enqueued` append does not wait for its batch to seal once it is admitted, but the backlog thresholds can delay admission.<br/>**It's only used when the provider is `experimental_object_store`**. |
| `wal.ack_mode` | String | `durable` | When an append to the object store WAL returns.<br/>- `durable`: an append returns after the object holding its entries is durable (the default).<br/>- `enqueued`: an append returns once its entries are admitted and their ids are assigned, and the object is created in the background. A crash loses the unpersisted backlog.<br/>**It's only used when the provider is `experimental_object_store`**. |
| `wal.max_unpersisted_bytes` | String | `64MB` | The size of the unpersisted backlog at which new appends stall until an upload completes, in `enqueued` mode.<br/>Nothing is dropped and nothing is rejected; the threshold does not bound what a crash during an outage can lose.<br/>**It's only used when the provider is `experimental_object_store` and `ack_mode` is `enqueued`**. |
| `wal.max_unpersisted_age` | String | `8s` | The age of the oldest unpersisted entry at which new appends stall until an upload completes, in `enqueued` mode.<br/>**It's only used when the provider is `experimental_object_store` and `ack_mode` is `enqueued`**. |
| `wal.on_corrupted_segment` | String | `skip` | What a read does with a segment that still does not decode after a second fetch, because its checksum does not match or its content disagrees with its footer entry.<br/>- `skip`: the segment is skipped and recorded as a WAL hole of its region, a metric is incremented and a warning is logged; the other regions of the object are unaffected (the default).<br/>- `fail`: the read fails, so the region does not open.<br/>**It's only used when the provider is `experimental_object_store`**. |
| `metadata_store` | -- | -- | Metadata storage options. |
| `metadata_store.file_size` | String | `64MB` | The size of the metadata store log file. |
@@ -599,7 +602,10 @@
| `wal.overwrite_entry_start_id` | Bool | `false` | Ignore missing entries during read WAL.<br/>**It's only used when the provider is `kafka`**.<br/><br/>This option ensures that when Kafka messages are deleted, the system<br/>can still successfully replay memtable data without throwing an<br/>out-of-range error.<br/>However, enabling this option might lead to unexpected data loss,<br/>as the system will skip over missing entries instead of treating<br/>them as critical errors. |
| `wal.storage_provider` | String | `""` | The name of the storage provider that holds the WAL objects, an empty name selects the default object store.<br/>**It's only used when the provider is `experimental_object_store`**. |
| `wal.prefix` | String | `wal` | The path prefix of the WAL objects inside the storage provider.<br/>The objects are written under `<prefix>/datanodes/<node_id>/epochs/<generation>`, which is derived from this prefix.<br/>**It's only used when the provider is `experimental_object_store`**. |
| `wal.flush_interval` | String | `100ms` | The interval of flushing buffered entries to the object store, at least `10ms`, defaults to `100ms`.<br/>Each non-empty timer-triggered flush creates one object, and a batch that reaches `max_batch_bytes` is flushed immediately, so under sustained load the batch seals on size and the interval no longer matters.<br/>When writes are sparse, a shorter interval lowers the acknowledgement latency of appends and raises the number of object requests: timer-triggered sealing creates at most one object per interval per node.<br/>**It's only used when the provider is `experimental_object_store`**. |
| `wal.flush_interval` | String | `100ms` | The interval of flushing buffered entries to the object store, at least `10ms`, defaults to `100ms`.<br/>Each non-empty timer-triggered flush creates one object, and a batch that reaches `max_batch_bytes` is flushed immediately, so under sustained load the batch seals on size and the interval no longer matters.<br/>When writes are sparse, a shorter interval lowers the acknowledgement latency of `durable` appends and raises the number of object requests: timer-triggered sealing creates at most one object per interval per node, and in `enqueued` mode the backlog thresholds can seal earlier. An `enqueued` append does not wait for its batch to seal once it is admitted, but the backlog thresholds can delay admission.<br/>**It's only used when the provider is `experimental_object_store`**. |
| `wal.ack_mode` | String | `durable` | When an append to the object store WAL returns.<br/>- `durable`: an append returns after the object holding its entries is durable (the default).<br/>- `enqueued`: an append returns once its entries are admitted and their ids are assigned, and the object is created in the background. A crash loses the unpersisted backlog.<br/>**It's only used when the provider is `experimental_object_store`**. |
| `wal.max_unpersisted_bytes` | String | `64MB` | The size of the unpersisted backlog at which new appends stall until an upload completes, in `enqueued` mode.<br/>Nothing is dropped and nothing is rejected; the threshold does not bound what a crash during an outage can lose.<br/>**It's only used when the provider is `experimental_object_store` and `ack_mode` is `enqueued`**. |
| `wal.max_unpersisted_age` | String | `8s` | The age of the oldest unpersisted entry at which new appends stall until an upload completes, in `enqueued` mode.<br/>**It's only used when the provider is `experimental_object_store` and `ack_mode` is `enqueued`**. |
| `wal.on_corrupted_segment` | String | `skip` | What a read does with a segment that still does not decode after a second fetch, because its checksum does not match or its content disagrees with its footer entry.<br/>- `skip`: the segment is skipped and recorded as a WAL hole of its region, a metric is incremented and a warning is logged; the other regions of the object are unaffected (the default).<br/>- `fail`: the read fails, so the region does not open.<br/>**It's only used when the provider is `experimental_object_store`**. |
| `query` | -- | -- | The query engine options. |
| `query.parallelism` | Integer | `0` | Parallelism of the query engine.<br/>Default to 0, which means the number of CPU cores. |
+16 -1
View File
@@ -231,10 +231,25 @@ overwrite_entry_start_id = false
## The interval of flushing buffered entries to the object store, at least `10ms`, defaults to `100ms`.
## Each non-empty timer-triggered flush creates one object, and a batch that reaches `max_batch_bytes` is flushed immediately, so under sustained load the batch seals on size and the interval no longer matters.
## When writes are sparse, a shorter interval lowers the acknowledgement latency of appends and raises the number of object requests: timer-triggered sealing creates at most one object per interval per node.
## When writes are sparse, a shorter interval lowers the acknowledgement latency of `durable` appends and raises the number of object requests: timer-triggered sealing creates at most one object per interval per node, and in `enqueued` mode the backlog thresholds can seal earlier. An `enqueued` append does not wait for its batch to seal once it is admitted, but the backlog thresholds can delay admission.
## **It's only used when the provider is `experimental_object_store`**.
#+ flush_interval = "100ms"
## When an append to the object store WAL returns.
## - `durable`: an append returns after the object holding its entries is durable (the default).
## - `enqueued`: an append returns once its entries are admitted and their ids are assigned, and the object is created in the background. A crash loses the unpersisted backlog.
## **It's only used when the provider is `experimental_object_store`**.
#+ ack_mode = "durable"
## The size of the unpersisted backlog at which new appends stall until an upload completes, in `enqueued` mode.
## Nothing is dropped and nothing is rejected; the threshold does not bound what a crash during an outage can lose.
## **It's only used when the provider is `experimental_object_store` and `ack_mode` is `enqueued`**.
#+ max_unpersisted_bytes = "64MB"
## The age of the oldest unpersisted entry at which new appends stall until an upload completes, in `enqueued` mode.
## **It's only used when the provider is `experimental_object_store` and `ack_mode` is `enqueued`**.
#+ max_unpersisted_age = "8s"
## What a read does with a segment that still does not decode after a second fetch, because its checksum does not match or its content disagrees with its footer entry.
## - `skip`: the segment is skipped and recorded as a WAL hole of its region, a metric is incremented and a warning is logged; the other regions of the object are unaffected (the default).
## - `fail`: the read fails, so the region does not open.
+16 -1
View File
@@ -414,10 +414,25 @@ overwrite_entry_start_id = false
## The interval of flushing buffered entries to the object store, at least `10ms`, defaults to `100ms`.
## Each non-empty timer-triggered flush creates one object, and a batch that reaches `max_batch_bytes` is flushed immediately, so under sustained load the batch seals on size and the interval no longer matters.
## When writes are sparse, a shorter interval lowers the acknowledgement latency of appends and raises the number of object requests: timer-triggered sealing creates at most one object per interval per node.
## When writes are sparse, a shorter interval lowers the acknowledgement latency of `durable` appends and raises the number of object requests: timer-triggered sealing creates at most one object per interval per node, and in `enqueued` mode the backlog thresholds can seal earlier. An `enqueued` append does not wait for its batch to seal once it is admitted, but the backlog thresholds can delay admission.
## **It's only used when the provider is `experimental_object_store`**.
#+ flush_interval = "100ms"
## When an append to the object store WAL returns.
## - `durable`: an append returns after the object holding its entries is durable (the default).
## - `enqueued`: an append returns once its entries are admitted and their ids are assigned, and the object is created in the background. A crash loses the unpersisted backlog.
## **It's only used when the provider is `experimental_object_store`**.
#+ ack_mode = "durable"
## The size of the unpersisted backlog at which new appends stall until an upload completes, in `enqueued` mode.
## Nothing is dropped and nothing is rejected; the threshold does not bound what a crash during an outage can lose.
## **It's only used when the provider is `experimental_object_store` and `ack_mode` is `enqueued`**.
#+ max_unpersisted_bytes = "64MB"
## The age of the oldest unpersisted entry at which new appends stall until an upload completes, in `enqueued` mode.
## **It's only used when the provider is `experimental_object_store` and `ack_mode` is `enqueued`**.
#+ max_unpersisted_age = "8s"
## What a read does with a segment that still does not decode after a second fetch, because its checksum does not match or its content disagrees with its footer entry.
## - `skip`: the segment is skipped and recorded as a WAL hole of its region, a metric is incremented and a warning is logged; the other regions of the object are unaffected (the default).
## - `fail`: the read fails, so the region does not open.
+19 -1
View File
@@ -138,7 +138,7 @@ mod tests {
use super::*;
use crate::TopicSelectorType;
use crate::config::object_store::CorruptedSegmentAction;
use crate::config::object_store::{AckMode, CorruptedSegmentAction};
use crate::config::{DatanodeKafkaConfig, MetasrvKafkaConfig};
#[test]
@@ -294,6 +294,9 @@ mod tests {
assert_eq!(config.prefix, "wal");
assert_eq!(config.flush_interval, Duration::from_millis(100));
assert_eq!(config.max_batch_bytes, ReadableSize::mb(8));
assert_eq!(config.ack_mode, AckMode::Durable);
assert_eq!(config.max_unpersisted_bytes, ReadableSize::mb(64));
assert_eq!(config.max_unpersisted_age, Duration::from_secs(8));
assert_eq!(config.on_corrupted_segment, CorruptedSegmentAction::Skip);
assert!(matches!(
MetasrvWalConfig::try_from(datanode_wal_config).unwrap_err(),
@@ -306,6 +309,9 @@ mod tests {
prefix = "cluster-a/wal"
flush_interval = "500ms"
max_batch_bytes = "4MB"
ack_mode = "enqueued"
max_unpersisted_bytes = "32MB"
max_unpersisted_age = "4s"
on_corrupted_segment = "fail"
"#;
let datanode_wal_config: DatanodeWalConfig = toml::from_str(toml_str).unwrap();
@@ -314,12 +320,24 @@ mod tests {
prefix: "cluster-a/wal".to_string(),
flush_interval: Duration::from_millis(500),
max_batch_bytes: ReadableSize::mb(4),
ack_mode: AckMode::Enqueued,
max_unpersisted_bytes: ReadableSize::mb(32),
max_unpersisted_age: Duration::from_secs(4),
on_corrupted_segment: CorruptedSegmentAction::Fail,
};
assert_eq!(
datanode_wal_config,
DatanodeWalConfig::ObjectStore(expected)
);
let serialized = toml::to_string(&datanode_wal_config).unwrap();
let table: toml::Table = toml::from_str(&serialized).unwrap();
assert_eq!(table["ack_mode"].as_str(), Some("enqueued"));
assert_eq!(table["max_unpersisted_bytes"].as_str(), Some("32MiB"));
assert_eq!(table["max_unpersisted_age"].as_str(), Some("4s"));
assert_eq!(
datanode_wal_config,
toml::from_str::<DatanodeWalConfig>(&serialized).unwrap()
);
// The persisted tag is not accepted as a config provider.
let toml_str = r#"
+24
View File
@@ -24,6 +24,18 @@ use serde::{Deserialize, Serialize};
/// node id and the generation from the metasrv.
pub const STANDALONE_GENERATION: u64 = 0;
/// When an append to the object store WAL returns.
#[derive(Debug, Clone, Copy, Default, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "snake_case")]
pub enum AckMode {
/// An append returns after the object holding its entries is durable.
#[default]
Durable,
/// An append returns once its entries are admitted and their ids are
/// assigned; the object is created in the background.
Enqueued,
}
/// What a read does with a segment that still does not decode after a
/// second fetch, because its checksum does not match or its content
/// disagrees with its footer entry.
@@ -56,6 +68,15 @@ pub struct ObjectStoreWalConfig {
pub flush_interval: Duration,
/// The max size of a single batch object, defaults to 8MiB.
pub max_batch_bytes: ReadableSize,
/// When an append returns, defaults to `durable`.
pub ack_mode: AckMode,
/// Size of the unpersisted backlog at which appends stall until an
/// upload completes in `enqueued` mode, defaults to 64MiB.
pub max_unpersisted_bytes: ReadableSize,
/// Age of the oldest unpersisted entry at which appends stall until an
/// upload completes in `enqueued` mode, defaults to 8s.
#[serde(with = "humantime_serde")]
pub max_unpersisted_age: Duration,
/// What a read does with a segment that still does not decode after a
/// second fetch, because its checksum does not match or its content
/// disagrees with its footer entry, defaults to `skip`.
@@ -69,6 +90,9 @@ impl Default for ObjectStoreWalConfig {
prefix: "wal".to_string(),
flush_interval: Duration::from_millis(100),
max_batch_bytes: ReadableSize::mb(8),
ack_mode: AckMode::Durable,
max_unpersisted_bytes: ReadableSize::mb(64),
max_unpersisted_age: Duration::from_secs(8),
on_corrupted_segment: CorruptedSegmentAction::Skip,
}
}
+32
View File
@@ -780,6 +780,22 @@ fn validate_object_store_wal_config(config: &ObjectStoreWalConfig) -> Result<()>
reason: "must be greater than 0",
}
);
ensure!(
config.max_unpersisted_bytes.as_bytes() > 0,
InvalidObjectStoreWalConfigSnafu {
field: "max_unpersisted_bytes",
value: config.max_unpersisted_bytes.to_string(),
reason: "must be greater than 0",
}
);
ensure!(
config.max_unpersisted_age > Duration::ZERO,
InvalidObjectStoreWalConfigSnafu {
field: "max_unpersisted_age",
value: format!("{:?}", config.max_unpersisted_age),
reason: "must be greater than 0",
}
);
Ok(())
}
@@ -1163,12 +1179,28 @@ mod tests {
},
"max_batch_bytes",
),
(
ObjectStoreWalConfig {
max_unpersisted_bytes: ReadableSize(0),
..Default::default()
},
"max_unpersisted_bytes",
),
(
ObjectStoreWalConfig {
max_unpersisted_age: Duration::ZERO,
..Default::default()
},
"max_unpersisted_age",
),
];
for (config, expected_field) in cases {
let expected_value = match expected_field {
"prefix" => config.prefix.clone(),
"flush_interval" => format!("{:?}", config.flush_interval),
"max_unpersisted_bytes" => config.max_unpersisted_bytes.to_string(),
"max_unpersisted_age" => format!("{:?}", config.max_unpersisted_age),
_ => config.max_batch_bytes.to_string(),
};
let data_home = create_temp_dir("object-store-wal-invalid-config");
@@ -145,6 +145,10 @@ impl OpenBatch {
self.entries.is_empty()
}
pub(crate) fn holds_region(&self, region_id: RegionId) -> bool {
self.positions.contains_key(&region_id)
}
/// Returns the estimated size of the admitted entries.
pub(crate) fn estimated_bytes(&self) -> usize {
self.estimated_bytes
File diff suppressed because it is too large Load Diff
+18
View File
@@ -105,6 +105,24 @@ pub trait LogStore: Send + Sync + 'static + std::fmt::Debug {
/// Returns the latest entry id in the log store.
fn latest_entry_id(&self, provider: &Provider) -> Result<EntryId, Self::Error>;
/// Waits until every entry of the provider's region with an id at or below
/// `entry_id` has the persistence this log store guarantees for an
/// acknowledged append, so that a caller can persist `entry_id` as a
/// replay watermark.
///
/// The default returns at once and keeps each log store's own
/// acknowledgement guarantee, whatever it is: Raft Engine with
/// `sync_write = false`, for example, acknowledges an append before its
/// periodic sync. A log store that acknowledges an append before its
/// entries have that persistence overrides this method.
async fn wait_durable(
&self,
_provider: &Provider,
_entry_id: EntryId,
) -> Result<(), Self::Error> {
Ok(())
}
}
/// The response of an `append` operation.