Compare commits

...
Author SHA1 Message Date
1df165d514 fix(volume): keep the TTL clock across a vacuum commit instead of rescanning (#11630)
* fix(volume): keep the TTL clock across a vacuum commit instead of rescanning

CommitCompact reloads the swapped files while holding dataFileAccessLock,
and for a vacuumed TTL volume that reload re-derived lastModifiedTsSeconds
by reading every live needle's append timestamp from the .dat: two random
reads per needle, with every read of the volume blocked behind them.

The in-memory clock is already current at that point. Every write since the
volume loaded moved it, and makeupDiff only replays writes that went through
that path. Carry it across the reload instead. This also stops an
over-budget scan from falling back to the new .dat's mtime and restarting
an expiring volume's TTL at the commit.

* volume: carry the append watermark as the TTL clock across a vacuum commit

The running append watermark is the clock the reload's recovery scan
recomputes, so the commit can keep it directly. Client-supplied needle
modified times can run ahead of or behind the append time; keeping
lastModifiedTsSeconds itself would let a forged or stale timestamp move
expiry through a vacuum, where the scan it replaces used server-side
append timestamps.

* storage: test that a vacuum commit keeps the append clock

A write's client supplied modified time can lie ahead of or behind its
append time; the commit must land the TTL clock on the append watermark,
the same value the recovery scan would have recomputed.

* volume: carry the append watermark as the TTL clock across a vacuum commit

Mirrors the Go volume server: the running append watermark is the clock
the reload's recovery scan recomputes, so the commit keeps it instead of
rescanning live needles under the write lock.

* volume: commit carries the last-write append time, not the latest append

lastAppendAtNs counts tombstone appends and is reseeded from the .dat
tail at every load, so it can sit ahead of the last write -- a delete
freshens the commit clock -- or behind it: a restarted vacuumed volume's
tail needle is not its newest write, and the commit would move the TTL
clock backward into premature expiry.

Track lastWriteAppendAtNs instead, bumped only on needle appends and
seeded by the recovery scan, so the commit lands the clock on the same
live-write maximum the rescan would have recomputed.

* volume: rescan at commit when the newest write was deleted

lastWriteAppendAtNs can hold a write the index no longer holds, so
carrying it extends the TTL clock past what recovery over the compacted
index would compute. Remember the key behind the watermark so its
tombstone or index rollback can send the reload back through
recoverLastModifiedTs, landing on the newest surviving write.

* volume: a tombstone retires the write rows beneath it in the last-write scan

A needle deleted after the compaction copy leaves its write row followed
by a tombstone in the committed index. The reverse scan skipped the
tombstone row but then counted the dead write, reseeding the watermark
and flag as if it were alive. Track keys whose latest row is a tombstone
so their earlier write rows stop counting, and exercise the
delete-inside-the-commit-window ordering in the tests.

---------

Co-authored-by: Chris Lu <chrislusf@users.noreply.github.com>
Co-authored-by: Chris Lu <chris.lu@gmail.com>
2026-10-08 22:58:44 +08:00
Chris LuGitHubDevin <158243242+devin-ai-integration[bot]@users.noreply.github.com>Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
0305e837fd iceberg maintenance: keep compacted files prunable (stats, bound order, row groups) (#11654)
* iceberg maintenance: record column statistics on compacted files

A compacted file's manifest entry was built with no column_sizes,
value_counts, null_value_counts, lower_bounds, upper_bounds or
split_offsets, so no reader could skip a compacted file on any
predicate. parquet-go already writes exact per-chunk min/max and null
counts into the footer; read that footer back after the merge and
record it on the data file, bounds as the spec's single-value
serialization with string and binary truncated as truncate(16). A
column whose bounds cannot be converted exactly gets none, and a
statistics failure is logged while the compaction commits anyway:
metrics are an optimization, not a correctness requirement.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* iceberg maintenance: merge bins in bound order so compacted files stay prunable

A bin's files were concatenated in manifest order, or largest-first
when a partition was split under the target size, so inputs disjoint
on a column came out as outputs that overlapped on it. Order each
bin's files by their bounds on one column before merging and split an
oversized partition into runs of consecutive files, so every output
covers one contiguous range. The column is the first identity field
of the table's sort order when it declares one (a descending order
sorts by upper bound), otherwise the first schema column every
candidate file has bounds for; detection resolves the same order so
it plans the bins execution builds. Files without bounds keep the old
behavior.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* iceberg maintenance: cap compacted files' row groups from table config

Neither merge writer set a row-group limit (parquet-go's default is
unlimited rows), so every compacted file was a single row group and
readers could not skip inside it either. Rows per row group now come
from the table's write.parquet.row-group-limit and
write.parquet.row-group-size-bytes, defaulting to PyIceberg's
1 048 576 rows and Iceberg's 128 MiB, with the byte size turned into
rows from the bin's inputs' compressed bytes per row and a floor of
1 024 rows. Both writers take the cap, and the statistics the entry
records list one split offset per row group.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* iceberg: keep compaction order eligibility per group, rescue stranded runs

The merge order resolved over all candidates, so one oversized or
non-Parquet file without bounds disabled ordering for files that could
participate. Resolve it per partition group over the eligible entries.

Ordered runs too short to merge were dropped entirely. Runs from an
inferred bounds order now fall back to size-based packing — ordering is
a preference there — while runs under a declared sort order are still
left for later passes so the sort contract holds.

An explicit write.parquet.row-group-limit is a cap, not a floor: values
below the estimate floor are now honored instead of being raised.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* iceberg: only use the declared sort order when its bounds are complete

An entry without bounds on the sort column sorted to the tail and merged
into an output claiming an order it cannot verify. Fall back to
bound inference instead of ordering by a later sort field alone.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* iceberg maintenance: keep ordered runs when the full repack yields nothing

The bestEffort leftover fallback removed the ordered runs before checking
whether repacking the whole bin produced any bins, discarding valid
compaction work.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

---------

Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
2026-10-08 22:06:44 +08:00
Chris LuGitHubDevin <158243242+devin-ai-integration[bot]@users.noreply.github.com>Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
8f80dac30f ec: strict_placement option so encode only runs while guarantees hold (#11656)
* ec: strict_placement option so encode only runs while guarantees hold

Shard placement during encode was best-effort (PlaceDurabilityFirst):
when the cluster could not satisfy the per-disk caps, anti-affinity,
replica-placement or per-rack caps, the constraints were relaxed and
the volume was encoded anyway, weaker than configured. A
strict_placement option on the erasure coding task switches planning
to PlaceStrict so the volume's planning fails instead, and the encode
is retried when capacity allows the guarantee.

Also documents the resilience rule in ec.encode help: a volume
survives losing any nodes or racks holding at most parity-shards
shards between them, and how -shardReplicaPlacement's rack and node
digits bound that loss.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* ec: expose strict_placement through the plugin form and persisted task policy

The admin UI, the admin.toml maintenance mapping, and the TaskPolicy
serialization all dropped the new flag; add the bool field to
ErasureCodingTaskConfig, the worker config form, and both conversion
directions.

* shell: describe shardReplicaPlacement as requested limits, not guarantees

ec.encode places shards best-effort, so the configured rack/node caps only
bound shard loss when the final placement actually satisfies them.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

---------

Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
2026-10-08 22:05:02 +08:00
Chris LuGitHubDevin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
30cf53265b deps: pin seaweedfs/goexif at v1.0.3 and stop dependabot re-bumping it (#11648)
* deps: pin github.com/seaweedfs/goexif back to v1.0.3

The fork's newest tagged release is v1.0.3; the v2.0.0+incompatible
requirement resolved to older, untagged code and breaks isolated
builds that fetch it from the proxy.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* ci: stop dependabot bumping seaweedfs/goexif past its latest tag

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* test/kafka: settle goexif at v1.0.3 too

The kafka test module recorded v2.0.0+incompatible in its own requires,
so MVS kept selecting it over the pinned v1.0.3 in the root module.

---------

Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
2026-10-08 22:04:36 +08:00
Chris LuGitHubDevin <158243242+devin-ai-integration[bot]@users.noreply.github.com>Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
4d1f49c638 filer.sync: sign proxied chunk I/O from the per-side security file (#11645)
* filer.sync: sign proxied chunk I/O from the per-side security file

The -a.security / -b.security files were used for gRPC TLS and the
HTTPS client but not for jwt.filer_signing, so filer-proxied chunk
reads and writes carried a token signed with the process-wide key and
failed authorization whenever the two clusters' keys differ.

LoadFilerJwtFromFile returns a FilerJwtProvider for each side's file,
which FilerSource and FilerSink now accept for proxied chunk reads and
writes. With no keys in the file or no flag, both fall back to the
process-wide jwt.filer_signing configuration as before.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* replication: use the side filer read key for manifest downloads and fall back per access level

Manifest chunk resolution still signed proxied downloads with the
process-wide read key, so a source filer requiring its own key 401'd on
manifest-bearing files. A side security file that set only one access
level also produced empty tokens for the other instead of inheriting the
process-wide key, and the side file loader ignored the WEED_ environment
overrides the filer itself honors.

ResolveChunkManifest/ResolveOneChunkManifest keep their signatures;
FilerJwt-aware variants thread the provider down to fetchWholeChunk,
which prefers it on proxy URLs. The side loader now applies the same
environment precedence and falls back to the process-wide signer per
missing access level.

* security: verify the configured filer token lifetimes

* security: reject negative filer token lifetimes

A negative expires_after_seconds reached GenJwtForFilerServer and produced
a token with no expiration claim. Also synchronize the Authorization-header
capture in the proxy test and restore the prior viper key on cleanup.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

---------

Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
2026-10-08 22:04:06 +08:00
Chris LuGitHubDevin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
14fdd61aea filer: stop the aggregated metadata subscribe loop rescanning an exhausted persisted log (#11644)
* fix(filer): gate the aggregated metadata disk pass on real change

A subscriber whose start position is past the end of the local persisted
log re-ran the whole persisted-log pass - store listings, file opens,
readahead - on every loop iteration. Each iteration is paced only by the
shortest wake (the 20ms hold floor on a busy watermark), so one parked
subscriber kept a full CPU core busy for the life of the stream.

The aggregated loop now mirrors the local loop's gate: the disk pass
runs on the first pass and afterwards only when something it cannot
miss changed - a local flush landed, the peers' flush low-watermark
advanced (more content admitted, or new files in a shared store), the
cursor moved, or a disk hold is pending (the ring read that follows an
empty pass parks internally, so skipping there would strand a held
entry).

Regression test: a subscriber parked past the persisted-log tail holds
the listing rate near zero and still delivers once peers report
progress.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* filer: re-arm the aggregated disk pass on unobserved change

Review found three staleness classes the gate could not see: the flush
low-watermark only catching rises (a joining peer lowers the minimum and
invalidates an earlier pass's proof), a peer past the minimum landing a
file without moving it, and a chunk subscriber's refs-stop bound
advancing with wall time. Re-read when the low-watermark moves in either
direction, when the chunk listing bound admits more files, and on a slow
re-probe cadence for files no watermark can signal.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* filer: unwind the parked ring read so the disk re-probe runs, and re-read on cursor rewinds

A caught-up subscriber parks inside LoopProcessLogData's wait loop, so
the re-probe interval in the outer disk gate could never elapse there;
the callback now unwinds the read once the cadence is due so the gate
re-evaluates. The cursor trigger also needs to notice rewinds, not just
advances, since ResumeFromDiskError moves the cursor backward.

---------

Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
2026-10-08 22:03:31 +08:00
github-actions[bot] 49c25890ad docs: regenerate star history chart 2026-10-08 00:52:01 +00:00
48 changed files with 3388 additions and 1252 deletions
+2
View File
@@ -10,3 +10,5 @@ updates:
directory: "/"
schedule:
interval: weekly
ignore:
- dependency-name: "github.com/seaweedfs/goexif"
+1 -2
View File
@@ -64,7 +64,7 @@ require (
github.com/prometheus/procfs v0.22.0
github.com/rcrowley/go-metrics v0.0.0-20201227073835-cf1acfcdf475 // indirect
github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec // indirect
github.com/seaweedfs/goexif v2.0.0+incompatible
github.com/seaweedfs/goexif v1.0.3
github.com/seaweedfs/raft v1.2.1
github.com/sirupsen/logrus v1.9.4 // indirect
github.com/spf13/afero v1.15.0 // indirect
@@ -258,7 +258,6 @@ require (
github.com/rclone/Proton-API-Bridge v1.0.5 // indirect
github.com/rclone/go-proton-api v1.0.4 // indirect
github.com/rogpeppe/go-internal v1.15.0 // indirect
github.com/rwcarlsen/goexif v0.0.0-20190401172101-9e8deecbddbd // indirect
github.com/ryanuber/go-glob v1.0.0 // indirect
github.com/sasha-s/go-deadlock v0.3.1 // indirect
github.com/smarty/assertions v1.15.0 // indirect
+2 -4
View File
@@ -1746,8 +1746,6 @@ github.com/rs/zerolog v1.34.0 h1:k43nTLIwcTVQAncfCw4KZ2VY6ukYoZaBPNOE8txlOeY=
github.com/rs/zerolog v1.34.0/go.mod h1:bJsvje4Z08ROH4Nhs5iH600c3IkWhwp44iRc54W6wYQ=
github.com/ruudk/golang-pdf417 v0.0.0-20181029194003-1af4ab5afa58/go.mod h1:6lfFZQK844Gfx8o5WFuvpxWRwnSoipWe/p622j1v06w=
github.com/ruudk/golang-pdf417 v0.0.0-20201230142125-a7e3863a1245/go.mod h1:pQAZKsJ8yyVxGRWYNEm9oFB8ieLgKFnamEyDmSA0BRk=
github.com/rwcarlsen/goexif v0.0.0-20190401172101-9e8deecbddbd h1:CmH9+J6ZSsIjUK3dcGsnCnO41eRBOnY12zwkn5qVwgc=
github.com/rwcarlsen/goexif v0.0.0-20190401172101-9e8deecbddbd/go.mod h1:hPqNNc0+uJM6H+SuU8sEs5K5IQeKccPqeSjfgcKGgPk=
github.com/ryanuber/go-glob v1.0.0 h1:iQh3xXAumdQ+4Ufa5b25cRpC5TYKlno6hsv6Cb3pkBk=
github.com/ryanuber/go-glob v1.0.0/go.mod h1:807d1WSdnB0XRJzKNil9Om6lcp/3a0v4qIHxIXzX/Yc=
github.com/sabhiram/go-gitignore v0.0.0-20210923224102-525f6e181f06 h1:OkMGxebDjyw0ULyrTYWeN0UNCCkmCWfjPnIA2W6oviI=
@@ -1766,8 +1764,8 @@ github.com/seaweedfs/cockroachdb-parser v0.0.0-20260225204133-2f342c5ea564 h1:Tg
github.com/seaweedfs/cockroachdb-parser v0.0.0-20260225204133-2f342c5ea564/go.mod h1:JSKCh6uCHBz91lQYFYHCyTrSVIPge4SUFVn28iwMNB0=
github.com/seaweedfs/go-fuse/v2 v2.9.4 h1:ACyloiuopdhRSjdLLeSWbsVaemMPskORaRF01TY6GyM=
github.com/seaweedfs/go-fuse/v2 v2.9.4/go.mod h1:zABdmWEa6A0bwaBeEOBUeUkGIZlxUhcdv+V1Dcc/U/I=
github.com/seaweedfs/goexif v2.0.0+incompatible h1:x8pckiT12QQhifwhDQpeISgDfsqmQ6VR4LFPQ64JRps=
github.com/seaweedfs/goexif v2.0.0+incompatible/go.mod h1:Oni780Z236sXpIQzk1XoJlTwqrJ02smEin9zQeff7Fk=
github.com/seaweedfs/goexif v1.0.3 h1:ve/OjI7dxPW8X9YQsv3JuVMaxEyF9Rvfd04ouL+Bz30=
github.com/seaweedfs/goexif v1.0.3/go.mod h1:Oni780Z236sXpIQzk1XoJlTwqrJ02smEin9zQeff7Fk=
github.com/seaweedfs/raft v1.2.1 h1:QgFl/aaPnagpUxYB6Bx+fFss1NyetVVcJmraMmnaQ5Q=
github.com/seaweedfs/raft v1.2.1/go.mod h1:fgs/rAVEzjQ7e04XMzG3eJhwZZRmBW+2uRtjakeCGeU=
github.com/secure-systems-lab/go-securesystemslib v0.11.0 h1:iuCR9kcMFD4QurdKrGvPLoKZLv9YvwPYVr0473BdtFs=
+1042 -1049
View File
File diff suppressed because it is too large Load Diff

Before

Width:  |  Height:  |  Size: 54 KiB

After

Width:  |  Height:  |  Size: 54 KiB

+225 -20
View File
@@ -1182,6 +1182,10 @@ pub struct Volume {
last_modified_ts_seconds: u64,
last_append_at_ns: u64,
last_write_append_at_ns: u64, // AppendAtNs of the newest write; tombstones don't move it
last_write_needle_key: NeedleId, // the write behind the watermark
last_write_deleted: bool, // that write was deleted, so recovery has to rescan
keep_last_modified_ts_on_load: bool,
pub last_disk_check_ns: Arc<std::sync::atomic::AtomicI64>, // for phantom volume detection cache
last_compact_index_offset: u64,
@@ -1281,6 +1285,10 @@ impl Volume {
location_disk_space_low: Arc::new(AtomicBool::new(false)),
last_modified_ts_seconds: 0,
last_append_at_ns: 0,
last_write_append_at_ns: 0,
last_write_needle_key: NeedleId(0),
last_write_deleted: false,
keep_last_modified_ts_on_load: false,
last_disk_check_ns: Arc::new(std::sync::atomic::AtomicI64::new(0)),
last_compact_index_offset: 0,
last_compact_revision: 0,
@@ -1327,6 +1335,10 @@ impl Volume {
location_disk_space_low: Arc::new(AtomicBool::new(false)),
last_modified_ts_seconds: 0,
last_append_at_ns: 0,
last_write_append_at_ns: 0,
last_write_needle_key: NeedleId(0),
last_write_deleted: false,
keep_last_modified_ts_on_load: false,
last_disk_check_ns: Arc::new(std::sync::atomic::AtomicI64::new(0)),
last_compact_index_offset: 0,
last_compact_revision: 0,
@@ -1442,12 +1454,14 @@ impl Volume {
Err(e) => return Err(e.into()),
}
self.last_modified_ts_seconds = metadata
.modified()
.unwrap_or(SystemTime::UNIX_EPOCH)
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
if !self.keep_last_modified_ts_on_load {
self.last_modified_ts_seconds = metadata
.modified()
.unwrap_or(SystemTime::UNIX_EPOCH)
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
}
if metadata.len() >= SUPER_BLOCK_SIZE as u64 {
already_has_super_block = true;
@@ -1549,7 +1563,9 @@ impl Volume {
"volumeDataIntegrityChecking failed"
);
}
self.recover_last_modified_ts();
if !self.keep_last_modified_ts_on_load {
self.recover_last_modified_ts();
}
// Structural check: no .idx entry may reference bytes past the
// end of .dat. The needle map's load walk above already
@@ -2339,6 +2355,13 @@ impl Volume {
.collect();
}
self.last_append_at_ns = last_append_at_ns;
for ((n, _), r) in run.iter().zip(&staged) {
if matches!(r, Ok(Some(_))) && n.append_at_ns > self.last_write_append_at_ns {
self.last_write_append_at_ns = n.append_at_ns;
self.last_write_needle_key = n.id;
self.last_write_deleted = false;
}
}
// A durable entry that fails to publish stops the volume taking
// writes, so the entries after it are refused the way a lone
@@ -2521,6 +2544,9 @@ impl Volume {
}
self.last_append_at_ns = n.append_at_ns;
self.last_write_append_at_ns = n.append_at_ns;
self.last_write_needle_key = n.id;
self.last_write_deleted = false;
self.publish_write(n, offset, fsync)?;
@@ -2839,6 +2865,11 @@ impl Volume {
if let Some(nm) = &mut self.nm {
nm.delete(n.id, Offset::from_actual_offset(offset as i64))?;
}
if n.id == self.last_write_needle_key {
self.last_write_append_at_ns = 0;
self.last_write_needle_key = NeedleId(0);
self.last_write_deleted = true;
}
let checkpoint_ok = self.maybe_checkpoint_index(false);
// Clear the EIO streak after a successful delete (tombstone append +
@@ -3106,8 +3137,15 @@ impl Volume {
return;
}
match self.find_last_write_append_at_ns() {
Ok(0) => {}
Ok(append_at_ns) => self.last_modified_ts_seconds = append_at_ns / 1_000_000_000,
Ok((0, _)) => {}
Ok((append_at_ns, key)) => {
self.last_modified_ts_seconds = append_at_ns / 1_000_000_000;
if append_at_ns > self.last_write_append_at_ns {
self.last_write_append_at_ns = append_at_ns;
self.last_write_needle_key = key;
self.last_write_deleted = false;
}
}
Err(e) => warn!(
volume_id = self.id.0,
error = %e,
@@ -3117,7 +3155,8 @@ impl Volume {
}
/// Scan the .idx backwards for the newest write — an entry that is not a
/// deletion tombstone — and return that needle's append timestamp. The .idx
/// deletion tombstone — and return that needle's append timestamp and key.
/// The .idx
/// and the .dat share an order, so an append-ordered volume answers with the
/// first write the scan reaches. Vacuum rewrites both in key order, which
/// tracks write order only because the master issues keys increasing: an
@@ -3126,19 +3165,21 @@ impl Volume {
/// nothing but tombstones, when a vacuumed volume holds more needles than
/// the scan budget, or for a volume older than version 3, whose needles
/// carry no append timestamp. Mirrors Go's findLastWriteAppendAtNs.
fn find_last_write_append_at_ns(&self) -> Result<u64, VolumeError> {
fn find_last_write_append_at_ns(&self) -> Result<(u64, NeedleId), VolumeError> {
let version = self.version();
if version != VERSION_3 {
return Ok(0);
return Ok((0, NeedleId(0)));
}
let idx_path = self.file_name(".idx");
let idx_size = fs::metadata(&idx_path).map(|m| m.len()).unwrap_or(0) as i64;
if idx_size == 0 || idx_size % NEEDLE_MAP_ENTRY_SIZE as i64 != 0 {
return Ok(0);
return Ok((0, NeedleId(0)));
}
let scan_every_write = self.super_block.compaction_revision > 0;
let mut entry_budget = Self::VACUUMED_LAST_WRITE_SCAN_ENTRIES;
let mut last_write_append_at_ns = 0u64;
let mut last_write_key = NeedleId(0);
let mut dead = HashSet::new();
let mut idx_file = File::open(&idx_path)?;
let mut block = vec![0u8; NEEDLE_MAP_ENTRY_SIZE * idx::ROWS_TO_READ];
let mut end = idx_size;
@@ -3149,7 +3190,13 @@ impl Volume {
idx_file.read_exact(entries)?;
for entry in entries.as_chunks::<NEEDLE_MAP_ENTRY_SIZE>().0.iter().rev() {
let (key, offset, size) = idx_entry_from_bytes(entry);
// The first row a key presents is its latest state: a tombstone
// there retires the write rows beneath it.
if dead.contains(&key) {
continue;
}
if offset.is_zero() || size.is_deleted() {
dead.insert(key);
continue;
}
let Some(needle_offset) =
@@ -3157,10 +3204,13 @@ impl Volume {
else {
continue;
};
last_write_append_at_ns = last_write_append_at_ns
.max(self.read_needle_append_at_ns(needle_offset, size)?);
let append_at_ns = self.read_needle_append_at_ns(needle_offset, size)?;
if append_at_ns > last_write_append_at_ns {
last_write_append_at_ns = append_at_ns;
last_write_key = key;
}
if !scan_every_write {
return Ok(last_write_append_at_ns);
return Ok((last_write_append_at_ns, last_write_key));
}
entry_budget -= 1;
if entry_budget == 0 {
@@ -3169,12 +3219,12 @@ impl Volume {
budget = Self::VACUUMED_LAST_WRITE_SCAN_ENTRIES,
"too many needles to scan for the last write, keeping the .dat mtime"
);
return Ok(0);
return Ok((0, NeedleId(0)));
}
}
end = start;
}
Ok(last_write_append_at_ns)
Ok((last_write_append_at_ns, last_write_key))
}
/// The .dat offset holding the needle an .idx entry describes, or None when
@@ -4294,6 +4344,9 @@ impl Volume {
// Update lastAppendAtNs (matches Go L352: v.lastAppendAtNs = appendAtNs)
self.last_append_at_ns = append_at_ns;
self.last_write_append_at_ns = append_at_ns;
self.last_write_needle_key = needle_id;
self.last_write_deleted = false;
// Update needle map index
let offset = Offset::from_actual_offset(dat_size);
@@ -4493,8 +4546,18 @@ impl Volume {
self.write_compact_commit_marker()?;
self.apply_compact_swap()?;
// Reload
self.load(true, false, 0, self.version())?;
// The write watermark already equals what recover_last_modified_ts
// would rescan, so keep the clock instead of paying for the scan
// under the lock. If its write was itself deleted, the reload
// recovers the newest surviving write instead.
if self.last_write_append_at_ns != 0 {
self.last_modified_ts_seconds = self.last_write_append_at_ns / 1_000_000_000;
}
self.keep_last_modified_ts_on_load =
self.last_modified_ts_seconds != 0 && !self.last_write_deleted;
let load_result = self.load(true, false, 0, self.version());
self.keep_last_modified_ts_on_load = false;
load_result?;
Ok(())
}
@@ -6381,6 +6444,7 @@ mod tests {
let prior = v.nm.as_ref().unwrap().get(NeedleId(1)).unwrap().unwrap();
let dat_len_before = dat_len(&v);
let last_append_before = v.last_append_at_ns;
let last_write_append_before = v.last_write_append_at_ns;
let last_modified_before = v.last_modified_ts_seconds;
let file_count_before = v.file_count();
@@ -6402,6 +6466,7 @@ mod tests {
);
assert_eq!(dat_len(&v), dat_len_before, "the run is off the .dat");
assert_eq!(v.last_append_at_ns, last_append_before);
assert_eq!(v.last_write_append_at_ns, last_write_append_before);
assert_eq!(v.last_modified_ts_seconds, last_modified_before);
let now = v.nm.as_ref().unwrap().get(NeedleId(1)).unwrap().unwrap();
assert_eq!((now.offset, now.size), (prior.offset, prior.size));
@@ -7037,6 +7102,146 @@ mod tests {
);
}
// Covers the reload that ends a vacuum commit: the clock must keep the
// last write's append time rather than be re-derived, and the intervening
// delete's tombstone must not freshen it.
#[test]
fn test_ttl_clock_carried_across_vacuum_commit() {
let tmp = TempDir::new().unwrap();
let dir = tmp.path().to_str().unwrap();
let ttl = crate::storage::needle::ttl::TTL::read("5m").unwrap();
let last_write_ns = (SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap()
.as_secs()
- 2 * 60 * 60)
* 1_000_000_000;
let mut v = make_ttl_volume(dir, ttl);
let mut written = Vec::new();
for i in 1..=3u64 {
let data = format!("data {}", i);
let mut n = Needle {
id: NeedleId(i),
cookie: Cookie(i as u32),
data: data.as_bytes().to_vec(),
data_size: data.len() as u32,
..Needle::default()
};
let (offset, _, _) = v.write_needle(&mut n, true, false).unwrap();
written.push((offset, n.size));
}
v.delete_needle(&mut Needle {
id: NeedleId(2),
cookie: Cookie(2),
..Needle::default()
})
.unwrap();
v.sync_to_disk().unwrap();
for (offset, size) in written {
backdate_append_at_ns(&v.dat_path(), offset, size, last_write_ns);
}
// Where a restart's recovery would have left the clock. The delete
// above pushed last_append_at_ns to ~now; the commit must not use it.
v.set_last_modified_ts_for_test(last_write_ns / 1_000_000_000);
v.last_write_append_at_ns = last_write_ns;
v.compact_by_index(0, 0, |_| true).unwrap();
v.commit_compact().unwrap();
assert_eq!(v.last_modified_ts(), last_write_ns / 1_000_000_000);
assert!(
v.is_expired(v.content_size(), 1024 * 1024),
"a TTL volume whose last write is 2h old must stay expired across a vacuum commit"
);
}
// A write's client supplied modified time can lie ahead of or behind when
// it was appended; the commit keeps the server-side write watermark the
// recovery scan would recompute.
#[test]
fn test_ttl_clock_at_commit_uses_append_time() {
let tmp = TempDir::new().unwrap();
let dir = tmp.path().to_str().unwrap();
let ttl = crate::storage::needle::ttl::TTL::read("5m").unwrap();
let future = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap()
.as_secs()
+ 24 * 60 * 60;
let mut v = make_ttl_volume(dir, ttl);
let data = b"data".to_vec();
let mut n = Needle {
id: NeedleId(1),
cookie: Cookie(1),
data,
data_size: 4,
last_modified: future,
..Needle::default()
};
v.write_needle(&mut n, true, false).unwrap();
assert!(v.last_modified_ts() >= future);
let append_watermark_sec = v.last_write_append_at_ns / 1_000_000_000;
v.compact_by_index(0, 0, |_| true).unwrap();
v.commit_compact().unwrap();
assert_eq!(v.last_modified_ts(), append_watermark_sec);
}
// A vacuum commit whose newest write was deleted first: the carried
// watermark belonged to that write, so the commit has to let the reload
// rescan and land on the newest surviving write rather than keep the
// volume alive on a deleted write's time.
#[test]
fn test_ttl_clock_at_commit_skips_deleted_write() {
let tmp = TempDir::new().unwrap();
let dir = tmp.path().to_str().unwrap();
let ttl = crate::storage::needle::ttl::TTL::read("5m").unwrap();
let now = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap()
.as_secs();
let old_write_ns = (now - 2 * 60 * 60) * 1_000_000_000;
let new_write_ns = (now - 60 * 60) * 1_000_000_000;
let mut v = make_ttl_volume(dir, ttl);
for (id, append_at_ns) in [(1u64, old_write_ns), (2, new_write_ns)] {
let data = format!("data {}", id);
let mut n = Needle {
id: NeedleId(id),
cookie: Cookie(id as u32),
data: data.as_bytes().to_vec(),
data_size: data.len() as u32,
..Needle::default()
};
let (offset, _, _) = v.write_needle(&mut n, true, false).unwrap();
v.sync_to_disk().unwrap();
backdate_append_at_ns(&v.dat_path(), offset, n.size, append_at_ns);
}
// The delete lands inside the commit window: makeup_diff replays its
// tombstone into the new .idx behind the write row the copy carried.
v.compact_by_index(0, 0, |_| true).unwrap();
v.delete_needle(&mut Needle {
id: NeedleId(2),
cookie: Cookie(2),
..Needle::default()
})
.unwrap();
assert!(
v.last_write_deleted,
"deleting the newest write must mark the watermark dead"
);
v.commit_compact().unwrap();
assert_eq!(v.last_modified_ts(), old_write_ns / 1_000_000_000);
assert!(
!v.last_write_deleted,
"the reload's rescan must reseed the watermark off the surviving write"
);
}
// Guard the destroy time an EC volume is reclaimed on: it was recomputed as
// now+TTL every time the .vif was written, so a read-only mark, a tier
// upload or an EC encode handed an already expiring volume another full TTL.
+7 -10
View File
@@ -7,20 +7,20 @@ require (
github.com/linkedin/goavro/v2 v2.15.0
github.com/seaweedfs/seaweedfs v0.0.0-00010101000000-000000000000
github.com/segmentio/kafka-go v0.4.49
github.com/stretchr/testify v1.11.1
github.com/stretchr/testify v1.12.1
google.golang.org/grpc v1.85.0-dev.0.20260915183914-4e49413dcab7
)
replace github.com/seaweedfs/seaweedfs => ../../
require (
github.com/andybalholm/brotli v1.2.2 // indirect
github.com/andybalholm/brotli v1.2.3 // indirect
github.com/aws/aws-sdk-go v1.55.8 // indirect
github.com/beorn7/perks v1.0.1 // indirect
github.com/cespare/xxhash/v2 v2.3.0 // indirect
github.com/cognusion/imaging v1.0.4 // indirect
github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc // indirect
github.com/dustin/go-humanize v1.0.1 // indirect
github.com/dustin/go-humanize v1.1.0 // indirect
github.com/eapache/go-resiliency v1.7.0 // indirect
github.com/eapache/go-xerial-snappy v0.0.0-20230731223053-c322873962e3 // indirect
github.com/eapache/queue v1.1.0 // indirect
@@ -61,17 +61,15 @@ require (
github.com/petermattis/goid v0.0.0-20260113132338-7c7de50cc741 // indirect
github.com/pierrec/lz4/v4 v4.1.29 // indirect
github.com/planetscale/vtprotobuf v0.6.1-0.20240319094008-0393e58bdf10 // indirect
github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 // indirect
github.com/power-devops/perfstat v0.0.0-20240221224432-82ca36839d55 // indirect
github.com/prometheus/client_golang v1.24.1 // indirect
github.com/prometheus/client_model v0.6.3 // indirect
github.com/prometheus/common v0.70.1 // indirect
github.com/prometheus/common v0.71.0 // indirect
github.com/prometheus/procfs v0.22.0 // indirect
github.com/rcrowley/go-metrics v0.0.0-20250401214520-65e299d6c5c9 // indirect
github.com/rdleal/intervalst v1.5.0 // indirect
github.com/rwcarlsen/goexif v0.0.0-20190401172101-9e8deecbddbd // indirect
github.com/sagikazarmark/locafero v0.11.0 // indirect
github.com/seaweedfs/goexif v2.0.0+incompatible // indirect
github.com/seaweedfs/goexif v1.0.3 // indirect
github.com/shirou/gopsutil/v4 v4.26.7 // indirect
github.com/sourcegraph/conc v0.3.1-0.20240121214520-5f936abd7ae8 // indirect
github.com/spf13/afero v1.15.0 // indirect
@@ -91,8 +89,8 @@ require (
github.com/xeipuuv/gojsonreference v0.0.0-20180127040603-bd5ef7bd5415 // indirect
github.com/xeipuuv/gojsonschema v1.2.0 // indirect
github.com/yusufpapurcu/wmi v1.2.4 // indirect
go.yaml.in/yaml/v3 v3.0.4 // indirect
golang.org/x/crypto v0.56.0 // indirect
go.yaml.in/yaml/v3 v3.0.5 // indirect
golang.org/x/crypto v0.57.0 // indirect
golang.org/x/image v0.46.0 // indirect
golang.org/x/net v0.58.0 // indirect
golang.org/x/sync v0.23.0 // indirect
@@ -101,5 +99,4 @@ require (
google.golang.org/genproto/googleapis/rpc v0.0.0-20260819154853-08b0e4226688 // indirect
google.golang.org/grpc/security/advancedtls v1.0.0 // indirect
google.golang.org/protobuf v1.36.12 // indirect
gopkg.in/yaml.v3 v3.0.1 // indirect
)
+16 -20
View File
@@ -10,8 +10,8 @@ github.com/alecthomas/assert/v2 v2.10.0 h1:jjRCHsj6hBJhkmhznrCzoNpbA3zqy0fYiUcYZ
github.com/alecthomas/assert/v2 v2.10.0/go.mod h1:Bze95FyfUr7x34QZrjL+XP+0qgp/zg8yS+TtBj1WA3k=
github.com/alecthomas/repr v0.4.0 h1:GhI2A8MACjfegCPVq9f1FLvIBS+DrQ2KQBFZP1iFzXc=
github.com/alecthomas/repr v0.4.0/go.mod h1:Fr0507jx4eOXV7AlPV6AVZLYrLIuIeSOWtW57eE/O/4=
github.com/andybalholm/brotli v1.2.2 h1:HzTuoo2ErYQqf5qvcJInB8uvqSVxRttzkFexPWtnceM=
github.com/andybalholm/brotli v1.2.2/go.mod h1:rzTDkvFWvIrjDXZHkuS16NPggd91W3kUSvPlQ1pLaKY=
github.com/andybalholm/brotli v1.2.3 h1:8H1qwOkl2LPfjf3YezB90JnCliZb6SInJ/OJkEbA5NQ=
github.com/andybalholm/brotli v1.2.3/go.mod h1:rzTDkvFWvIrjDXZHkuS16NPggd91W3kUSvPlQ1pLaKY=
github.com/aws/aws-sdk-go v1.55.8 h1:JRmEUbU52aJQZ2AjX4q4Wu7t4uZjOu71uyNmaWlUkJQ=
github.com/aws/aws-sdk-go v1.55.8/go.mod h1:ZkViS9AqA6otK+JBBNH2++sx1sgxrPKcSzPPvQkUtXk=
github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM=
@@ -28,8 +28,8 @@ github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSs
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc h1:U9qPSI2PIWSS1VwoXQT9A3Wy9MM3WgvqSxFWenqJduM=
github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/dustin/go-humanize v1.0.1 h1:GzkhY7T5VNhEkwH0PVJgjz+fX1rhBrR7pRT3mDkpeCY=
github.com/dustin/go-humanize v1.0.1/go.mod h1:Mu1zIs6XwVuF/gI1OepvI0qD18qycQx+mFykh5fBlto=
github.com/dustin/go-humanize v1.1.0 h1:dbKTrvD0klcbBV/h4AWJdMuZogJACoMlvWIWZ5b2xWg=
github.com/dustin/go-humanize v1.1.0/go.mod h1:hc1CvRkJMsgxqjmjMQF3QNRAZBwY8AXBAzKYoSX9sFI=
github.com/eapache/go-resiliency v1.7.0 h1:n3NRTnBn5N0Cbi/IeOHuQn9s2UwVUH7Ga0ZWcP+9JTA=
github.com/eapache/go-resiliency v1.7.0/go.mod h1:5yPzW0MIvSe0JDsv0v+DvcjEv2FyD6iZYSs1ZI+iQho=
github.com/eapache/go-xerial-snappy v0.0.0-20230731223053-c322873962e3 h1:Oy0F4ALJ04o5Qqpdz8XLIpNA3WM/iSIXqxtqo7UGVws=
@@ -174,8 +174,8 @@ github.com/prometheus/client_golang v1.24.1 h1:JnJkREXzWxUdCuPFpIWZiPispT9xVV59u
github.com/prometheus/client_golang v1.24.1/go.mod h1:F+oSRECHg4sse5ucfYpYDeIv/hu68Zo0uoHKetWnzcE=
github.com/prometheus/client_model v0.6.3 h1:O0jaTVAYNxTHYInEPFJt5I3+sN8zqBtVMPTB1qyxiEo=
github.com/prometheus/client_model v0.6.3/go.mod h1:gpN5P9S7Rr6Yr92PiQ+Ixvhf6JZEkF1dnxsYL2aPBEM=
github.com/prometheus/common v0.70.1 h1:1HvjP4D5oL3t8RsPlwxA9onvvStjtIHYE5XuuwOi/PY=
github.com/prometheus/common v0.70.1/go.mod h1:VdFUQDMZK3VLkurFUVhia6uys/0suUp86TJz5qbJRhc=
github.com/prometheus/common v0.71.0 h1:9KDAKb7Mj3HEVKyFCK6Dc/HIwlBzZIN2l7/lrHl3KK8=
github.com/prometheus/common v0.71.0/go.mod h1:CLJ5H8TEsGX8bl31BdMkfhIZ+QmZ9tBPPotUxUbfcmk=
github.com/prometheus/procfs v0.22.0 h1:6q9+/JL9IKAPbCmBrv9n5O5Ty3NKnciV5X7YGw0oics=
github.com/prometheus/procfs v0.22.0/go.mod h1:CvmFr/GVhIjIvWJZW3tgkODBQMRIf0EyWMQLHCHab58=
github.com/rcrowley/go-metrics v0.0.0-20250401214520-65e299d6c5c9 h1:bsUq1dX0N8AOIL7EB/X911+m4EHsnWEHeJ0c+3TTBrg=
@@ -184,12 +184,10 @@ github.com/rdleal/intervalst v1.5.0 h1:SEB9bCFz5IqD1yhfH1Wv8IBnY/JQxDplwkxHjT6ha
github.com/rdleal/intervalst v1.5.0/go.mod h1:xO89Z6BC+LQDH+IPQQw/OESt5UADgFD41tYMUINGpxQ=
github.com/rogpeppe/go-internal v1.15.0 h1:D0RCU5rMAp+SpgkiNdrjfJ+LX4J1M32V2NeCY7EJ6hc=
github.com/rogpeppe/go-internal v1.15.0/go.mod h1:DrUVZyrJU+txYW5/1kwtXQSMFio52ZOxX7yM1VHvnxs=
github.com/rwcarlsen/goexif v0.0.0-20190401172101-9e8deecbddbd h1:CmH9+J6ZSsIjUK3dcGsnCnO41eRBOnY12zwkn5qVwgc=
github.com/rwcarlsen/goexif v0.0.0-20190401172101-9e8deecbddbd/go.mod h1:hPqNNc0+uJM6H+SuU8sEs5K5IQeKccPqeSjfgcKGgPk=
github.com/sagikazarmark/locafero v0.11.0 h1:1iurJgmM9G3PA/I+wWYIOw/5SyBtxapeHDcg+AAIFXc=
github.com/sagikazarmark/locafero v0.11.0/go.mod h1:nVIGvgyzw595SUSUE6tvCp3YYTeHs15MvlmU87WwIik=
github.com/seaweedfs/goexif v2.0.0+incompatible h1:x8pckiT12QQhifwhDQpeISgDfsqmQ6VR4LFPQ64JRps=
github.com/seaweedfs/goexif v2.0.0+incompatible/go.mod h1:Oni780Z236sXpIQzk1XoJlTwqrJ02smEin9zQeff7Fk=
github.com/seaweedfs/goexif v1.0.3 h1:ve/OjI7dxPW8X9YQsv3JuVMaxEyF9Rvfd04ouL+Bz30=
github.com/seaweedfs/goexif v1.0.3/go.mod h1:Oni780Z236sXpIQzk1XoJlTwqrJ02smEin9zQeff7Fk=
github.com/segmentio/kafka-go v0.4.49 h1:GJiNX1d/g+kG6ljyJEoi9++PUMdXGAxb7JGPiDCuNmk=
github.com/segmentio/kafka-go v0.4.49/go.mod h1:Y1gn60kzLEEaW28YshXyk2+VCUKbJ3Qr6DrnT3i4+9E=
github.com/shirou/gopsutil/v4 v4.26.7 h1:IXzpHz/dkMRYAhKkOXr1HB6SuzWU3eoyyeWe7g3bNZc=
@@ -216,8 +214,8 @@ github.com/stretchr/testify v1.7.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/
github.com/stretchr/testify v1.7.5/go.mod h1:yNjHg4UonilssWZ8iaSj1OCr/vHnekPRkoO+kdMU+MU=
github.com/stretchr/testify v1.8.0/go.mod h1:yNjHg4UonilssWZ8iaSj1OCr/vHnekPRkoO+kdMU+MU=
github.com/stretchr/testify v1.8.1/go.mod h1:w2LPCIKwWwSfY2zedu0+kehJoqGctiVI29o6fzry7u4=
github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U=
github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U=
github.com/stretchr/testify v1.12.1 h1:EuwCh5fleGS7H32xRwO3wRGT7DxrDhLAT6FF8MpWDWE=
github.com/stretchr/testify v1.12.1/go.mod h1:MDEgiDPPsNp5cuIrHPPCyornHKgEVbtFUmoNlxoYthg=
github.com/subosito/gotenv v1.6.0 h1:9NlTDc1FTs4qu0DDq7AEtTPNw6SVm7uBMsUCUjABIf8=
github.com/subosito/gotenv v1.6.0/go.mod h1:Dk4QP5c2W3ibzajGcXpNraDfq2IrhjMIvMSWPKKo0FU=
github.com/syndtr/goleveldb v1.0.1-0.20190318030020-c3a204f8e965 h1:1oFLiOyVl+W7bnBzGhf7BbIv9loSFQcieWWYIjLqcAw=
@@ -263,15 +261,15 @@ go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto=
go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE=
go.yaml.in/yaml/v2 v2.4.4 h1:tuyd0P+2Ont/d6e2rl3be67goVK4R6deVxCUX5vyPaQ=
go.yaml.in/yaml/v2 v2.4.4/go.mod h1:gMZqIpDtDqOfM0uNfy0SkpRhvUryYH0Z6wdMYcacYXQ=
go.yaml.in/yaml/v3 v3.0.4 h1:tfq32ie2Jv2UxXFdLJdh3jXuOzWiL1fo0bu/FbuKpbc=
go.yaml.in/yaml/v3 v3.0.4/go.mod h1:DhzuOOF2ATzADvBadXxruRBLzYTpT36CKvDb3+aBEFg=
go.yaml.in/yaml/v3 v3.0.5 h1:N6y/pJk8buWs9NY5ERU2HSMfm+IuD/OtfdAnq6kESPw=
go.yaml.in/yaml/v3 v3.0.5/go.mod h1:HVTZu1O7/Vkt2N+BFy8Zza+lnLsABggaTM2ZpNIGuKg=
golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w=
golang.org/x/crypto v0.0.0-20191011191535-87dc89f01550/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI=
golang.org/x/crypto v0.0.0-20200622213623-75b288015ac9/go.mod h1:LzIPMQfyMNhhGPhUkYOs5KpL4U8rLKemX1yGLhDgUto=
golang.org/x/crypto v0.0.0-20210921155107-089bfa567519/go.mod h1:GvvjBRRGRdwPK5ydBHafDWAxML/pGHZbMvKqRZ5+Abc=
golang.org/x/crypto v0.6.0/go.mod h1:OFC/31mSvZgRz0V1QTNCzfAI1aIRzbiufJtkMIlEp58=
golang.org/x/crypto v0.56.0 h1:GUh5Ii4J5jtcseSMiRqr1jXCNHoxjeV9Fmekc2oLy6Y=
golang.org/x/crypto v0.56.0/go.mod h1:OMW5y6CY9l38uPLmxU6l6pwcXp1obtLo3e6gT7gQR2I=
golang.org/x/crypto v0.57.0 h1:3ZVCjf8Ggz7zneR/EHRVx68Ctf+2pmIMP2UFhh9cC6M=
golang.org/x/crypto v0.57.0/go.mod h1:Fdz0i5U6CoizGwLda9DttjSk6qlZo25zYNtR+ycvuZA=
golang.org/x/image v0.46.0 h1:b1+oYj0Jbp6K5MDT4i4/eZpYlk3V8SJhhDKh6LBHAyQ=
golang.org/x/image v0.46.0/go.mod h1:3B3W05VGVQyuXucLINLjXKrqISASfi4Xj+iCVkLMwew=
golang.org/x/mod v0.3.0/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA=
@@ -288,8 +286,8 @@ golang.org/x/net v0.6.0/go.mod h1:2Tu9+aMcznHK/AK1HMvgo6xiTLG5rD5rZLDS+rp2Bjs=
golang.org/x/net v0.7.0/go.mod h1:2Tu9+aMcznHK/AK1HMvgo6xiTLG5rD5rZLDS+rp2Bjs=
golang.org/x/net v0.58.0 h1:ynWG7rqYi4ccpTEuPZ2QGWHktVEM9DMCj9yzDE0Q7To=
golang.org/x/net v0.58.0/go.mod h1:YwCddHnFlT7eLQqVprV19OnhLGtc5xOKgE0RyqgfWAU=
golang.org/x/oauth2 v0.36.0 h1:peZ/1z27fi9hUOFCAZaHyrpWG5lwe0RJEEEeH0ThlIs=
golang.org/x/oauth2 v0.36.0/go.mod h1:YDBUJMTkDnJS+A4BP4eZBjCqtokkg1hODuPjwiGPO7Q=
golang.org/x/oauth2 v0.37.0 h1:JUlcxA8oAtauLfiH8FX2/FkAWHAdi0QtGCGc+hofE98=
golang.org/x/oauth2 v0.37.0/go.mod h1:IxwZNxUULJmpBFf9K/9NTMSIfZZuvuTy1gGxhigP/58=
golang.org/x/sync v0.0.0-20180314180146-1d60e4601c6f/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/sync v0.0.0-20201020160332-67f06af15bc9/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
@@ -351,8 +349,6 @@ google.golang.org/protobuf v1.23.0/go.mod h1:EGpADcykh3NcUnDUJcl1+ZksZNG86OlYog2
google.golang.org/protobuf v1.36.12 h1:pJOKDDOyeXErUroCihFAd5LQuwXBSpVnKGrj5o/fwxc=
google.golang.org/protobuf v1.36.12/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco=
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk=
gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q=
gopkg.in/fsnotify.v1 v1.4.7/go.mod h1:Tz8NjZHkW78fSQdbUxIjBTcgA1z1m8ZHf0WmKUhAMys=
gopkg.in/tomb.v1 v1.0.0-20141024135613-dd632973f1e7 h1:uRGJdciOHaEIrze2W8Q3AKkepLTh2hOroT7a+7czfdQ=
gopkg.in/tomb.v1 v1.0.0-20141024135613-dd632973f1e7/go.mod h1:dt/ZhP58zS4L8KSrWDmTeBkI65Dw0HsyUHuEVlX15mw=
+9
View File
@@ -98,6 +98,10 @@ func (cp *ConfigPersistence) ApplyMaintenanceConfigFromToml(v TomlConfig) error
ecConf.ReplicaPlacement = v.GetString(k)
ecChanged = true
}
if k := "maintenance.erasure_coding.strict_placement"; v.IsSet(k) {
ecConf.StrictPlacement = v.GetBool(k)
ecChanged = true
}
if !maintenanceChanged && !vacuumChanged && !balanceChanged && !ecChanged {
return nil
@@ -224,6 +228,7 @@ var pluginConfigSections = []pluginConfigSection{
"min_size_mb": int64Value,
"preferred_tags": stringListValue,
"replica_placement": stringValue,
"strict_placement": boolValue,
},
// workers read collection_filter from the admin values, not the worker values
adminKeys: map[string]func(v TomlConfig, key string) *plugin_pb.ConfigValue{
@@ -240,6 +245,10 @@ func int64Value(v TomlConfig, key string) *plugin_pb.ConfigValue {
return &plugin_pb.ConfigValue{Kind: &plugin_pb.ConfigValue_Int64Value{Int64Value: int64(v.GetInt(key))}}
}
func boolValue(v TomlConfig, key string) *plugin_pb.ConfigValue {
return &plugin_pb.ConfigValue{Kind: &plugin_pb.ConfigValue_BoolValue{BoolValue: v.GetBool(key)}}
}
func stringValue(v TomlConfig, key string) *plugin_pb.ConfigValue {
return &plugin_pb.ConfigValue{Kind: &plugin_pb.ConfigValue_StringValue{StringValue: v.GetString(key)}}
}
+21 -3
View File
@@ -174,17 +174,24 @@ func runFilerSynchronize(cmd *Command, args []string) bool {
// per-cluster HTTPS clients for volume server connections
var httpClientA, httpClientB *util_http_client.HTTPClient
var jwtForFilerA, jwtForFilerB security.FilerJwtProvider
if *syncOptions.aSecurity != "" {
var err error
if httpClientA, err = security.LoadHTTPClientFromFile(*syncOptions.aSecurity); err != nil {
glog.Fatalf("load HTTPS client config for filer A: %v", err)
}
if jwtForFilerA, err = security.LoadFilerJwtFromFile(*syncOptions.aSecurity); err != nil {
glog.Fatalf("load filer JWT config for filer A: %v", err)
}
}
if *syncOptions.bSecurity != "" {
var err error
if httpClientB, err = security.LoadHTTPClientFromFile(*syncOptions.bSecurity); err != nil {
glog.Fatalf("load HTTPS client config for filer B: %v", err)
}
if jwtForFilerB, err = security.LoadFilerJwtFromFile(*syncOptions.bSecurity); err != nil {
glog.Fatalf("load filer JWT config for filer B: %v", err)
}
}
grace.SetupProfiling(*syncCpuProfile, *syncMemProfile)
@@ -263,7 +270,9 @@ func runFilerSynchronize(cmd *Command, args []string) bool {
bFilerSignature,
&syncStateA2B,
httpClientA,
httpClientB)
httpClientB,
jwtForFilerA,
jwtForFilerB)
if err != nil {
glog.Errorf("sync from %s to %s: %v", *syncOptions.filerA, *syncOptions.filerB, err)
time.Sleep(1747 * time.Millisecond)
@@ -306,7 +315,9 @@ func runFilerSynchronize(cmd *Command, args []string) bool {
aFilerSignature,
&syncStateB2A,
httpClientB,
httpClientA)
httpClientA,
jwtForFilerB,
jwtForFilerA)
if err != nil {
glog.Errorf("sync from %s to %s: %v", *syncOptions.filerB, *syncOptions.filerA, err)
time.Sleep(2147 * time.Millisecond)
@@ -336,7 +347,8 @@ func initOffsetFromTsMs(grpcDialOption grpc.DialOption, targetFiler pb.ServerAdd
func doSubscribeFilerMetaChanges(clientId int32, clientEpoch int32, sourceGrpcDialOption grpc.DialOption, sourceFiler pb.ServerAddress, sourcePath string, sourceExcludePaths []string, sourceReadChunkFromFiler bool, targetGrpcDialOption grpc.DialOption, targetFiler pb.ServerAddress, targetPath string,
replicationStr, collection string, ttlSec int, sinkWriteChunkByFiler bool, diskType string, debug bool, concurrency int, chunkConcurrency int, doDeleteFiles bool, sourceFilerSignature int32, targetFilerSignature int32, statePtr *atomic.Pointer[syncState],
sourceHttpClient *util_http_client.HTTPClient, sinkHttpClient *util_http_client.HTTPClient) error {
sourceHttpClient *util_http_client.HTTPClient, sinkHttpClient *util_http_client.HTTPClient,
sourceJwtProvider security.FilerJwtProvider, sinkJwtProvider security.FilerJwtProvider) error {
// if first time, start from now
// if has previously synced, resume from that point of time
@@ -357,12 +369,18 @@ func doSubscribeFilerMetaChanges(clientId int32, clientEpoch int32, sourceGrpcDi
if sourceHttpClient != nil {
filerSource.SetHttpClient(sourceHttpClient)
}
if sourceJwtProvider != nil {
filerSource.SetFilerJwtProvider(sourceJwtProvider)
}
filerSink := &filersink.FilerSink{}
filerSink.DoInitialize(targetFiler.ToHttpAddress(), targetFiler.ToGrpcAddress(), targetPath, replicationStr, collection, ttlSec, diskType, targetGrpcDialOption, sinkWriteChunkByFiler)
filerSink.SetChunkConcurrency(chunkConcurrency)
if sinkHttpClient != nil {
filerSink.SetUploader(operation.NewUploaderWithHttpClient(sinkHttpClient))
}
if sinkJwtProvider != nil {
filerSink.SetFilerJwtProvider(sinkJwtProvider)
}
filerSink.SetSourceFiler(filerSource)
persistEventFn := genProcessFunction(sourcePath, targetPath, sourceExcludePaths, nil, nil, nil, filerSink, doDeleteFiles, debug)
+2
View File
@@ -64,6 +64,8 @@
# preferred_tags = ["fast", "ssd"]
# EC shard placement constraint, e.g. "020"; empty uses the master default replication
# replica_placement = ""
# fail planning when the placement constraints cannot be met instead of relaxing them
# strict_placement = false
# max retry attempts for a failed erasure coding job
# retry_limit = 1
# seconds to wait between retry attempts
+1 -1
View File
@@ -240,7 +240,7 @@ func (group *ChunkGroup) SetChunks(chunks []*filer_pb.FileChunk) error {
continue
}
resolvedChunks, err := resolveOneChunkManifest(context.Background(), group.lookupFn, chunk, group.cacheInvalidator, group.manifestCache)
resolvedChunks, err := resolveOneChunkManifest(context.Background(), group.lookupFn, chunk, group.cacheInvalidator, group.manifestCache, nil)
if err != nil {
group.resolveErr = err
return err
+35 -9
View File
@@ -16,6 +16,7 @@ import (
"github.com/seaweedfs/seaweedfs/weed/glog"
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
"github.com/seaweedfs/seaweedfs/weed/security"
"github.com/seaweedfs/seaweedfs/weed/util"
util_http "github.com/seaweedfs/seaweedfs/weed/util/http"
)
@@ -55,7 +56,16 @@ func SeparateManifestChunks(chunks []*filer_pb.FileChunk) (manifestChunks, nonMa
}
func ResolveChunkManifest(ctx context.Context, lookupFileIdFn wdclient.LookupFileIdFunctionType, chunks []*filer_pb.FileChunk, startOffset, stopOffset int64, invalidator CacheInvalidator) (dataChunks, manifestChunks []*filer_pb.FileChunk, manifestResolveErr error) {
resolver := newChunkManifestResolver(ctx, lookupFileIdFn, invalidator)
resolver := newChunkManifestResolver(ctx, lookupFileIdFn, invalidator, nil)
defer resolver.close()
return resolver.resolve(chunks, startOffset, stopOffset)
}
// ResolveChunkManifestWithFilerJwt is ResolveChunkManifest that signs proxied
// manifest downloads with filerJwtFn instead of the process-wide filer read
// key, for readers carrying per-source credentials.
func ResolveChunkManifestWithFilerJwt(ctx context.Context, lookupFileIdFn wdclient.LookupFileIdFunctionType, filerJwtFn security.FilerJwtProvider, chunks []*filer_pb.FileChunk, startOffset, stopOffset int64, invalidator CacheInvalidator) (dataChunks, manifestChunks []*filer_pb.FileChunk, manifestResolveErr error) {
resolver := newChunkManifestResolver(ctx, lookupFileIdFn, invalidator, filerJwtFn)
defer resolver.close()
return resolver.resolve(chunks, startOffset, stopOffset)
}
@@ -81,6 +91,7 @@ type chunkManifestResolver struct {
cancel context.CancelFunc
lookupFileIdFn wdclient.LookupFileIdFunctionType
invalidator CacheInvalidator
filerJwtFn security.FilerJwtProvider
jobs chan chunkManifestResolveJob
overflowSem chan struct{}
workers sync.WaitGroup
@@ -88,7 +99,7 @@ type chunkManifestResolver struct {
started bool
}
func newChunkManifestResolver(ctx context.Context, lookupFileIdFn wdclient.LookupFileIdFunctionType, invalidator CacheInvalidator) *chunkManifestResolver {
func newChunkManifestResolver(ctx context.Context, lookupFileIdFn wdclient.LookupFileIdFunctionType, invalidator CacheInvalidator, filerJwtFn security.FilerJwtProvider) *chunkManifestResolver {
workCtx, cancel := context.WithCancel(ctx)
resolver := &chunkManifestResolver{
ctx: workCtx,
@@ -96,6 +107,7 @@ func newChunkManifestResolver(ctx context.Context, lookupFileIdFn wdclient.Looku
cancel: cancel,
lookupFileIdFn: lookupFileIdFn,
invalidator: invalidator,
filerJwtFn: filerJwtFn,
jobs: make(chan chunkManifestResolveJob, chunkManifestResolveJobBufferSize),
overflowSem: make(chan struct{}, maxChunkManifestResolveWorkers),
}
@@ -103,7 +115,7 @@ func newChunkManifestResolver(ctx context.Context, lookupFileIdFn wdclient.Looku
}
func (r *chunkManifestResolver) executeJob(job chunkManifestResolveJob) {
job.result.chunks, job.result.err = ResolveOneChunkManifest(job.batchCtx, r.lookupFileIdFn, job.chunk, r.invalidator)
job.result.chunks, job.result.err = resolveOneChunkManifest(job.batchCtx, r.lookupFileIdFn, job.chunk, r.invalidator, nil, r.filerJwtFn)
if job.result.err != nil && r.parentCtx.Err() == nil {
if job.batchCtx.Err() != nil && errors.Is(job.result.err, context.Canceled) {
job.result.internalCancel = true
@@ -293,14 +305,23 @@ func (r *chunkManifestResolver) resolve(chunks []*filer_pb.FileChunk, startOffse
// Keeping this signature stable preserves the existing four-argument contract
// for external callers.
func ResolveOneChunkManifest(ctx context.Context, lookupFileIdFn wdclient.LookupFileIdFunctionType, chunk *filer_pb.FileChunk, invalidator CacheInvalidator) (dataChunks []*filer_pb.FileChunk, manifestResolveErr error) {
return resolveOneChunkManifest(ctx, lookupFileIdFn, chunk, invalidator, nil)
return resolveOneChunkManifest(ctx, lookupFileIdFn, chunk, invalidator, nil, nil)
}
// ResolveOneChunkManifestWithFilerJwt is ResolveOneChunkManifest that signs
// proxied manifest downloads with filerJwtFn instead of the process-wide
// filer read key, for readers carrying per-source credentials.
func ResolveOneChunkManifestWithFilerJwt(ctx context.Context, lookupFileIdFn wdclient.LookupFileIdFunctionType, filerJwtFn security.FilerJwtProvider, chunk *filer_pb.FileChunk, invalidator CacheInvalidator) (dataChunks []*filer_pb.FileChunk, manifestResolveErr error) {
return resolveOneChunkManifest(ctx, lookupFileIdFn, chunk, invalidator, nil, filerJwtFn)
}
// resolveOneChunkManifest is the cache-aware implementation. cache may be nil,
// in which case the manifest is fetched and validated on every call, matching
// the historical uncached behavior. A non-nil cache is owned by a single mount
// (WFS) and coalesces concurrent cold misses via singleflight.
func resolveOneChunkManifest(ctx context.Context, lookupFileIdFn wdclient.LookupFileIdFunctionType, chunk *filer_pb.FileChunk, invalidator CacheInvalidator, cache *ChunkManifestCache) (dataChunks []*filer_pb.FileChunk, manifestResolveErr error) {
// (WFS) and coalesces concurrent cold misses via singleflight. filerJwtFn
// overrides the filer read credential for proxied downloads; nil means the
// process-wide key.
func resolveOneChunkManifest(ctx context.Context, lookupFileIdFn wdclient.LookupFileIdFunctionType, chunk *filer_pb.FileChunk, invalidator CacheInvalidator, cache *ChunkManifestCache, filerJwtFn security.FilerJwtProvider) (dataChunks []*filer_pb.FileChunk, manifestResolveErr error) {
if !chunk.IsChunkManifest {
return
}
@@ -318,7 +339,7 @@ func resolveOneChunkManifest(ctx context.Context, lookupFileIdFn wdclient.Lookup
bytesBuffer := bytesBufferPool.Get().(*bytes.Buffer)
bytesBuffer.Reset()
defer bytesBufferPool.Put(bytesBuffer)
if err := fetchWholeChunk(ctx, bytesBuffer, lookupFileIdFn, key.fileID, chunk.CipherKey, chunk.IsCompressed, invalidator); err != nil {
if err := fetchWholeChunk(ctx, bytesBuffer, lookupFileIdFn, key.fileID, chunk.CipherKey, chunk.IsCompressed, invalidator, filerJwtFn); err != nil {
return nil, fmt.Errorf("fail to read manifest %s: %w", key.fileID, err)
}
// Copy before the buffer returns to the pool so concurrent callers
@@ -357,13 +378,18 @@ func resolveOneChunkManifest(ctx context.Context, lookupFileIdFn wdclient.Lookup
return m.Chunks, nil
}
func fetchWholeChunk(ctx context.Context, bytesBuffer *bytes.Buffer, lookupFileIdFn wdclient.LookupFileIdFunctionType, fileId string, cipherKey []byte, isGzipped bool, invalidator CacheInvalidator) error {
func fetchWholeChunk(ctx context.Context, bytesBuffer *bytes.Buffer, lookupFileIdFn wdclient.LookupFileIdFunctionType, fileId string, cipherKey []byte, isGzipped bool, invalidator CacheInvalidator, filerJwtFn security.FilerJwtProvider) error {
urlStrings, err := lookupFileIdFn(ctx, fileId)
if err != nil {
glog.ErrorfCtx(ctx, "operation LookupFileId %s failed, err: %v", fileId, err)
return err
}
jwt := ChunkReadJwt(urlStrings, fileId)
var jwt string
if filerJwtFn != nil && len(urlStrings) > 0 && util_http.IsProxyChunkUrl(urlStrings[0]) {
jwt = string(filerJwtFn(false))
} else {
jwt = ChunkReadJwt(urlStrings, fileId)
}
if _, err = retriedStreamFetchChunkData(ctx, bytesBuffer, urlStrings, jwt, cipherKey, isGzipped, true, 0, 0, refreshUrls(ctx, invalidator, lookupFileIdFn, fileId)); err == nil {
return nil
}
@@ -23,7 +23,7 @@ func BenchmarkManifestResolutionRepeatedOpen(b *testing.B) {
}
chunk := newManifestCacheTestChunk("benchmark-cached")
cache := NewChunkManifestCache(MaxMountChunkManifestCacheEntries, MaxMountChunkManifestCacheBytes)
_, err := resolveOneChunkManifest(context.Background(), lookup, chunk, nil, cache)
_, err := resolveOneChunkManifest(context.Background(), lookup, chunk, nil, cache, nil)
if err != nil {
b.Fatal(err)
}
@@ -33,7 +33,7 @@ func BenchmarkManifestResolutionRepeatedOpen(b *testing.B) {
b.ReportAllocs()
b.ResetTimer()
for i := 0; i < b.N; i++ {
if _, err := resolveOneChunkManifest(context.Background(), lookup, chunk, nil, cache); err != nil {
if _, err := resolveOneChunkManifest(context.Background(), lookup, chunk, nil, cache, nil); err != nil {
b.Fatal(err)
}
}
+3 -3
View File
@@ -258,7 +258,7 @@ func TestResolveOneChunkManifestHonorsCanceledContextOnCacheHit(t *testing.T) {
return nil, errors.New("lookup should not be called")
}
_, err := resolveOneChunkManifest(ctx, lookup, chunk, nil, cache)
_, err := resolveOneChunkManifest(ctx, lookup, chunk, nil, cache, nil)
require.ErrorIs(t, err, context.Canceled)
require.False(t, lookupCalled, "a canceled cache hit must not issue a lookup")
}
@@ -278,7 +278,7 @@ func TestResolveOneChunkManifestCanceledWaiterReturnsDuringCoalescedMiss(t *test
defer leaderCancel()
leaderDone := make(chan error, 1)
go func() {
_, err := resolveOneChunkManifest(leaderCtx, fixture.lookup, chunk, nil, cache)
_, err := resolveOneChunkManifest(leaderCtx, fixture.lookup, chunk, nil, cache, nil)
leaderDone <- err
}()
@@ -287,7 +287,7 @@ func TestResolveOneChunkManifestCanceledWaiterReturnsDuringCoalescedMiss(t *test
// for the leader's result.
waiterCtx, waiterCancel := context.WithCancel(context.Background())
waiterCancel()
_, err := resolveOneChunkManifest(waiterCtx, fixture.lookup, chunk, nil, cache)
_, err := resolveOneChunkManifest(waiterCtx, fixture.lookup, chunk, nil, cache, nil)
require.ErrorIs(t, err, context.Canceled)
// The leader must still complete successfully and populate the cache.
+47 -5
View File
@@ -16,6 +16,8 @@ import (
"google.golang.org/protobuf/proto"
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
"github.com/seaweedfs/seaweedfs/weed/security"
util_http "github.com/seaweedfs/seaweedfs/weed/util/http"
)
func TestDoMaybeManifestize(t *testing.T) {
@@ -621,7 +623,7 @@ func TestFetchWholeChunkRetriesFreshLocations(t *testing.T) {
inv := &countingInvalidator{}
bytesBuffer := fetchManifestBuffer(t)
assert.NoError(t, fetchWholeChunk(context.Background(), bytesBuffer, lookup.lookup, "5,stale", nil, false, inv))
assert.NoError(t, fetchWholeChunk(context.Background(), bytesBuffer, lookup.lookup, "5,stale", nil, false, inv, nil))
assert.Equal(t, int32(1), inv.invalidations.Load())
assert.Equal(t, int32(2), lookup.calls.Load())
@@ -655,7 +657,7 @@ func TestFetchWholeChunkRefreshesLocationsAfterPartialFailure(t *testing.T) {
inv := &countingInvalidator{}
bytesBuffer := fetchManifestBuffer(t)
assert.NoError(t, fetchWholeChunk(context.Background(), bytesBuffer, lookup.lookup, "5,abc", nil, false, inv))
assert.NoError(t, fetchWholeChunk(context.Background(), bytesBuffer, lookup.lookup, "5,abc", nil, false, inv, nil))
assert.Equal(t, int32(1), inv.invalidations.Load())
assert.Equal(t, int32(2), lookup.calls.Load())
decoded := &filer_pb.FileChunkManifest{}
@@ -676,7 +678,7 @@ func TestFetchWholeChunkWithoutInvalidator(t *testing.T) {
freshUrls: []string{"http://unused:8080/5,abc"},
}
assert.Error(t, fetchWholeChunk(context.Background(), fetchManifestBuffer(t), lookup.lookup, "5,abc", nil, false, nil))
assert.Error(t, fetchWholeChunk(context.Background(), fetchManifestBuffer(t), lookup.lookup, "5,abc", nil, false, nil, nil))
assert.Equal(t, int32(1), lookup.calls.Load())
}
@@ -694,7 +696,7 @@ func TestFetchWholeChunkUnchangedLocations(t *testing.T) {
}
inv := &countingInvalidator{}
assert.Error(t, fetchWholeChunk(context.Background(), fetchManifestBuffer(t), lookup.lookup, "5,abc", nil, false, inv))
assert.Error(t, fetchWholeChunk(context.Background(), fetchManifestBuffer(t), lookup.lookup, "5,abc", nil, false, inv, nil))
assert.Equal(t, int32(2), lookup.calls.Load())
assert.Equal(t, int32(1), inv.invalidations.Load())
}
@@ -711,7 +713,7 @@ func TestFetchWholeChunkCancelledKeepsLocations(t *testing.T) {
ctx, cancel := context.WithCancel(context.Background())
cancel()
err := fetchWholeChunk(ctx, fetchManifestBuffer(t), lookup.lookup, "5,abc", nil, false, inv)
err := fetchWholeChunk(ctx, fetchManifestBuffer(t), lookup.lookup, "5,abc", nil, false, inv, nil)
assert.ErrorIs(t, err, context.Canceled)
assert.Equal(t, int32(0), inv.invalidations.Load())
assert.Equal(t, int32(1), lookup.calls.Load())
@@ -730,3 +732,43 @@ func TestFetchWholeChunkCancelledKeepsLocations(t *testing.T) {
})
assert.ErrorIs(t, noInvalidator, context.Canceled)
}
// TestFetchWholeChunkUsesProvidedFilerJwt covers a replicating reader whose
// source filer authenticates proxied downloads with its own read key: the
// supplied provider's token must reach the server, not the process-wide one.
func TestFetchWholeChunkUsesProvidedFilerJwt(t *testing.T) {
manifestBytes, err := proto.Marshal(&filer_pb.FileChunkManifest{
Chunks: []*filer_pb.FileChunk{{FileId: "100,abc", Offset: 0, Size: 8}},
})
assert.NoError(t, err)
gotAuth := make(chan string, 1)
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
gotAuth <- r.Header.Get("Authorization")
w.Header().Set("Content-Length", strconv.Itoa(len(manifestBytes)))
w.Write(manifestBytes)
}))
t.Cleanup(srv.Close)
lookup := func(ctx context.Context, fileId string) ([]string, error) {
return []string{srv.URL + "/?" + util_http.ProxyChunkIdParam + "=" + fileId}, nil
}
jwtFn := func(isWrite bool) security.EncodedJwt {
assert.False(t, isWrite)
return "side-read-jwt"
}
bytesBuffer := fetchManifestBuffer(t)
assert.NoError(t, fetchWholeChunk(context.Background(), bytesBuffer, lookup, "5,abc", nil, false, nil, jwtFn))
assert.Equal(t, security.BearerPrefix+"side-read-jwt", <-gotAuth)
// non-proxy URLs keep the volume-server credential and never call the provider
volumeURL := manifestServer(t, manifestBytes).URL + "/5,abc"
volumeLookup := func(ctx context.Context, fileId string) ([]string, error) {
return []string{volumeURL}, nil
}
bytesBuffer.Reset()
assert.NoError(t, fetchWholeChunk(context.Background(), bytesBuffer, volumeLookup, "5,abc", nil, false, nil, func(bool) security.EncodedJwt {
t.Fatal("provider must not be consulted for a volume url")
return ""
}))
}
+1 -1
View File
@@ -190,7 +190,7 @@ func loadLogFileEntries(masterClient *wdclient.MasterClient, chunk *filer_pb.Fil
lookupFileIdFn := func(ctx context.Context, fileId string) (targetUrls []string, err error) {
return masterClient.LookupFileId(ctx, fileId)
}
if fetchErr := fetchWholeChunk(context.Background(), bytesBuffer, lookupFileIdFn, chunk.GetFileIdString(), chunk.CipherKey, chunk.IsCompressed, masterClient); fetchErr != nil {
if fetchErr := fetchWholeChunk(context.Background(), bytesBuffer, lookupFileIdFn, chunk.GetFileIdString(), chunk.CipherKey, chunk.IsCompressed, masterClient, nil); fetchErr != nil {
return nil, false, fetchErr
}
return decodeLogRecords(bytesBuffer.Bytes())
+6 -1
View File
@@ -62,6 +62,7 @@ type UploadOption struct {
SourceUrl string // optional: for logging when reading from a remote source
MaxAttempts int // <=0 uses the default
GenUploadUrl func(host, fileId string) string // if nil → fallback "http://{host}/{fileId}"
FilerJwt security.FilerJwtProvider // credential for proxy chunk URLs; nil → process-wide jwt.filer_signing
}
type UploadResult struct {
@@ -217,7 +218,11 @@ func (uploader *Uploader) uploadWithRetryData(assignFn func() (fileId string, ho
// The request addresses the filer, which authorizes it and mints the
// volume credential itself. The AssignVolume token is not a filer
// credential and gets the caller nowhere here.
uploadOption.Jwt = security.EncodedJwt(util_http.JwtForFilerServer(true))
if uploadOption.FilerJwt != nil {
uploadOption.Jwt = uploadOption.FilerJwt(true)
} else {
uploadOption.Jwt = security.EncodedJwt(util_http.JwtForFilerServer(true))
}
}
uploadResult, err = uploader.retriedUploadData(context.Background(), data, uploadOption)
+1
View File
@@ -380,6 +380,7 @@ message ErasureCodingTaskConfig {
string collection_filter = 4; // Only process volumes from specific collections
repeated string preferred_tags = 5; // Disk tags to prioritize for EC shard placement
string replica_placement = 6; // EC shard replica placement (e.g. "020"); empty falls back to master default replication
bool strict_placement = 7; // fail planning instead of relaxing placement constraints
}
// BalanceTaskConfig contains balance-specific configuration
+11 -2
View File
@@ -2980,6 +2980,7 @@ type ErasureCodingTaskConfig struct {
CollectionFilter string `protobuf:"bytes,4,opt,name=collection_filter,json=collectionFilter,proto3" json:"collection_filter,omitempty"` // Only process volumes from specific collections
PreferredTags []string `protobuf:"bytes,5,rep,name=preferred_tags,json=preferredTags,proto3" json:"preferred_tags,omitempty"` // Disk tags to prioritize for EC shard placement
ReplicaPlacement string `protobuf:"bytes,6,opt,name=replica_placement,json=replicaPlacement,proto3" json:"replica_placement,omitempty"` // EC shard replica placement (e.g. "020"); empty falls back to master default replication
StrictPlacement bool `protobuf:"varint,7,opt,name=strict_placement,json=strictPlacement,proto3" json:"strict_placement,omitempty"` // fail planning instead of relaxing placement constraints
unknownFields protoimpl.UnknownFields
sizeCache protoimpl.SizeCache
}
@@ -3056,6 +3057,13 @@ func (x *ErasureCodingTaskConfig) GetReplicaPlacement() string {
return ""
}
func (x *ErasureCodingTaskConfig) GetStrictPlacement() bool {
if x != nil {
return x.StrictPlacement
}
return false
}
// BalanceTaskConfig contains balance-specific configuration
type BalanceTaskConfig struct {
state protoimpl.MessageState `protogen:"open.v1"`
@@ -4285,14 +4293,15 @@ const file_worker_proto_rawDesc = "" +
"\x10VacuumTaskConfig\x12+\n" +
"\x11garbage_threshold\x18\x01 \x01(\x01R\x10garbageThreshold\x12/\n" +
"\x14min_volume_age_hours\x18\x02 \x01(\x05R\x11minVolumeAgeHours\x120\n" +
"\x14min_interval_seconds\x18\x03 \x01(\x05R\x12minIntervalSeconds\"\x9a\x02\n" +
"\x14min_interval_seconds\x18\x03 \x01(\x05R\x12minIntervalSeconds\"\xc5\x02\n" +
"\x17ErasureCodingTaskConfig\x12%\n" +
"\x0efullness_ratio\x18\x01 \x01(\x01R\rfullnessRatio\x12*\n" +
"\x11quiet_for_seconds\x18\x02 \x01(\x05R\x0fquietForSeconds\x12+\n" +
"\x12min_volume_size_mb\x18\x03 \x01(\x05R\x0fminVolumeSizeMb\x12+\n" +
"\x11collection_filter\x18\x04 \x01(\tR\x10collectionFilter\x12%\n" +
"\x0epreferred_tags\x18\x05 \x03(\tR\rpreferredTags\x12+\n" +
"\x11replica_placement\x18\x06 \x01(\tR\x10replicaPlacement\"\x9b\x01\n" +
"\x11replica_placement\x18\x06 \x01(\tR\x10replicaPlacement\x12)\n" +
"\x10strict_placement\x18\a \x01(\bR\x0fstrictPlacement\"\x9b\x01\n" +
"\x11BalanceTaskConfig\x12/\n" +
"\x13imbalance_threshold\x18\x01 \x01(\x01R\x12imbalanceThreshold\x12(\n" +
"\x10min_server_count\x18\x02 \x01(\x05R\x0eminServerCount\x12+\n" +
+15
View File
@@ -125,6 +125,21 @@ func ReadIntConfig(values map[string]*plugin_pb.ConfigValue, field string, fallb
return int(v)
}
// ReadBoolConfig reads a bool-valued plugin config field.
func ReadBoolConfig(values map[string]*plugin_pb.ConfigValue, field string, fallback bool) bool {
if values == nil {
return fallback
}
value := values[field]
if value == nil {
return fallback
}
if kind, ok := value.Kind.(*plugin_pb.ConfigValue_BoolValue); ok {
return kind.BoolValue
}
return fallback
}
// ReadBytesConfig reads a bytes-valued plugin config field, returning nil when
// the value is missing or of a different kind.
func ReadBytesConfig(values map[string]*plugin_pb.ConfigValue, field string) []byte {
@@ -195,7 +195,7 @@ func (fs *FilerSink) replicateOneManifestChunk(ctx context.Context, sourceChunk
resolveName := fmt.Sprintf("resolve manifest %s", sourceChunk.GetFileIdString())
missingGate := fs.newMissingSourceChunkGate(sourceChunk.GetFileIdString())
err := util.RetryUntil(resolveName, func() error {
rc, e := filer.ResolveOneChunkManifest(ctx, fs.filerSource.LookupFileId, sourceChunk, nil)
rc, e := filer.ResolveOneChunkManifestWithFilerJwt(ctx, fs.filerSource.LookupFileId, fs.filerSource.FilerJwt(), sourceChunk, nil)
if e != nil {
return e
}
@@ -315,6 +315,7 @@ func (fs *FilerSink) uploadManifestChunk(path string, sourceMtimeNs int64, sourc
}
if fs.writeChunkByFiler {
uploadOption.GenUploadUrl = operation.GenUploadUrlProxy(fs.address)
uploadOption.FilerJwt = fs.jwtForFiler
}
currentFileId, uploadResult, uploadErr, _ := uploader.UploadWithRetry(
fs,
@@ -435,6 +436,7 @@ func (fs *FilerSink) fetchAndWrite(sourceChunk *filer_pb.FileChunk, path string,
}
if fs.writeChunkByFiler {
uploadOption.GenUploadUrl = operation.GenUploadUrlProxy(fs.address)
uploadOption.FilerJwt = fs.jwtForFiler
}
currentFileId, uploadResult, uploadErr, _ := uploader.UploadWithRetry(
fs,
+11 -4
View File
@@ -62,6 +62,7 @@ type FilerSink struct {
signature int32
activeTransfers sync.Map // chunkFileId -> *ChunkTransferStatus
uploader *operation.Uploader
jwtForFiler security.FilerJwtProvider
// lastServedFileId is the most recent chunk the source did serve, the probe
// sourceStillServesChunks re-checks before writing an entry off.
lastServedFileId atomic.Pointer[string]
@@ -116,6 +117,12 @@ func (fs *FilerSink) SetUploader(uploader *operation.Uploader) {
fs.uploader = uploader
}
// SetFilerJwtProvider sets the filer API credential for proxied chunk writes.
// Must be called during initialization, before any replication goroutines start.
func (fs *FilerSink) SetFilerJwtProvider(provider security.FilerJwtProvider) {
fs.jwtForFiler = provider
}
func (fs *FilerSink) getUploader() (*operation.Uploader, error) {
if fs.uploader != nil {
return fs.uploader, nil
@@ -358,7 +365,7 @@ func (fs *FilerSink) UpdateEntry(key string, oldEntry *filer_pb.Entry, newParent
existingEntry.RemoteEntry = newEntry.RemoteEntry
default:
// source-side chunks resolve via source filer; sink volume IDs may collide.
deletedChunks, newChunks, err := compareChunks(context.Background(), filer.LookupFn(fs.filerSource), oldEntry, newEntry)
deletedChunks, newChunks, err := compareChunks(context.Background(), filer.LookupFn(fs.filerSource), fs.filerSource.FilerJwt(), oldEntry, newEntry)
if err != nil {
return true, fmt.Errorf("replicate %s compare chunks error: %w", key, err)
}
@@ -404,12 +411,12 @@ func (fs *FilerSink) UpdateEntry(key string, oldEntry *filer_pb.Entry, newParent
})
}
func compareChunks(ctx context.Context, lookupFileIdFn wdclient.LookupFileIdFunctionType, oldEntry, newEntry *filer_pb.Entry) (deletedChunks, newChunks []*filer_pb.FileChunk, err error) {
aData, aMeta, aErr := filer.ResolveChunkManifest(ctx, lookupFileIdFn, oldEntry.GetChunks(), 0, math.MaxInt64, nil)
func compareChunks(ctx context.Context, lookupFileIdFn wdclient.LookupFileIdFunctionType, filerJwtFn security.FilerJwtProvider, oldEntry, newEntry *filer_pb.Entry) (deletedChunks, newChunks []*filer_pb.FileChunk, err error) {
aData, aMeta, aErr := filer.ResolveChunkManifestWithFilerJwt(ctx, lookupFileIdFn, filerJwtFn, oldEntry.GetChunks(), 0, math.MaxInt64, nil)
if aErr != nil {
return nil, nil, aErr
}
bData, bMeta, bErr := filer.ResolveChunkManifest(ctx, lookupFileIdFn, newEntry.GetChunks(), 0, math.MaxInt64, nil)
bData, bMeta, bErr := filer.ResolveChunkManifestWithFilerJwt(ctx, lookupFileIdFn, filerJwtFn, newEntry.GetChunks(), 0, math.MaxInt64, nil)
if bErr != nil {
return nil, nil, bErr
}
+16 -1
View File
@@ -34,6 +34,7 @@ type FilerSource struct {
dataCenter string
signature int32
httpClient *util_http_client.HTTPClient
jwtForFiler security.FilerJwtProvider
}
func (fs *FilerSource) Initialize(configuration util.Configuration, prefix string) error {
@@ -67,6 +68,16 @@ func (fs *FilerSource) SetHttpClient(client *util_http_client.HTTPClient) {
fs.httpClient = client
}
func (fs *FilerSource) SetFilerJwtProvider(provider security.FilerJwtProvider) {
fs.jwtForFiler = provider
}
// FilerJwt returns the side-specific filer API credential, or nil when the
// process-wide jwt.filer_signing configuration applies.
func (fs *FilerSource) FilerJwt() security.FilerJwtProvider {
return fs.jwtForFiler
}
func (fs *FilerSource) LookupFileId(ctx context.Context, part string) (fileUrls []string, err error) {
vid2Locations := make(map[string]*filer_pb.Locations)
@@ -126,7 +137,11 @@ func (fs *FilerSource) ReadPart(fileId string, offset int64) (filename string, h
if fs.proxyByFiler {
fileUrl := util_http.ProxyChunkUrl(fs.address, fileId)
filename, header, resp, err = downloadFn(fileUrl, util_http.JwtForFilerServer(false), offset)
jwt := util_http.JwtForFilerServer(false)
if fs.jwtForFiler != nil {
jwt = string(fs.jwtForFiler(false))
}
filename, header, resp, err = downloadFn(fileUrl, jwt, offset)
if err == nil {
err = readPartStatusError(fileUrl, resp)
}
+56
View File
@@ -8,6 +8,8 @@ import (
jwt "github.com/golang-jwt/jwt/v5"
"github.com/seaweedfs/seaweedfs/weed/glog"
"github.com/seaweedfs/seaweedfs/weed/util"
"github.com/spf13/viper"
)
type EncodedJwt string
@@ -151,3 +153,57 @@ func DecodeJwt(signingKey SigningKey, tokenString EncodedJwt, claims jwt.Claims)
return []byte(signingKey), nil
})
}
// FilerJwtProvider signs the credential a filer's HTTP API expects. A nil
// provider means the process-wide jwt.filer_signing configuration applies.
type FilerJwtProvider func(isWrite bool) EncodedJwt
// LoadFilerJwtFromFile reads jwt.filer_signing from a security file the way
// LoadClientTLSFromFile reads the TLS section, honoring the same WEED_
// environment precedence. A nil provider means the file configures no filer
// signing keys and the process-wide configuration applies. A file that sets
// only one access level falls back to the process-wide key for the other.
func LoadFilerJwtFromFile(configFile string) (FilerJwtProvider, error) {
v := viper.New()
v.SetConfigFile(configFile)
v.AutomaticEnv()
v.SetEnvPrefix("weed")
v.SetEnvKeyReplacer(strings.NewReplacer(".", "_"))
if err := v.ReadInConfig(); err != nil {
return nil, fmt.Errorf("failed to read security config %s: %v", configFile, err)
}
signingKey := SigningKey(v.GetString("jwt.filer_signing.key"))
readSigningKey := SigningKey(v.GetString("jwt.filer_signing.read.key"))
if len(signingKey) == 0 && len(readSigningKey) == 0 {
return nil, nil
}
signingKeyExpires := v.GetInt("jwt.filer_signing.expires_after_seconds")
readSigningKeyExpires := v.GetInt("jwt.filer_signing.read.expires_after_seconds")
if len(signingKey) == 0 || len(readSigningKey) == 0 {
gv := util.GetViper()
if len(signingKey) == 0 {
signingKey = SigningKey(gv.GetString("jwt.filer_signing.key"))
signingKeyExpires = gv.GetInt("jwt.filer_signing.expires_after_seconds")
}
if len(readSigningKey) == 0 {
readSigningKey = SigningKey(gv.GetString("jwt.filer_signing.read.key"))
readSigningKeyExpires = gv.GetInt("jwt.filer_signing.read.expires_after_seconds")
}
}
if signingKeyExpires < 0 || readSigningKeyExpires < 0 {
return nil, fmt.Errorf("jwt.filer_signing lifetimes must not be negative")
}
if signingKeyExpires == 0 {
signingKeyExpires = 10
}
if readSigningKeyExpires == 0 {
readSigningKeyExpires = 60
}
return func(isWrite bool) EncodedJwt {
if isWrite {
return GenJwtForFilerServer(signingKey, signingKeyExpires)
}
return GenJwtForFilerServer(readSigningKey, readSigningKeyExpires)
}, nil
}
+146
View File
@@ -0,0 +1,146 @@
package security
import (
"os"
"path/filepath"
"testing"
"time"
"github.com/seaweedfs/seaweedfs/weed/util"
)
func TestLoadFilerJwtFromFile(t *testing.T) {
configFile := filepath.Join(t.TempDir(), "security.toml")
config := `
[jwt]
[jwt.filer_signing]
key = "side-write-key"
expires_after_seconds = 30
[jwt.filer_signing.read]
key = "side-read-key"
expires_after_seconds = 90
`
if err := os.WriteFile(configFile, []byte(config), 0644); err != nil {
t.Fatal(err)
}
provider, err := LoadFilerJwtFromFile(configFile)
if err != nil {
t.Fatal(err)
}
if provider == nil {
t.Fatal("config with filer signing keys gave a nil provider")
}
for _, tc := range []struct {
name string
isWrite bool
signedBy string
otherKey string
expires int64
}{
{"read", false, "side-read-key", "side-write-key", 90},
{"write", true, "side-write-key", "side-read-key", 30},
} {
t.Run(tc.name, func(t *testing.T) {
before := time.Now()
token := provider(tc.isWrite)
after := time.Now()
claims := &SeaweedFilerClaims{}
if _, err := DecodeJwt(SigningKey(tc.signedBy), token, claims); err != nil {
t.Fatalf("token does not validate against the %s key: %v", tc.name, err)
}
if claims.ExpiresAt == nil {
t.Fatal("token never expires")
}
expiresIn := claims.ExpiresAt.Time
if expiresIn.Before(before.Add(time.Duration(tc.expires-1)*time.Second)) || expiresIn.After(after.Add(time.Duration(tc.expires+1)*time.Second)) {
t.Fatalf("token expires at %v, want %ds after %v", expiresIn, tc.expires, before)
}
if _, err := DecodeJwt(SigningKey(tc.otherKey), token, &SeaweedFilerClaims{}); err == nil {
t.Fatal("token also validates against the other access level's key")
}
})
}
}
func TestLoadFilerJwtFromFileWithoutKeys(t *testing.T) {
configFile := filepath.Join(t.TempDir(), "security.toml")
if err := os.WriteFile(configFile, []byte("[grpc.client]\n"), 0644); err != nil {
t.Fatal(err)
}
provider, err := LoadFilerJwtFromFile(configFile)
if err != nil {
t.Fatal(err)
}
if provider != nil {
t.Fatal("config without filer signing keys gave a provider")
}
}
func TestLoadFilerJwtFromFileMissing(t *testing.T) {
if _, err := LoadFilerJwtFromFile(filepath.Join(t.TempDir(), "none.toml")); err == nil {
t.Fatal("missing config file loaded without error")
}
}
func TestLoadFilerJwtFromFilePartialKeys(t *testing.T) {
configFile := filepath.Join(t.TempDir(), "security.toml")
config := `
[jwt.filer_signing.read]
key = "side-read-key"
`
if err := os.WriteFile(configFile, []byte(config), 0644); err != nil {
t.Fatal(err)
}
gv := util.GetViper()
priorKey := gv.GetString("jwt.filer_signing.key")
gv.Set("jwt.filer_signing.key", "global-write-key")
t.Cleanup(func() { gv.Set("jwt.filer_signing.key", priorKey) })
provider, err := LoadFilerJwtFromFile(configFile)
if err != nil {
t.Fatal(err)
}
if provider == nil {
t.Fatal("config with a filer signing key gave a nil provider")
}
if _, err := DecodeJwt(SigningKey("global-write-key"), provider(true), &SeaweedFilerClaims{}); err != nil {
t.Fatalf("write token does not validate against the global write key: %v", err)
}
if _, err := DecodeJwt(SigningKey("side-read-key"), provider(false), &SeaweedFilerClaims{}); err != nil {
t.Fatalf("read token does not validate against the side read key: %v", err)
}
}
func TestLoadFilerJwtFromFileEnvOverride(t *testing.T) {
configFile := filepath.Join(t.TempDir(), "security.toml")
config := `
[jwt.filer_signing]
key = "side-write-key"
[jwt.filer_signing.read]
key = "side-read-key"
`
if err := os.WriteFile(configFile, []byte(config), 0644); err != nil {
t.Fatal(err)
}
t.Setenv("WEED_JWT_FILER_SIGNING_READ_KEY", "env-read-key")
provider, err := LoadFilerJwtFromFile(configFile)
if err != nil {
t.Fatal(err)
}
if provider == nil {
t.Fatal("config with filer signing keys gave a nil provider")
}
if _, err := DecodeJwt(SigningKey("env-read-key"), provider(false), &SeaweedFilerClaims{}); err != nil {
t.Fatalf("read token does not validate against the env override key: %v", err)
}
if _, err := DecodeJwt(SigningKey("side-write-key"), provider(true), &SeaweedFilerClaims{}); err != nil {
t.Fatalf("write token does not validate against the side write key: %v", err)
}
}
+129 -68
View File
@@ -27,6 +27,11 @@ var (
// (possibly-unflushed) gap, in case the flush notification is missed.
unflushedGapRetryInterval = 2 * time.Second
// aggDiskReprobeInterval paces the aggregated persisted-log re-listing for
// files no watermark signals: a peer past the flush low-watermark can land
// a file without moving the minimum.
aggDiskReprobeInterval = 2 * time.Second
// gapStallWarnInterval paces the warning for a subscriber that stays parked.
gapStallWarnInterval = time.Minute
@@ -666,8 +671,9 @@ func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest,
var lastHeartbeatNs int64
baseEachLogEntryFn := eachLogEntryFn(req, sender, eachEventNotificationFn, &unsyncedEvents)
// heldAtTsNs remembers the entry a read was held at (for the log line);
// the rewind target is the last entry actually delivered.
var heldAtTsNs int64
// diskHeldAtTsNs is the same marker for the disk pass alone: a pending
// disk hold keeps the pass re-reading until the entry is served.
var heldAtTsNs, diskHeldAtTsNs int64
// Each read path holds at its own watermark: persisted logs are complete
// only up to every peer's flush watermark, the ring only up to every
// peer's delivery watermark.
@@ -698,7 +704,14 @@ func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest,
// What the last disk pass proved covered: flushed on every peer AND inside
// the pass's listing, so an empty pass proves (cursor, proven] empty.
var diskPassProvenTsNs int64
diskEachLogEntryFn := guardedEachLogEntryFn(func() int64 { return diskPassHoldTsNs })
diskBaseEachLogEntryFn := guardedEachLogEntryFn(func() int64 { return diskPassHoldTsNs })
diskEachLogEntryFn := func(logEntry *filer_pb.LogEntry) (bool, error) {
isDone, err := diskBaseEachLogEntryFn(logEntry)
if errors.Is(err, errHeldByPeerWatermark) {
diskHeldAtTsNs = logEntry.TsNs
}
return isDone, err
}
memEachLogEntryFn := guardedEachLogEntryFn(holdMemTsNs)
// waitHeld pauses a held read until a peer reports further progress, or
// the retry interval elapses (a peer dropped past its grace, or a log file
@@ -733,6 +746,11 @@ func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest,
var readPersistedLogErr error
var readInMemoryLogErr error
var isDone bool
var lastCheckedFlushTsNs int64 = -1 // Track the last local flush we read the disk under
var lastCheckedFlushLowTsNs int64 = -1 // Track the last peer flush low-watermark we read the disk under
var lastDiskReadTsNs int64 = -1 // Track the last read position we used for disk read
var lastDiskRefsStopTsNs int64 = -1 // Track the last chunk listing bound we read the disk under
var lastDiskPassAt time.Time // Paces re-probes for files no watermark signals
sentRefs := make(map[string]sentRefState)
aggBuffer := fs.filer.MetaAggregator.MetaLogBuffer
@@ -765,78 +783,110 @@ func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest,
// diskPassHoldTsNs above).
diskPassFlushLowTsNs = fs.filer.MetaAggregator.PeerLowFlushWatermarkTsNs()
diskPassHoldTsNs = resolveAggReadHoldTsNs(diskPassFlushLowTsNs, time.Now().UnixNano(), metadataGapSettledHorizon)
diskPassProvenTsNs = diskPassFlushLowTsNs
// Re-read the disk only when something changed it cannot miss: a local
// flush landed, the peers' flush low-watermark moved in either
// direction (a joining peer invalidates what an earlier pass proved),
// the cursor moved, or a disk hold is pending (the ring read after an
// empty pass parks internally, so skipping would strand the held
// entry). A peer past the low-watermark can still land a file without
// moving it, so the listing is also re-probed at a slow cadence;
// chunk listings re-arm as soon as their read bound admits more files.
currentFlushTsNs := fs.filer.LocalMetaLogBuffer.GetLastFlushTsNs()
currentReadTsNs := lastReadTime.Time.UnixNano()
var currentRefsStopTsNs int64
if req.ClientSupportsMetadataChunks {
refsStopTsNs := chunkRefsStopTsNs(diskPassHoldTsNs, req.UntilNs)
// Nothing above the listing bound is proven by this pass.
if refsStopTsNs < diskPassProvenTsNs {
diskPassProvenTsNs = refsStopTsNs
}
if refsStopTsNs > lastReadTime.Time.UnixNano() {
processedTsNs, isDone, readPersistedLogErr = fs.chunkDiskPass(ctx, sender, lastReadTime, refsStopTsNs, sentRefs, nil)
} else {
processedTsNs, isDone, readPersistedLogErr = 0, false, nil
}
} else {
processedTsNs, isDone, readPersistedLogErr = fs.filer.ReadPersistedLogBuffer(ctx, lastReadTime, req.UntilNs, diskEachLogEntryFn)
currentRefsStopTsNs = chunkRefsStopTsNs(diskPassHoldTsNs, req.UntilNs)
}
if errors.Is(readPersistedLogErr, errHeldByPeerWatermark) {
// Stay at the last delivered entry; the held entry is re-read (and
// re-checked) by the next pass.
if processedTsNs > 0 {
shouldReadFromDisk := lastCheckedFlushTsNs == -1 ||
currentFlushTsNs > lastCheckedFlushTsNs ||
diskPassFlushLowTsNs != lastCheckedFlushLowTsNs ||
currentReadTsNs != lastDiskReadTsNs ||
currentRefsStopTsNs > lastDiskRefsStopTsNs ||
diskHeldAtTsNs != 0 ||
time.Since(lastDiskPassAt) >= aggDiskReprobeInterval
diskAdvanced := false
if shouldReadFromDisk {
lastCheckedFlushTsNs = currentFlushTsNs
lastCheckedFlushLowTsNs = diskPassFlushLowTsNs
lastDiskReadTsNs = currentReadTsNs
lastDiskRefsStopTsNs = currentRefsStopTsNs
lastDiskPassAt = time.Now()
diskHeldAtTsNs = 0
diskPassProvenTsNs = diskPassFlushLowTsNs
if req.ClientSupportsMetadataChunks {
refsStopTsNs := currentRefsStopTsNs
// Nothing above the listing bound is proven by this pass.
if refsStopTsNs < diskPassProvenTsNs {
diskPassProvenTsNs = refsStopTsNs
}
if refsStopTsNs > lastReadTime.Time.UnixNano() {
processedTsNs, isDone, readPersistedLogErr = fs.chunkDiskPass(ctx, sender, lastReadTime, refsStopTsNs, sentRefs, nil)
} else {
processedTsNs, isDone, readPersistedLogErr = 0, false, nil
}
} else {
processedTsNs, isDone, readPersistedLogErr = fs.filer.ReadPersistedLogBuffer(ctx, lastReadTime, req.UntilNs, diskEachLogEntryFn)
}
if errors.Is(readPersistedLogErr, errHeldByPeerWatermark) {
// Stay at the last delivered entry; the held entry is re-read (and
// re-checked) by the next pass.
if processedTsNs > 0 {
lastReadTime = log_buffer.NewMessagePosition(processedTsNs, gapResumeCursorOffset)
if processedTsNs > diskAnchorTsNs {
diskAnchorTsNs = processedTsNs
}
}
// A hold is not a gap: clear any stale ResumeFromDiskError so the
// next pass's disk-miss handling cannot skip past the held entry.
readInMemoryLogErr = nil
if !waitHeld("disk", flushChan) {
return nil
}
continue
}
if readPersistedLogErr != nil {
return fmt.Errorf("reading from persisted logs: %w", readPersistedLogErr)
}
if isDone {
return nil
}
glog.V(4).Infof("processed to %v: %v", clientName, processedTsNs)
diskAdvanced = diskReadAdvanced(processedTsNs, lastReadTime)
// Read after the disk read (an eviction landing mid-read must count) and
// in received-ts space: the ring's bumped stopTimes exceed anything on
// any peer's disk, and gating disk cursors on them parks subscribers
// that drained every peer's log.
lastEvictedTsNs := fs.filer.MetaAggregator.MetaLogBuffer.GetLastEvictedOriginalTsNs()
if diskAdvanced {
gapStall.resumed()
reportUnprovenAggregatedCrossing(cursorBeforeDiskTsNs, processedTsNs, lastEvictedTsNs, diskPassFlushLowTsNs, clientName, req.PathPrefix)
lastReadTime = log_buffer.NewMessagePosition(processedTsNs, gapResumeCursorOffset)
if processedTsNs > diskAnchorTsNs {
diskAnchorTsNs = processedTsNs
}
}
// A hold is not a gap: clear any stale ResumeFromDiskError so the
// next pass's disk-miss handling cannot skip past the held entry.
readInMemoryLogErr = nil
if !waitHeld("disk", flushChan) {
return nil
}
continue
}
if readPersistedLogErr != nil {
return fmt.Errorf("reading from persisted logs: %w", readPersistedLogErr)
}
if isDone {
return nil
}
glog.V(4).Infof("processed to %v: %v", clientName, processedTsNs)
diskAdvanced := diskReadAdvanced(processedTsNs, lastReadTime)
// Read after the disk read (an eviction landing mid-read must count) and
// in received-ts space: the ring's bumped stopTimes exceed anything on
// any peer's disk, and gating disk cursors on them parks subscribers
// that drained every peer's log.
lastEvictedTsNs := fs.filer.MetaAggregator.MetaLogBuffer.GetLastEvictedOriginalTsNs()
if diskAdvanced {
gapStall.resumed()
reportUnprovenAggregatedCrossing(cursorBeforeDiskTsNs, processedTsNs, lastEvictedTsNs, diskPassFlushLowTsNs, clientName, req.PathPrefix)
lastReadTime = log_buffer.NewMessagePosition(processedTsNs, gapResumeCursorOffset)
if processedTsNs > diskAnchorTsNs {
diskAnchorTsNs = processedTsNs
}
} else if readInMemoryLogErr == nil {
// Nothing on disk and memory never spoke: scan forward for the next
// day that has logs.
nextDayTs := util.GetNextDayTsNano(lastReadTime.Time.UnixNano())
// The day jump delivers nothing; stay put until the hold point
// covers the skipped range.
if nextDayTs <= diskPassHoldTsNs {
position := log_buffer.NewMessagePosition(nextDayTs, gapResumeCursorOffset)
found, err := fs.filer.HasPersistedLogFiles(position)
if err != nil {
return fmt.Errorf("checking persisted log files: %w", err)
}
if found {
gapStall.resumed()
reportUnprovenAggregatedCrossing(cursorBeforeDiskTsNs, nextDayTs, lastEvictedTsNs, diskPassFlushLowTsNs, clientName, req.PathPrefix)
lastReadTime = position
if nextDayTs > diskAnchorTsNs {
diskAnchorTsNs = nextDayTs
} else if readInMemoryLogErr == nil {
// Nothing on disk and memory never spoke: scan forward for the next
// day that has logs.
nextDayTs := util.GetNextDayTsNano(lastReadTime.Time.UnixNano())
// The day jump delivers nothing; stay put until the hold point
// covers the skipped range.
if nextDayTs <= diskPassHoldTsNs {
position := log_buffer.NewMessagePosition(nextDayTs, gapResumeCursorOffset)
found, err := fs.filer.HasPersistedLogFiles(position)
if err != nil {
return fmt.Errorf("checking persisted log files: %w", err)
}
if found {
gapStall.resumed()
reportUnprovenAggregatedCrossing(cursorBeforeDiskTsNs, nextDayTs, lastEvictedTsNs, diskPassFlushLowTsNs, clientName, req.PathPrefix)
lastReadTime = position
if nextDayTs > diskAnchorTsNs {
diskAnchorTsNs = nextDayTs
}
}
}
}
@@ -867,6 +917,7 @@ func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest,
// every event whose original timestamp is at or below this.
preMemDeliveryLowTsNs := fs.filer.MetaAggregator.PeerLowWatermarkTsNs()
diskReprobeDue := false
lastReadTime, isDone, readInMemoryLogErr = fs.filer.MetaAggregator.MetaLogBuffer.LoopProcessLogData(aggReaderName, lastReadTime, req.UntilNs, func() bool {
select {
case <-ctx.Done():
@@ -876,6 +927,13 @@ func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest,
if !fs.hasClient(req.ClientId, req.ClientEpoch) {
return false
}
// Caught-up readers park in the inner wait loop; the outer disk
// gate never runs again unless this read returns, so unwind to
// re-probe the persisted logs on the slow cadence.
if time.Since(lastDiskPassAt) >= aggDiskReprobeInterval {
diskReprobeDue = true
return false
}
// Contiguous and caught up: advance the anchor to the delivery
// low-watermark so long live tails keep eviction rewinds short.
// Only once the run is connected to the ring - the empty-ring
@@ -919,6 +977,9 @@ func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest,
}
}
if isDone {
if diskReprobeDue {
continue
}
return nil
}
if !fs.hasClient(req.ClientId, req.ClientEpoch) {
+79
View File
@@ -0,0 +1,79 @@
package weed_server
// The persisted log may end before a subscriber's start position (a filer
// whose own journal is older than the position a backup client resumes from).
// With nothing on disk and every ring entry held by a peer watermark, the
// aggregated loop used to re-list and re-read the persisted log on every wake
// - about a full CPU core per such subscriber. The disk pass may only re-run
// when something it cannot miss has changed.
import (
"context"
"strings"
"sync/atomic"
"testing"
"time"
"github.com/seaweedfs/seaweedfs/weed/filer"
"github.com/seaweedfs/seaweedfs/weed/util"
)
type countingStore struct {
filer.FilerStore
logLists *atomic.Int64
}
func (s *countingStore) ListDirectoryPrefixedEntries(ctx context.Context, dirPath util.FullPath, startFileName string, includeStartFile bool, limit int64, prefix string, eachEntryFunc filer.ListEachEntryFunc) (lastFileName string, err error) {
if strings.HasPrefix(string(dirPath), filer.SystemLogDir) {
s.logLists.Add(1)
}
return s.FilerStore.ListDirectoryPrefixedEntries(ctx, dirPath, startFileName, includeStartFile, limit, prefix, eachEntryFunc)
}
func TestSubscribeLoop_AggregatedNoPersistedEntryAfterStart(t *testing.T) {
h := newSubscribeHarness(t)
lists := &atomic.Int64{}
h.f.SetStore(&countingStore{FilerStore: h.f.GetStore(), logLists: lists})
// Persisted log ends at T1: one flushed window, nothing after.
h.append(h.tsAt(0, 0))
h.append(h.tsAt(0, 1))
h.f.LocalMetaLogBuffer.ForceFlush()
waitForFlushedFiles(t, h, h.tsAt(0, 1))
// Client cursor sits just past the last persisted entry.
cursor := h.tsAt(0, 1) + int64(time.Millisecond)
ma := h.startAggregator()
// Aggregated ring holds only much newer entries (peer events).
recent := time.Now().UnixNano()
h.appendAggregated(recent)
h.appendAggregated(recent + int64(time.Millisecond))
// Peers' watermarks are stuck at the old log tail.
reportPeersAt(ma, h.tsAt(0, 1), h.tsAt(0, 1))
r := h.subscribeAggregated(cursor)
// Warm up: the first pass plus the cursor-move re-read after the gap
// machinery re-arms the cursor are legitimate.
time.Sleep(150 * time.Millisecond)
before := lists.Load()
time.Sleep(500 * time.Millisecond)
rate := float64(lists.Load()-before) / 0.5
if rate > 4 {
t.Fatalf("%.1f persisted-log listings per second while parked; the disk pass re-ran on every wake", rate)
}
// Peer progress through the held entries releases the read: the held
// events are delivered without another disk pass.
lists.Store(0)
reportPeersAt(ma, recent+int64(time.Millisecond), h.tsAt(0, 1))
waitForEvents(t, r, []int64{recent, recent + int64(time.Millisecond)}, 3*time.Second)
if got := lists.Load(); got > 4 {
t.Fatalf("%d listings while draining held entries; delivery should come from the ring", got)
}
}
+9
View File
@@ -47,6 +47,15 @@ func (c *commandEcEncode) Help() string {
If you only have less than 4 volume servers, with erasure coding, at least you can afford to
have 4 corrupted shard files.
The guarantee follows from where shards land: a volume survives the loss of any
nodes (or racks) that hold at most parityShards (4) shards between them. Spread
is best-effort; -shardReplicaPlacement requests limits: its rack digit
sets the requested shards per rack, and its node digit the requested shards
per node (the data-center digit is not used for EC). For example
-shardReplicaPlacement=021 requests at most 1 shard per node and 2 per rack.
Only when the final placement meets these limits does losing one rack cost
at most 2 shards; the command does not guarantee that the limits are met.
The -collection parameter is a comma-separated list of collection names, with
"*" and "?" wildcards, and regex patterns:
- One collection: ec.encode -collection="mybucket"
+6 -2
View File
@@ -54,8 +54,12 @@ type Volume struct {
asyncRequestsChan chan *needle.AsyncRequest
asyncWorkerClosed bool
lastModifiedTsSeconds uint64 // unix time in seconds
lastAppendAtNs uint64 // unix time in nanoseconds
lastModifiedTsSeconds uint64 // unix time in seconds
lastAppendAtNs uint64 // unix time in nanoseconds
lastWriteAppendAtNs uint64 // AppendAtNs of the newest write; tombstones don't move it
lastWriteNeedleKey types.NeedleId // the write behind the watermark
lastWriteDeleted bool // that write was deleted, so recovery has to rescan
keepLastModifiedTsOnLoad bool
lastCompactIndexOffset uint64
lastCompactRevision uint16
+34 -18
View File
@@ -264,7 +264,7 @@ func (v *Volume) recoverLastModifiedTs(indexFile *os.File) {
if err != nil || indexSize == 0 {
return
}
appendAtNs, err := findLastWriteAppendAtNs(v, indexFile, indexSize)
appendAtNs, key, err := findLastWriteAppendAtNs(v, indexFile, indexSize)
if err != nil {
glog.Warningf("volume %d recover last write from %s: %v", v.Id, indexFile.Name(), err)
return
@@ -273,6 +273,11 @@ func (v *Volume) recoverLastModifiedTs(indexFile *os.File) {
return
}
v.lastModifiedTsSeconds = appendAtNs / uint64(time.Second)
if appendAtNs > v.lastWriteAppendAtNs {
v.lastWriteAppendAtNs = appendAtNs
v.lastWriteNeedleKey = key
v.lastWriteDeleted = false
}
}
// vacuumedLastWriteScanEntries bounds the work a vacuumed volume's recovery
@@ -285,25 +290,27 @@ var vacuumedLastWriteScanEntries = 1 << 16
// findLastWriteAppendAtNs scans the .idx backwards for the newest write -- an
// entry that is not a deletion tombstone -- and returns that needle's append
// timestamp. The .idx and the .dat share an order, so an append-ordered volume
// answers with the first write the scan reaches. Vacuum rewrites both in key
// order, which tracks write order only because the master issues keys
// increasing: an overwrite keeps its original, lower key, so a vacuumed volume
// has to take the maximum over every write it indexes. Returns 0 when the .idx
// holds nothing but tombstones, when a vacuumed volume holds more needles than
// the scan budget, or for a volume older than version 3, whose needles carry no
// append timestamp.
func findLastWriteAppendAtNs(v *Volume, indexFile *os.File, indexSize int64) (uint64, error) {
// timestamp and key. The .idx and the .dat share an order, so an
// append-ordered volume answers with the first write the scan reaches. Vacuum
// rewrites both in key order, which tracks write order only because the master
// issues keys increasing: an overwrite keeps its original, lower key, so a
// vacuumed volume has to take the maximum over every write it indexes. Returns
// 0 when the .idx holds nothing but tombstones, when a vacuumed volume holds
// more needles than the scan budget, or for a volume older than version 3,
// whose needles carry no append timestamp.
func findLastWriteAppendAtNs(v *Volume, indexFile *os.File, indexSize int64) (uint64, types.NeedleId, error) {
version := v.Version()
if version != needle.Version3 {
return 0, nil
return 0, 0, nil
}
scanEveryWrite := v.SuperBlock.CompactionRevision > 0
entryBudget := vacuumedLastWriteScanEntries
if scanEveryWrite && !affordableVacuumedScan(indexFile, indexSize, v.Id, v.FileName(".dat")) {
return 0, nil
return 0, 0, nil
}
var lastWriteAppendAtNs uint64
var lastWriteKey types.NeedleId
dead := make(map[types.NeedleId]struct{})
block := make([]byte, types.NeedleMapEntrySize*idx.RowsToRead)
for end := indexSize; end > 0; {
start := max(end-int64(len(block)), 0)
@@ -313,11 +320,17 @@ func findLastWriteAppendAtNs(v *Volume, indexFile *os.File, indexSize int64) (ui
err = nil
}
if err != nil {
return 0, fmt.Errorf("read %s at %d: %v", indexFile.Name(), start, err)
return 0, 0, fmt.Errorf("read %s at %d: %v", indexFile.Name(), start, err)
}
for i := len(entries) - types.NeedleMapEntrySize; i >= 0; i -= types.NeedleMapEntrySize {
key, offset, size := idx.IdxFileEntry(entries[i : i+types.NeedleMapEntrySize])
// The first row a key presents is its latest state: a tombstone
// there retires the write rows beneath it.
if _, gone := dead[key]; gone {
continue
}
if offset.IsZero() || size.IsDeleted() {
dead[key] = struct{}{}
continue
}
needleOffset := findNeedleOffset(v.DataBackend, version, offset.ToActualOffset(), key, size)
@@ -326,21 +339,24 @@ func findLastWriteAppendAtNs(v *Volume, indexFile *os.File, indexSize int64) (ui
}
appendAtNs, err := readNeedleAppendAtNs(v.DataBackend, needleOffset, size)
if err != nil {
return 0, err
return 0, 0, err
}
if appendAtNs > lastWriteAppendAtNs {
lastWriteAppendAtNs = appendAtNs
lastWriteKey = key
}
lastWriteAppendAtNs = max(lastWriteAppendAtNs, appendAtNs)
if !scanEveryWrite {
return lastWriteAppendAtNs, nil
return lastWriteAppendAtNs, lastWriteKey, nil
}
if entryBudget--; entryBudget == 0 {
glog.V(0).Infof("volume %d: more than %d needles to scan for its last write, keeping the %s mtime",
v.Id, vacuumedLastWriteScanEntries, v.FileName(".dat"))
return 0, nil
return 0, 0, nil
}
}
end = start
}
return lastWriteAppendAtNs, nil
return lastWriteAppendAtNs, lastWriteKey, nil
}
// affordableVacuumedScan reports whether a vacuumed volume's recovery scan fits
+7 -3
View File
@@ -172,7 +172,7 @@ func (v *Volume) load(alsoLoadIndex bool, createDatIfMissing bool, needleMapKind
return fmt.Errorf("load remote file %v: %w", v.volumeInfo, err)
}
// Set lastModifiedTsSeconds from remote file to prevent premature expiry on startup
if len(v.volumeInfo.GetFiles()) > 0 {
if len(v.volumeInfo.GetFiles()) > 0 && !v.keepLastModifiedTsOnLoad {
remoteFileModifiedTime := v.volumeInfo.GetFiles()[0].GetModifiedTime()
if remoteFileModifiedTime > 0 {
v.lastModifiedTsSeconds = remoteFileModifiedTime
@@ -201,7 +201,9 @@ func (v *Volume) load(alsoLoadIndex bool, createDatIfMissing bool, needleMapKind
if err != nil {
return datFileLoadError(v.FileName(".dat"), err)
}
v.lastModifiedTsSeconds = uint64(modifiedTime.Unix())
if !v.keepLastModifiedTsOnLoad {
v.lastModifiedTsSeconds = uint64(modifiedTime.Unix())
}
if fileSize >= super_block.SuperBlockSize {
alreadyHasSuperBlock = true
}
@@ -292,7 +294,9 @@ func (v *Volume) load(alsoLoadIndex bool, createDatIfMissing bool, needleMapKind
v.noWriteOrDelete = true
glog.V(0).Infof("volumeDataIntegrityChecking failed %v", err)
}
v.recoverLastModifiedTs(indexFile)
if !v.keepLastModifiedTsOnLoad {
v.recoverLastModifiedTs(indexFile)
}
}
// The post-load structural check below uses the in-memory needle map
+147 -1
View File
@@ -250,7 +250,7 @@ func TestVolumeTtlClockSkipsUnaffordableScanWithoutDatReads(t *testing.T) {
defer func(budget int) { vacuumedLastWriteScanEntries = budget }(vacuumedLastWriteScanEntries)
vacuumedLastWriteScanEntries = 2
appendAtNs, err := findLastWriteAppendAtNs(v, indexFile, indexStat.Size())
appendAtNs, _, err := findLastWriteAppendAtNs(v, indexFile, indexStat.Size())
if err != nil {
t.Fatalf("recover last write: %v", err)
}
@@ -262,6 +262,152 @@ func TestVolumeTtlClockSkipsUnaffordableScanWithoutDatReads(t *testing.T) {
}
}
// TestVolumeTtlClockCarriedAcrossVacuumCommit covers the reload that ends a
// vacuum commit: the clock must keep the last write's append time rather than
// be re-derived, and the intervening delete's tombstone must not freshen it.
func TestVolumeTtlClockCarriedAcrossVacuumCommit(t *testing.T) {
dir := t.TempDir()
ttl, err := needle.ReadTTL("5m")
if err != nil {
t.Fatalf("read ttl: %v", err)
}
v, err := NewVolume(dir, dir, "", 1, NeedleMapInMemory, &super_block.ReplicaPlacement{}, ttl, 0, needle.GetCurrentVersion(), 0, 0)
if err != nil {
t.Fatalf("volume creation: %v", err)
}
defer v.Close()
lastWriteNs := uint64(time.Now().Add(-2 * time.Hour).UnixNano())
for i := 1; i <= 3; i++ {
n := newRandomNeedle(uint64(i))
offset, _, _, err := v.writeNeedle2(n, true, false, false)
if err != nil {
t.Fatalf("write needle %d: %v", i, err)
}
backdateAppendAtNs(t, v, int64(offset), n.Size, lastWriteNs)
}
if _, err := v.doDeleteRequest(newEmptyNeedle(2)); err != nil {
t.Fatalf("delete needle 2: %v", err)
}
// Where a restart's recovery would have left the clock. The delete above
// pushed lastAppendAtNs to ~now; the commit must not consult it.
v.lastModifiedTsSeconds = lastWriteNs / uint64(time.Second)
v.lastWriteAppendAtNs = lastWriteNs
defer func(budget int) { vacuumedLastWriteScanEntries = budget }(vacuumedLastWriteScanEntries)
vacuumedLastWriteScanEntries = 1
if err := v.CompactByIndex(nil); err != nil {
t.Fatalf("compact: %v", err)
}
if err := v.CommitCompact(); err != nil {
t.Fatalf("commit compact: %v", err)
}
if v.SuperBlock.CompactionRevision == 0 {
t.Fatal("vacuum must bump CompactionRevision for this test to exercise the vacuumed path")
}
if got, want := v.lastModifiedTsSeconds, lastWriteNs/uint64(time.Second); got != want {
t.Errorf("TTL clock after commit is %d, want the last write at %d", got, want)
}
if !v.expired(v.ContentSize(), 1024*1024) {
t.Error("a TTL volume whose last write is 2h old must stay expired across a vacuum commit")
}
}
// TestVolumeTtlClockAtCommitUsesAppendTime covers a write whose client
// supplied modified time does not match when it was appended: the commit
// keeps the server-side append watermark the recovery scan would recompute.
func TestVolumeTtlClockAtCommitUsesAppendTime(t *testing.T) {
dir := t.TempDir()
ttl, err := needle.ReadTTL("5m")
if err != nil {
t.Fatalf("read ttl: %v", err)
}
v, err := NewVolume(dir, dir, "", 1, NeedleMapInMemory, &super_block.ReplicaPlacement{}, ttl, 0, needle.GetCurrentVersion(), 0, 0)
if err != nil {
t.Fatalf("volume creation: %v", err)
}
defer v.Close()
future := uint64(time.Now().Add(24 * time.Hour).Unix())
n := newRandomNeedle(1)
n.LastModified = future
if _, _, _, err := v.writeNeedle2(n, true, false, false); err != nil {
t.Fatalf("write needle: %v", err)
}
if v.lastModifiedTsSeconds < future {
t.Fatalf("clock %d did not follow the needle's modified time %d", v.lastModifiedTsSeconds, future)
}
appendWatermarkSec := v.lastWriteAppendAtNs / uint64(time.Second)
if err := v.CompactByIndex(nil); err != nil {
t.Fatalf("compact: %v", err)
}
if err := v.CommitCompact(); err != nil {
t.Fatalf("commit compact: %v", err)
}
if got := v.lastModifiedTsSeconds; got != appendWatermarkSec {
t.Errorf("TTL clock after commit is %d, want the append watermark %d", got, appendWatermarkSec)
}
}
// TestVolumeTtlClockAtCommitSkipsDeletedWrite covers a vacuum commit whose
// newest write was deleted first: the carried watermark belonged to that
// write, so the commit has to let the reload rescan and land on the newest
// surviving write rather than keep the volume alive on a deleted write's time.
func TestVolumeTtlClockAtCommitSkipsDeletedWrite(t *testing.T) {
dir := t.TempDir()
ttl, err := needle.ReadTTL("5m")
if err != nil {
t.Fatalf("read ttl: %v", err)
}
v, err := NewVolume(dir, dir, "", 1, NeedleMapInMemory, &super_block.ReplicaPlacement{}, ttl, 0, needle.GetCurrentVersion(), 0, 0)
if err != nil {
t.Fatalf("volume creation: %v", err)
}
defer v.Close()
oldWriteNs := uint64(time.Now().Add(-2 * time.Hour).UnixNano())
newWriteNs := uint64(time.Now().Add(-time.Hour).UnixNano())
for _, w := range []struct {
id uint64
ns uint64
}{{1, oldWriteNs}, {2, newWriteNs}} {
n := newRandomNeedle(w.id)
offset, _, _, err := v.writeNeedle2(n, true, false, false)
if err != nil {
t.Fatalf("write needle %d: %v", w.id, err)
}
backdateAppendAtNs(t, v, int64(offset), n.Size, w.ns)
}
// The delete lands inside the commit window: makeupDiff replays its
// tombstone into the new .idx behind the write row the copy carried.
if err := v.CompactByIndex(nil); err != nil {
t.Fatalf("compact: %v", err)
}
if _, err := v.doDeleteRequest(newEmptyNeedle(2)); err != nil {
t.Fatalf("delete needle 2: %v", err)
}
if !v.lastWriteDeleted {
t.Fatal("deleting the newest write must mark the watermark dead")
}
if err := v.CommitCompact(); err != nil {
t.Fatalf("commit compact: %v", err)
}
if got, want := v.lastModifiedTsSeconds, oldWriteNs/uint64(time.Second); got != want {
t.Errorf("TTL clock after commit is %d, want the surviving write at %d", got, want)
}
if v.lastWriteDeleted {
t.Error("the reload's rescan must reseed the watermark off the surviving write")
}
}
// TestVolumeExpireAtSecCountsFromLastWrite guards the destroy time an EC volume
// is reclaimed on (erasure_coding.EcVolume.IsTimeToDestroy). It was recomputed
// as now+TTL on every .vif write, so a read-only mark, a tier upload or an EC
+11 -1
View File
@@ -222,7 +222,17 @@ func (v *Volume) CommitCompact() error {
//time.Sleep(20 * time.Second)
glog.V(3).Infof("Loading volume %d commit file...", v.Id)
if e := v.load(true, false, v.needleMapKind, 0, v.Version()); e != nil {
// The write watermark already equals what recoverLastModifiedTs would
// rescan, so keep the clock instead of paying for the scan under the lock.
// If its write was itself deleted, the reload recovers the newest
// surviving write instead.
if v.lastWriteAppendAtNs != 0 {
v.lastModifiedTsSeconds = v.lastWriteAppendAtNs / uint64(time.Second)
}
v.keepLastModifiedTsOnLoad = v.lastModifiedTsSeconds != 0 && !v.lastWriteDeleted
e := v.load(true, false, v.needleMapKind, 0, v.Version())
v.keepLastModifiedTsOnLoad = false
if e != nil {
return e
}
glog.V(3).Infof("Finish committing volume %d", v.Id)
+20
View File
@@ -329,6 +329,10 @@ func (v *Volume) rollbackUnflushedWrite(n *needle.Needle, offset uint64, end int
}
}
}
if err == nil && n.Id == v.lastWriteNeedleKey {
v.lastWriteAppendAtNs, v.lastWriteNeedleKey = 0, 0
v.lastWriteDeleted = true
}
if err != nil {
recoveryErr = errors.Join(recoveryErr,
fmt.Errorf("roll back the index of needle %d in volume %d: %w", n.Id, v.Id, err))
@@ -409,6 +413,9 @@ func (v *Volume) doWriteRequest(n *needle.Needle, checkCookie bool, fsync bool)
return
}
v.lastAppendAtNs = n.AppendAtNs
v.lastWriteAppendAtNs = n.AppendAtNs
v.lastWriteNeedleKey = n.Id
v.lastWriteDeleted = false
// add to needle map
if !ok || uint64(nv.Offset.ToActualOffset()) < offset {
@@ -490,6 +497,10 @@ func (v *Volume) doDeleteRequest(n *needle.Needle) (Size, error) {
if err = v.nm.Delete(n.Id, ToOffset(int64(offset))); err != nil {
return size, err
}
if n.Id == v.lastWriteNeedleKey {
v.lastWriteAppendAtNs, v.lastWriteNeedleKey = 0, 0
v.lastWriteDeleted = true
}
return size, err
}
return 0, nil
@@ -524,6 +535,9 @@ func (v *Volume) processBatch(currentRequests []*needle.AsyncRequest) {
}
indexEnd := int64(v.nm.IndexFileSize())
batchLastAppendAtNs := v.lastAppendAtNs
batchLastWriteAppendAtNs := v.lastWriteAppendAtNs
batchLastWriteNeedleKey := v.lastWriteNeedleKey
batchLastWriteDeleted := v.lastWriteDeleted
batchLastModifiedTsSeconds := v.lastModifiedTsSeconds
for i := 0; i < len(currentRequests); i++ {
needleID := currentRequests[i].N.Id
@@ -563,6 +577,9 @@ func (v *Volume) processBatch(currentRequests []*needle.AsyncRequest) {
if syncErr := v.DataBackend.Sync(); syncErr != nil {
v.checkReadWriteError(syncErr)
v.lastAppendAtNs = batchLastAppendAtNs
v.lastWriteAppendAtNs = batchLastWriteAppendAtNs
v.lastWriteNeedleKey = batchLastWriteNeedleKey
v.lastWriteDeleted = batchLastWriteDeleted
v.lastModifiedTsSeconds = batchLastModifiedTsSeconds
batchErr := syncErr
if recoveryErr := v.rollbackBatch(end, indexEnd, orderedSnapshots, metricRollbacker, batchMetrics); recoveryErr != nil {
@@ -702,6 +719,9 @@ func (v *Volume) WriteNeedleBlob(needleId NeedleId, needleBlob []byte, size Size
return err
}
v.lastAppendAtNs = appendAtNs
v.lastWriteAppendAtNs = appendAtNs
v.lastWriteNeedleKey = needleId
v.lastWriteDeleted = false
// add to needle map
if err = v.nm.Put(needleId, ToOffset(int64(offset)), size); err != nil {
@@ -18,6 +18,7 @@ type Config struct {
MinSizeMB int `json:"min_size_mb"`
PreferredTags []string `json:"preferred_tags"`
ReplicaPlacement string `json:"replica_placement"` // e.g. "020"; empty falls back to the master default replication
StrictPlacement bool `json:"strict_placement"` // fail planning instead of relaxing placement constraints
}
// NewDefaultConfig creates a new default erasure coding configuration
@@ -171,6 +172,18 @@ func GetConfigSpec() base.ConfigSpec {
InputType: "text",
CSSClasses: "form-control",
},
{
Name: "strict_placement",
JSONName: "strict_placement",
Type: config.FieldTypeBool,
DefaultValue: false,
Required: false,
DisplayName: "Strict Placement",
Description: "Refuse to encode a volume when the placement constraints can't be satisfied",
HelpText: "When enabled, a volume is only encoded if every shard can be placed within the per-disk, anti-affinity, replica-placement and per-rack caps, so the configured resilience is preserved. When disabled (default), unsatisfiable constraints are relaxed and noted in the log",
InputType: "checkbox",
CSSClasses: "form-check-input",
},
},
}
}
@@ -192,6 +205,7 @@ func (c *Config) ToTaskPolicy() *worker_pb.TaskPolicy {
CollectionFilter: c.CollectionFilter,
PreferredTags: preferredTagsCopy,
ReplicaPlacement: c.ReplicaPlacement,
StrictPlacement: c.StrictPlacement,
},
},
}
@@ -216,6 +230,7 @@ func (c *Config) FromTaskPolicy(policy *worker_pb.TaskPolicy) error {
c.CollectionFilter = ecConfig.CollectionFilter
c.PreferredTags = append([]string(nil), ecConfig.PreferredTags...)
c.ReplicaPlacement = ecConfig.ReplicaPlacement
c.StrictPlacement = ecConfig.StrictPlacement
}
return nil
+10 -3
View File
@@ -487,8 +487,11 @@ func buildNodeAddressMap(at *topology.ActiveTopology) map[string]string {
//
// Encode is lenient (PlaceDurabilityFirst): it relaxes caps/anti-affinity/RP and,
// last, the total-shards-per-rack cap as needed, failing only when no eligible
// disk has room. It prefers the source disk type but spills if that type can't
// hold every shard. rp is the resolved replica placement (may be nil).
// disk has room. ecConfig.StrictPlacement instead fails the volume's planning
// when the constraints cannot be satisfied (PlaceStrict), so a volume is only
// encoded while its configured resilience is preserved. The task prefers the
// source disk type but spills if that type can't hold every shard. rp is the
// resolved replica placement (may be nil).
func planECDestinations(snap *ecbalancer.Topology, nodeAddresses map[string]string, metric *types.VolumeHealthMetrics, ecConfig *Config, rp *super_block.ReplicaPlacement, dataShards, parityShards int) (*topology.MultiDestinationPlan, [][]uint32, error) {
if snap == nil {
return nil, nil, fmt.Errorf("EC placement snapshot not available")
@@ -511,13 +514,17 @@ func planECDestinations(snap *ecbalancer.Topology, nodeAddresses map[string]stri
for i := range need {
need[i] = i
}
mode := ecbalancer.PlaceDurabilityFirst
if ecConfig.StrictPlacement {
mode = ecbalancer.PlaceStrict
}
res, err := snap.Place(metric.VolumeID, metric.Collection, need, ecbalancer.Constraints{
DiskType: metric.DiskType,
DiskTypePolicy: ecbalancer.DiskTypePrefer,
PreferredTags: ecConfig.PreferredTags,
ReplicaPlacement: rp,
Ratio: func(string) (int, int) { return dataShards, parityShards },
}, ecbalancer.PlaceDurabilityFirst)
}, mode)
if err != nil {
return nil, nil, err
}
@@ -0,0 +1,50 @@
package erasure_coding
import (
"testing"
"github.com/seaweedfs/seaweedfs/weed/storage/erasure_coding"
"github.com/seaweedfs/seaweedfs/weed/storage/erasure_coding/ecbalancer"
"github.com/seaweedfs/seaweedfs/weed/storage/super_block"
"github.com/seaweedfs/seaweedfs/weed/worker/types"
"github.com/stretchr/testify/require"
)
// StrictPlacement refuses a volume whose configured caps cannot be met,
// rather than relaxing them: here one rack must hold all 14 shards under
// a 2-shards-per-rack replica placement.
func TestPlanECDestinationsStrictPlacement(t *testing.T) {
activeTopology := buildActiveTopology(t, 7, []string{"hdd"}, 100, 0, "")
metric := &types.VolumeHealthMetrics{
VolumeID: 1,
Server: "10.0.0.1:8080",
Size: 100 * 1024 * 1024,
}
rp, err := super_block.NewReplicaPlacementFromString("020")
require.NoError(t, err)
nodeAddresses := buildNodeAddressMap(activeTopology)
snap := ecbalancer.FromActiveTopology(activeTopology, erasure_coding.DataShardsCount)
cfg := NewDefaultConfig()
plan, shardsPerPlan, err := planECDestinations(snap, nodeAddresses, metric, cfg, rp, erasure_coding.DataShardsCount, erasure_coding.ParityShardsCount)
require.NoError(t, err, "lenient placement relaxes the unsatisfiable rack cap")
requireAllShardsPlaced(t, plan, shardsPerPlan)
snap = ecbalancer.FromActiveTopology(activeTopology, erasure_coding.DataShardsCount)
cfg.StrictPlacement = true
_, _, err = planECDestinations(snap, nodeAddresses, metric, cfg, rp, erasure_coding.DataShardsCount, erasure_coding.ParityShardsCount)
require.Error(t, err, "strict placement must refuse rather than weaken the rack cap")
}
func TestStrictPlacementRoundTripsThroughTaskPolicy(t *testing.T) {
cfg := NewDefaultConfig()
cfg.StrictPlacement = true
restored := NewDefaultConfig()
require.NoError(t, restored.FromTaskPolicy(cfg.ToTaskPolicy()))
require.True(t, restored.StrictPlacement, "strict placement must survive the persisted policy round trip")
restored.StrictPlacement = false
require.NoError(t, restored.FromTaskPolicy(restored.ToTaskPolicy()))
require.False(t, restored.StrictPlacement)
}
@@ -146,6 +146,13 @@ func (h *ErasureCodingHandler) Descriptor() *plugin_pb.JobTypeDescriptor {
FieldType: plugin_pb.ConfigFieldType_CONFIG_FIELD_TYPE_STRING,
Widget: plugin_pb.ConfigWidget_CONFIG_WIDGET_TEXT,
},
{
Name: "strict_placement",
Label: "Strict Placement",
Description: "Fail EC planning when the placement constraints cannot be met instead of relaxing them.",
FieldType: plugin_pb.ConfigFieldType_CONFIG_FIELD_TYPE_BOOL,
Widget: plugin_pb.ConfigWidget_CONFIG_WIDGET_TOGGLE,
},
},
},
},
@@ -165,6 +172,9 @@ func (h *ErasureCodingHandler) Descriptor() *plugin_pb.JobTypeDescriptor {
"replica_placement": {
Kind: &plugin_pb.ConfigValue_StringValue{StringValue: ""},
},
"strict_placement": {
Kind: &plugin_pb.ConfigValue_BoolValue{BoolValue: false},
},
},
},
AdminRuntimeDefaults: &plugin_pb.AdminRuntimeDefaults{
@@ -624,6 +634,8 @@ func deriveErasureCodingWorkerConfig(values map[string]*plugin_pb.ConfigValue) *
taskConfig.ReplicaPlacement = strings.TrimSpace(pluginworker.ReadStringConfig(values, "replica_placement", taskConfig.ReplicaPlacement))
taskConfig.StrictPlacement = pluginworker.ReadBoolConfig(values, "strict_placement", taskConfig.StrictPlacement)
return &erasureCodingWorkerConfig{
TaskConfig: taskConfig,
}
+90 -7
View File
@@ -189,9 +189,11 @@ func (h *Handler) compactDataFiles(
return "", nil, err
}
// Build compaction bins: group small data files by partition.
// Build compaction bins: group small data files by partition, each bin's
// files in the order of their bounds so the merged file keeps the
// ordering its inputs had.
targetSize := compactionTargetSizeForPlan(config, rewritePlan)
bins := buildCompactionBins(candidateEntries, targetSize, minInputFiles)
bins := buildCompactionBins(candidateEntries, targetSize, minInputFiles, meta)
initialBinCount := len(bins)
bins = filterCompactionBinsByPlan(bins, config, rewritePlan)
if len(bins) == 0 {
@@ -269,10 +271,11 @@ func (h *Handler) compactDataFiles(
var mergedData []byte
var recordCount int64
rowGroupRows := rowsPerRowGroup(bin, config)
if rewritePlan != nil && rewritePlan.strategy == "sort" {
mergedData, recordCount, err = mergeParquetFilesSorted(ctx, filerClient, bucketName, dataPath, bin.Entries, positionDeletes, eqDeleteGroups, schema, rewritePlan)
mergedData, recordCount, err = mergeParquetFilesSorted(ctx, filerClient, bucketName, dataPath, bin.Entries, positionDeletes, eqDeleteGroups, schema, rewritePlan, rowGroupRows)
} else {
mergedData, recordCount, err = mergeParquetFiles(ctx, filerClient, bucketName, dataPath, bin.Entries, positionDeletes, eqDeleteGroups, schema)
mergedData, recordCount, err = mergeParquetFiles(ctx, filerClient, bucketName, dataPath, bin.Entries, positionDeletes, eqDeleteGroups, schema, rowGroupRows)
}
if err != nil {
glog.Warningf("iceberg compact: failed to merge bin %d (%d files): %v", binIdx, len(bin.Entries), err)
@@ -314,6 +317,14 @@ func (h *Handler) compactDataFiles(
}
writtenArtifacts = append(writtenArtifacts, artifact{dir: dataDir, fileName: mergedFileName})
// Record the column statistics readers prune files by; without
// them every scan has to read every compacted file.
if stats, statsErr := collectColumnStats(mergedData, schema); statsErr != nil {
glog.Warningf("iceberg compact: no column statistics for %s: %v", mergedFileName, statsErr)
} else {
stats.applyTo(dfBuilder)
}
mergedDataFile := dfBuilder.Build()
summary.addFile(mergedDataFile)
newEntry := iceberg.NewManifestEntry(
@@ -527,7 +538,13 @@ func (h *Handler) compactDataFiles(
// buildCompactionBins groups small data files by partition for bin-packing.
// A file is "small" if it's below targetSize. A bin must have at least
// minFiles entries to be worth compacting.
func buildCompactionBins(entries []iceberg.ManifestEntry, targetSize int64, minFiles int) []compactionBin {
//
// With an order, every bin lists its files by their bounds on the ordering
// column and an oversized partition is split into runs of consecutive files,
// so each merged file covers one contiguous range of that column. Without one
// the files keep manifest order and an oversized partition is packed
// largest-first.
func buildCompactionBins(entries []iceberg.ManifestEntry, targetSize int64, minFiles int, meta table.Metadata) []compactionBin {
if minFiles < 2 {
minFiles = 2
}
@@ -560,14 +577,37 @@ func buildCompactionBins(entries []iceberg.ManifestEntry, targetSize int64, minF
bin.TotalSize += df.FileSizeBytes()
}
// Filter to bins with enough files, splitting oversized bins
// Filter to bins with enough files, splitting oversized bins. The merge
// order resolves per group over the files that survived the eligibility
// filter, so a file that can never participate (oversized, another
// format) cannot veto ordering for the ones that can.
var result []compactionBin
for _, bin := range groups {
if len(bin.Entries) < minFiles {
continue
}
order := resolveCompactionOrder(meta, bin.Entries)
order.sortEntries(bin.Entries)
if bin.TotalSize <= targetSize {
result = append(result, *bin)
} else if order != nil {
runs, leftover := splitOrderedBin(*bin, targetSize, minFiles)
result = append(result, runs...)
if order.bestEffort && len(leftover) > 0 {
runtBin := compactionBin{PartitionKey: bin.PartitionKey, Partition: bin.Partition, SpecID: bin.SpecID, Entries: leftover}
packed := splitOversizedBin(runtBin, targetSize, minFiles)
if len(packed) == 0 {
// Stragglers cannot merge even with each other; repacking
// the whole bin lets them pair across the ordered runs.
// Only replace the runs when the repack yields bins —
// otherwise the runs stay valid compaction work.
if full := splitOversizedBin(*bin, targetSize, minFiles); len(full) > 0 {
result = result[:len(result)-len(runs)]
packed = full
}
}
result = append(result, packed...)
}
} else {
result = append(result, splitOversizedBin(*bin, targetSize, minFiles)...)
}
@@ -918,6 +958,46 @@ func resolveEqualityColIndices(pqSchema *parquet.Schema, fieldIDs []int, iceberg
return indices, nil
}
const (
// Iceberg's write.parquet.row-group-size-bytes default and PyIceberg's
// write.parquet.row-group-limit default.
defaultRowGroupSizeBytes = 128 * 1024 * 1024
defaultRowGroupRowLimit = 1048576
minRowGroupRows = 1024
)
// rowsPerRowGroup is the most rows one row group of a compacted file holds:
// the row limit, lowered so a row group stays near the byte size, estimated
// from the bin's inputs' bytes per row. parquet-go's writers have no limit of
// their own, so without one every compacted file is a single row group. The
// floor only protects the estimate — a configured row limit always wins.
func rowsPerRowGroup(bin compactionBin, config Config) int64 {
rows, sizeBytes := config.RowGroupRowLimit, config.RowGroupSizeBytes
explicit := rows > 0
if !explicit {
rows = defaultRowGroupRowLimit
}
if sizeBytes <= 0 {
sizeBytes = defaultRowGroupSizeBytes
}
var totalRows int64
for _, entry := range bin.Entries {
totalRows += entry.DataFile().Count()
}
if totalRows > 0 && bin.TotalSize > 0 {
bytesPerRow := max(bin.TotalSize/totalRows, 1)
estimated := sizeBytes / bytesPerRow
if !explicit {
estimated = max(estimated, minRowGroupRows)
}
rows = min(rows, estimated)
}
if !explicit {
rows = max(rows, minRowGroupRows)
}
return rows
}
// mergeParquetFiles reads multiple small Parquet files and merges them into
// a single Parquet file, optionally filtering out rows matching position or
// equality deletes. Files are processed one at a time to keep memory usage
@@ -930,6 +1010,7 @@ func mergeParquetFiles(
positionDeletes map[string][]int64,
eqDeleteGroups []equalityDeleteGroup,
icebergSchema *iceberg.Schema,
rowGroupRows int64,
) ([]byte, int64, error) {
if len(entries) == 0 {
return nil, 0, fmt.Errorf("no entries to merge")
@@ -954,7 +1035,7 @@ func mergeParquetFiles(
}
var outputBuf bytes.Buffer
writer := parquet.NewWriter(&outputBuf, parquetSchema)
writer := parquet.NewWriter(&outputBuf, parquetSchema, parquet.MaxRowsPerRowGroup(rowGroupRows))
drainReader := func(reader *parquet.Reader, source string) (int64, error) {
return visitFilteredParquetRows(ctx, reader, source, bucketName, dataPath, positionDeletes, resolvedEqGroups, func(filtered []parquet.Row) error {
@@ -1119,6 +1200,7 @@ func mergeParquetFilesSorted(
eqDeleteGroups []equalityDeleteGroup,
icebergSchema *iceberg.Schema,
rewritePlan *compactionRewritePlan,
rowGroupRows int64,
) ([]byte, int64, error) {
if len(entries) == 0 {
return nil, 0, fmt.Errorf("no entries to merge")
@@ -1168,6 +1250,7 @@ func mergeParquetFilesSorted(
var outputBuf bytes.Buffer
writer := parquet.NewSortingWriter[any](&outputBuf, sortBufferRows, parquetSchema,
parquet.MaxRowsPerRowGroup(rowGroupRows),
parquet.SortingWriterConfig(
parquet.SortingColumns(sortingColumns...),
parquet.SortingBuffers(parquet.NewFileBufferPool(spillDir, "seaweedfs-iceberg-sort-*")),
+201
View File
@@ -0,0 +1,201 @@
package iceberg
import (
"sort"
"github.com/apache/iceberg-go"
"github.com/apache/iceberg-go/table"
)
// compactionOrder names the column by whose bounds a bin's input files are
// ordered before they are merged: the first identity field of the table's
// sort order when it declares one, otherwise the first column every candidate
// file carries bounds for. Rows are copied file by file, so a merged file is
// as ordered as the sequence of its inputs.
type compactionOrder struct {
fieldID int
typ iceberg.PrimitiveType
descending bool
// bestEffort marks an order inferred from column bounds rather than the
// table's declared sort order; runs it cannot form may still merge by
// size since contiguous ranges are a preference, not a requirement.
bestEffort bool
}
func resolveCompactionOrder(meta table.Metadata, entries []iceberg.ManifestEntry) *compactionOrder {
if meta == nil || meta.CurrentSchema() == nil || len(entries) == 0 {
return nil
}
schema := meta.CurrentSchema()
if sortOrder := meta.SortOrder(); !sortOrder.IsUnsorted() {
for _, sortField := range sortOrder.Fields() {
if _, ok := sortField.Transform.(iceberg.IdentityTransform); !ok {
continue
}
field, ok := schema.FindFieldByID(sortField.SourceID())
if !ok {
continue
}
typ, ok := field.Type.(iceberg.PrimitiveType)
if !ok {
continue
}
candidate := &compactionOrder{fieldID: field.ID, typ: typ, descending: sortField.Direction == table.SortDESC}
// The declared order only holds when every entry carries a bound
// for it; otherwise merging would claim sortedness it cannot
// verify. Fall back to inference instead of ordering by a later
// sort field alone.
if candidate.coversEntries(entries) {
return candidate
}
break
}
}
// A column only orders the merge when every entry in the group carries a
// bound for it. The check is scoped to the entries handed in — a bin's
// eligible files — so an unrelated file without bounds (oversized, a
// different format, or another partition) cannot disable ordering here.
for _, field := range schema.Fields() {
typ, ok := field.Type.(iceberg.PrimitiveType)
if !ok {
continue
}
candidate := &compactionOrder{fieldID: field.ID, typ: typ, bestEffort: true}
if candidate.coversEntries(entries) {
return candidate
}
}
return nil
}
func (o *compactionOrder) coversEntries(entries []iceberg.ManifestEntry) bool {
for _, entry := range entries {
if _, ok := o.key(entry); !ok {
return false
}
}
return true
}
// key is the bound a file is ordered by: its lower bound on the column for an
// ascending order, its upper bound for a descending one.
func (o *compactionOrder) key(entry iceberg.ManifestEntry) (iceberg.Literal, bool) {
df := entry.DataFile()
raw := df.LowerBoundValues()[o.fieldID]
if o.descending {
raw = df.UpperBoundValues()[o.fieldID]
}
if raw == nil {
return nil, false
}
lit, err := iceberg.LiteralFromBytes(o.typ, raw)
return lit, err == nil
}
// sortEntries orders a bin's files by their key. Files without a key sort
// after the rest, in their original order.
func (o *compactionOrder) sortEntries(entries []iceberg.ManifestEntry) {
if o == nil || len(entries) < 2 {
return
}
type keyed struct {
key iceberg.Literal
ok bool
}
keys := make([]keyed, len(entries))
idx := make([]int, len(entries))
for i, entry := range entries {
k, ok := o.key(entry)
keys[i], idx[i] = keyed{k, ok}, i
}
sort.SliceStable(idx, func(i, j int) bool {
a, b := keys[idx[i]], keys[idx[j]]
if a.ok != b.ok {
return a.ok
}
if !a.ok {
return false
}
c, ok := compareLiterals(a.key, b.key)
if o.descending {
c = -c
}
return ok && c < 0
})
sorted := make([]iceberg.ManifestEntry, len(entries))
for i, j := range idx {
sorted[i] = entries[j]
}
copy(entries, sorted)
}
// splitOrderedBin splits a bin whose files are in bound order into runs of
// consecutive files that stay under targetSize, so every output covers one
// contiguous range of the ordering column. The returned entries are the runs
// too short to reach minFiles: the caller can still merge them by size when
// contiguous ranges are only a preference, or leave them for a later pass
// when order must hold.
func splitOrderedBin(bin compactionBin, targetSize int64, minFiles int) (valid []compactionBin, leftover []iceberg.ManifestEntry) {
newBin := func() compactionBin {
return compactionBin{PartitionKey: bin.PartitionKey, Partition: bin.Partition, SpecID: bin.SpecID}
}
current := newBin()
flush := func() {
if len(current.Entries) >= minFiles {
valid = append(valid, current)
} else {
leftover = append(leftover, current.Entries...)
}
current = newBin()
}
for _, entry := range bin.Entries {
size := entry.DataFile().FileSizeBytes()
if current.TotalSize > 0 && current.TotalSize+size > targetSize {
flush()
}
current.Entries = append(current.Entries, entry)
current.TotalSize += size
}
flush()
return valid, leftover
}
// compareLiterals orders two literals of the same type; false when they are
// not comparable.
func compareLiterals(a, b iceberg.Literal) (int, bool) {
switch x := a.(type) {
case iceberg.TypedLiteral[bool]:
return compareTyped(x, b)
case iceberg.TypedLiteral[int32]:
return compareTyped(x, b)
case iceberg.TypedLiteral[int64]:
return compareTyped(x, b)
case iceberg.TypedLiteral[float32]:
return compareTyped(x, b)
case iceberg.TypedLiteral[float64]:
return compareTyped(x, b)
case iceberg.TypedLiteral[iceberg.Date]:
return compareTyped(x, b)
case iceberg.TypedLiteral[iceberg.Time]:
return compareTyped(x, b)
case iceberg.TypedLiteral[iceberg.Timestamp]:
return compareTyped(x, b)
case iceberg.TypedLiteral[iceberg.TimestampNano]:
return compareTyped(x, b)
case iceberg.TypedLiteral[string]:
return compareTyped(x, b)
case iceberg.TypedLiteral[[]byte]:
return compareTyped(x, b)
case iceberg.TypedLiteral[iceberg.Decimal]:
return compareTyped(x, b)
}
return 0, false
}
func compareTyped[T iceberg.LiteralType](a iceberg.TypedLiteral[T], b iceberg.Literal) (int, bool) {
bt, ok := b.(iceberg.TypedLiteral[T])
if !ok {
return 0, false
}
return a.Comparator()(a.Value(), bt.Value()), true
}
@@ -0,0 +1,388 @@
package iceberg
import (
"bytes"
"context"
"encoding/json"
"fmt"
"path"
"strings"
"testing"
"time"
"github.com/apache/iceberg-go"
"github.com/apache/iceberg-go/table"
filer_pb "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
"github.com/seaweedfs/seaweedfs/weed/s3api/s3tables"
)
type boundedFile struct {
Name string
IDs []int64
Lower, Upper int64
SizeBytes int64
NoBounds bool
}
// populateBoundedTable is populateTableWithDeleteFiles narrowed to data
// files whose manifest entries carry bounds on the id column.
func populateBoundedTable(t *testing.T, fs *fakeFilerServer, setup tableSetup, files []boundedFile) {
t.Helper()
populateBoundedTableSorted(t, fs, setup, files, table.UnsortedSortOrder)
}
func populateBoundedTableSorted(t *testing.T, fs *fakeFilerServer, setup tableSetup, files []boundedFile, sortOrder table.SortOrder) {
t.Helper()
schema := newTestSchema()
spec := *iceberg.UnpartitionedSpec
meta, err := table.NewMetadata(schema, &spec, sortOrder, "s3://"+setup.BucketName+"/"+setup.dataPath(), nil)
if err != nil {
t.Fatalf("create metadata: %v", err)
}
bucketPath := path.Join(s3tables.TablesPath, setup.BucketName)
nsPath := path.Join(bucketPath, setup.Namespace)
tableFilerPath := path.Join(bucketPath, setup.dataPath())
metaDir := path.Join(tableFilerPath, "metadata")
dataDir := path.Join(tableFilerPath, "data")
version := meta.Version()
var entries []iceberg.ManifestEntry
for _, f := range files {
rows := make([]struct {
ID int64
Name string
}, len(f.IDs))
for i, id := range f.IDs {
rows[i] = struct {
ID int64
Name string
}{id, fmt.Sprintf("r%d", id)}
}
data := writeTestParquetFile(t, fs, dataDir, f.Name, rows)
size := int64(len(data))
if f.SizeBytes > 0 {
size = f.SizeBytes
}
dfb, err := iceberg.NewDataFileBuilder(spec, iceberg.EntryContentData, setup.fileRef("data", f.Name),
iceberg.ParquetFile, map[int]any{}, nil, nil, int64(len(f.IDs)), size)
if err != nil {
t.Fatalf("build data file %s: %v", f.Name, err)
}
if !f.NoBounds {
lo, _ := iceberg.Int64Literal(f.Lower).MarshalBinary()
hi, _ := iceberg.Int64Literal(f.Upper).MarshalBinary()
dfb.LowerBoundValues(map[int][]byte{1: lo}).UpperBoundValues(map[int][]byte{1: hi})
}
snapID := int64(1)
entries = append(entries, iceberg.NewManifestEntry(iceberg.EntryStatusADDED, &snapID, nil, nil, dfb.Build()))
}
var manifestBuf bytes.Buffer
mf, err := iceberg.WriteManifest(setup.fileRef("metadata", "data-manifest-1.avro"), &manifestBuf,
version, spec, schema, 1, entries)
if err != nil {
t.Fatalf("write manifest: %v", err)
}
fs.putEntry(metaDir, "data-manifest-1.avro", &filer_pb.Entry{
Name: "data-manifest-1.avro", Content: manifestBuf.Bytes(),
Attributes: &filer_pb.FuseAttributes{Mtime: time.Now().Unix(), FileSize: uint64(manifestBuf.Len())},
})
var mlBuf bytes.Buffer
seqNum := int64(1)
if err := iceberg.WriteManifestList(version, &mlBuf, 1, nil, &seqNum, 0, []iceberg.ManifestFile{mf}); err != nil {
t.Fatalf("write manifest list: %v", err)
}
fs.putEntry(metaDir, "snap-1.avro", &filer_pb.Entry{
Name: "snap-1.avro", Content: mlBuf.Bytes(),
Attributes: &filer_pb.FuseAttributes{Mtime: time.Now().Unix(), FileSize: uint64(mlBuf.Len())},
})
snap := table.Snapshot{SnapshotID: 1, TimestampMs: time.Now().UnixMilli(),
ManifestList: setup.fileRef("metadata", "snap-1.avro")}
builder, err := table.MetadataBuilderFromBase(meta, "s3://"+setup.BucketName+"/"+setup.dataPath())
if err != nil {
t.Fatalf("metadata builder: %v", err)
}
if err := builder.AddSnapshot(&snap); err != nil {
t.Fatalf("add snapshot: %v", err)
}
if err := builder.SetSnapshotRef(table.MainBranch, snap.SnapshotID, table.BranchRef); err != nil {
t.Fatalf("set snapshot ref: %v", err)
}
meta, err = builder.Build()
if err != nil {
t.Fatalf("build metadata: %v", err)
}
fullMetadataJSON, _ := json.Marshal(meta)
xattr, _ := json.Marshal(map[string]interface{}{
"metadataVersion": 1,
"metadataLocation": setup.fileRef("metadata", "v1.metadata.json"),
"metadata": map[string]interface{}{"fullMetadata": json.RawMessage(fullMetadataJSON)},
})
fs.putEntry(path.Join(s3tables.TablesPath), setup.BucketName, &filer_pb.Entry{
Name: setup.BucketName, IsDirectory: true,
Extended: map[string][]byte{s3tables.ExtendedKeyTableBucket: []byte("true")},
})
fs.putEntry(bucketPath, setup.Namespace, &filer_pb.Entry{Name: setup.Namespace, IsDirectory: true})
fs.putEntry(nsPath, setup.TableName, &filer_pb.Entry{
Name: setup.TableName, IsDirectory: true,
Extended: map[string][]byte{s3tables.ExtendedKeyMetadata: xattr},
})
}
// A bin's files are merged in the order of their bounds on the
// ordering column, not manifest order, so a merged file is as ordered
// as the sequence of its inputs.
func TestCompactDataFilesKeepsInputOrderByBounds(t *testing.T) {
fs, client := startFakeFiler(t)
setup := tableSetup{BucketName: "tb", Namespace: "ns", TableName: "tbl"}
// manifest order is the reverse of bound order
populateBoundedTable(t, fs, setup, []boundedFile{
{Name: "hi.parquet", IDs: []int64{3, 4}, Lower: 3, Upper: 4},
{Name: "lo.parquet", IDs: []int64{1, 2}, Lower: 1, Upper: 2},
})
handler := NewHandler(nil)
config := Config{TargetFileSizeBytes: 256 * 1024 * 1024, MinInputFiles: 2, MaxCommitRetries: 3}
if _, _, err := handler.compactDataFiles(context.Background(), client, setup.BucketName, setup.tablePath(), config, nil); err != nil {
t.Fatalf("compactDataFiles: %v", err)
}
df := compactedDataFile(t, client, setup)
ids := compactedIDs(t, client, setup, df.FilePath())
want := []int64{1, 2, 3, 4}
if fmt.Sprint(ids) != fmt.Sprint(want) {
t.Errorf("merged rows = %v, want %v", ids, want)
}
}
// An oversized partition is split into runs of consecutive files in
// bound order, so each output covers one contiguous range.
func TestCompactDataFilesSplitsOrderedRuns(t *testing.T) {
fs, client := startFakeFiler(t)
setup := tableSetup{BucketName: "tb", Namespace: "ns", TableName: "tbl"}
// four disjoint ranges, each "4 MB": one bin, split into two runs of two
populateBoundedTable(t, fs, setup, []boundedFile{
{Name: "q4.parquet", IDs: []int64{7, 8}, Lower: 7, Upper: 8, SizeBytes: 4 << 20},
{Name: "q1.parquet", IDs: []int64{1, 2}, Lower: 1, Upper: 2, SizeBytes: 4 << 20},
{Name: "q3.parquet", IDs: []int64{5, 6}, Lower: 5, Upper: 6, SizeBytes: 4 << 20},
{Name: "q2.parquet", IDs: []int64{3, 4}, Lower: 3, Upper: 4, SizeBytes: 4 << 20},
})
handler := NewHandler(nil)
config := Config{TargetFileSizeBytes: 8 << 20, MinInputFiles: 2, MaxCommitRetries: 3}
result, _, err := handler.compactDataFiles(context.Background(), client, setup.BucketName, setup.tablePath(), config, nil)
if err != nil {
t.Fatalf("compactDataFiles: %v", err)
}
if !strings.Contains(result, "compacted 4 files into 2") {
t.Fatalf("expected two ordered runs, got %q", result)
}
// each output must cover one contiguous range
state, err := loadCurrentMetadata(context.Background(), client, setup.BucketName, setup.tablePath())
if err != nil {
t.Fatalf("loadCurrentMetadata: %v", err)
}
manifests, err := loadCurrentManifests(context.Background(), client, setup.BucketName, state.DataPath, state.Metadata)
if err != nil {
t.Fatalf("loadCurrentManifests: %v", err)
}
var ranges [][2]int64
for _, mf := range manifests {
if mf.ManifestContent() != iceberg.ManifestContentData {
continue
}
manifestData, err := loadFileByIcebergPath(context.Background(), client, setup.BucketName, state.DataPath, mf.FilePath())
if err != nil {
t.Fatalf("load manifest: %v", err)
}
entries, err := iceberg.ReadManifest(mf, bytes.NewReader(manifestData), true)
if err != nil {
t.Fatalf("read manifest: %v", err)
}
for _, entry := range entries {
if entry.Status() == iceberg.EntryStatusDELETED {
continue
}
lo, _ := iceberg.LiteralFromBytes(iceberg.PrimitiveTypes.Int64, entry.DataFile().LowerBoundValues()[1])
hi, _ := iceberg.LiteralFromBytes(iceberg.PrimitiveTypes.Int64, entry.DataFile().UpperBoundValues()[1])
ranges = append(ranges, [2]int64{int64(lo.(iceberg.Int64Literal)), int64(hi.(iceberg.Int64Literal))})
}
}
if len(ranges) != 2 {
t.Fatalf("expected 2 output files, got %d", len(ranges))
}
// ranges must be contiguous runs, not interleaved: {1..4} and {5..8}
for _, r := range ranges {
if r[1]-r[0] != 3 {
t.Errorf("output covers non-contiguous range %v", r)
}
}
}
// An oversized file without bounds can never join a merge, so it must not
// disable bounds ordering for the files that can.
func TestCompactDataFilesOrdersPastIneligibleFile(t *testing.T) {
fs, client := startFakeFiler(t)
setup := tableSetup{BucketName: "tb", Namespace: "ns", TableName: "tbl"}
populateBoundedTable(t, fs, setup, []boundedFile{
{Name: "hi.parquet", IDs: []int64{3, 4}, Lower: 3, Upper: 4},
{Name: "lo.parquet", IDs: []int64{1, 2}, Lower: 1, Upper: 2},
{Name: "huge.parquet", IDs: []int64{9}, SizeBytes: 512 << 20, NoBounds: true},
})
handler := NewHandler(nil)
config := Config{TargetFileSizeBytes: 256 << 20, MinInputFiles: 2, MaxCommitRetries: 3}
result, _, err := handler.compactDataFiles(context.Background(), client, setup.BucketName, setup.tablePath(), config, nil)
if err != nil {
t.Fatalf("compactDataFiles: %v", err)
}
if !strings.Contains(result, "compacted 2 files into 1") {
t.Fatalf("expected only the two eligible files to merge, got %q", result)
}
var merged iceberg.DataFile
for _, df := range liveDataFiles(t, client, setup) {
if !strings.HasSuffix(df.FilePath(), "huge.parquet") {
merged = df
}
}
ids := compactedIDs(t, client, setup, merged.FilePath())
want := []int64{1, 2, 3, 4}
if fmt.Sprint(ids) != fmt.Sprint(want) {
t.Errorf("merged rows = %v, want %v", ids, want)
}
}
// Runs too short to merge in bound order still merge by size, so files that
// cannot form a contiguous run are not stranded for every later pass.
func TestCompactDataFilesMergesStrandedRun(t *testing.T) {
fs, client := startFakeFiler(t)
setup := tableSetup{BucketName: "tb", Namespace: "ns", TableName: "tbl"}
populateBoundedTable(t, fs, setup, []boundedFile{
{Name: "r1.parquet", IDs: []int64{1}, Lower: 1, Upper: 2, SizeBytes: 4 << 20},
{Name: "r2.parquet", IDs: []int64{3}, Lower: 3, Upper: 4, SizeBytes: 4 << 20},
{Name: "r3.parquet", IDs: []int64{5}, Lower: 5, Upper: 6, SizeBytes: 6 << 20},
{Name: "r4.parquet", IDs: []int64{7}, Lower: 7, Upper: 8, SizeBytes: 3 << 20},
{Name: "r5.parquet", IDs: []int64{9}, Lower: 9, Upper: 10, SizeBytes: 3 << 20},
})
handler := NewHandler(nil)
config := Config{TargetFileSizeBytes: 10 << 20, MinInputFiles: 3, MaxCommitRetries: 3}
result, _, err := handler.compactDataFiles(context.Background(), client, setup.BucketName, setup.tablePath(), config, nil)
if err != nil {
t.Fatalf("compactDataFiles: %v", err)
}
if !strings.Contains(result, "compacted 3 files into 1") {
t.Fatalf("expected stranded files to merge by size, got %q", result)
}
}
// A declared sort order is preserved at the cost of compaction progress:
// runs too short to merge are left for later passes instead of being
// repacked out of order.
func TestCompactDataFilesKeepsStrandedRunForSortedTable(t *testing.T) {
fs, client := startFakeFiler(t)
setup := tableSetup{BucketName: "tb", Namespace: "ns", TableName: "tbl"}
sortOrder, err := table.NewSortOrder(1, []table.SortField{{
SourceIDs: []int{1},
Transform: iceberg.IdentityTransform{},
Direction: table.SortASC,
NullOrder: table.NullsFirst,
}})
if err != nil {
t.Fatalf("new sort order: %v", err)
}
populateBoundedTableSorted(t, fs, setup, []boundedFile{
{Name: "r1.parquet", IDs: []int64{1}, Lower: 1, Upper: 2, SizeBytes: 4 << 20},
{Name: "r2.parquet", IDs: []int64{3}, Lower: 3, Upper: 4, SizeBytes: 4 << 20},
{Name: "r3.parquet", IDs: []int64{5}, Lower: 5, Upper: 6, SizeBytes: 6 << 20},
{Name: "r4.parquet", IDs: []int64{7}, Lower: 7, Upper: 8, SizeBytes: 3 << 20},
{Name: "r5.parquet", IDs: []int64{9}, Lower: 9, Upper: 10, SizeBytes: 3 << 20},
}, sortOrder)
handler := NewHandler(nil)
config := Config{TargetFileSizeBytes: 10 << 20, MinInputFiles: 3, MaxCommitRetries: 3}
result, _, err := handler.compactDataFiles(context.Background(), client, setup.BucketName, setup.tablePath(), config, nil)
if err != nil {
t.Fatalf("compactDataFiles: %v", err)
}
if !strings.Contains(result, "no files eligible for compaction") {
t.Fatalf("expected underfilled ordered runs to be kept, got %q", result)
}
}
// A declared sort order only applies when every file carries a bound for its
// column; a file without one falls back to inference or unordered packing.
func TestResolveCompactionOrderSkipsUncoveredSortField(t *testing.T) {
schema := newTestSchema()
spec := *iceberg.UnpartitionedSpec
sortOrder, err := table.NewSortOrder(1, []table.SortField{{
SourceIDs: []int{1},
Transform: iceberg.IdentityTransform{},
Direction: table.SortASC,
NullOrder: table.NullsFirst,
}})
if err != nil {
t.Fatalf("new sort order: %v", err)
}
meta, err := table.NewMetadata(schema, &spec, sortOrder, "s3://b/p", nil)
if err != nil {
t.Fatalf("create metadata: %v", err)
}
entry := func(name string, bounded bool) iceberg.ManifestEntry {
dfb, err := iceberg.NewDataFileBuilder(spec, iceberg.EntryContentData,
"s3://b/data/"+name, iceberg.ParquetFile, map[int]any{}, nil, nil, 2, 16)
if err != nil {
t.Fatalf("build data file: %v", err)
}
if bounded {
lo, _ := iceberg.Int64Literal(1).MarshalBinary()
hi, _ := iceberg.Int64Literal(2).MarshalBinary()
dfb.LowerBoundValues(map[int][]byte{1: lo}).UpperBoundValues(map[int][]byte{1: hi})
}
snapID := int64(1)
return iceberg.NewManifestEntry(iceberg.EntryStatusADDED, &snapID, nil, nil, dfb.Build())
}
if order := resolveCompactionOrder(meta, []iceberg.ManifestEntry{
entry("a.parquet", true), entry("b.parquet", true),
}); order == nil || order.bestEffort {
t.Fatalf("covered sort field: got %+v, want strict order", order)
}
if order := resolveCompactionOrder(meta, []iceberg.ManifestEntry{
entry("a.parquet", true), entry("b.parquet", false),
}); order != nil {
t.Fatalf("uncovered sort field: got %+v, want nil", order)
}
}
// A full-bin repack that produces nothing must not discard the ordered runs
// it was meant to replace — they are still valid compaction work.
func TestCompactDataFilesKeepsRunsWhenRepackFails(t *testing.T) {
fs, client := startFakeFiler(t)
setup := tableSetup{BucketName: "tb", Namespace: "ns", TableName: "tbl"}
populateBoundedTable(t, fs, setup, []boundedFile{
{Name: "r1.parquet", IDs: []int64{1}, Lower: 1, Upper: 2, SizeBytes: 4 << 20},
{Name: "r2.parquet", IDs: []int64{3}, Lower: 3, Upper: 4, SizeBytes: 3 << 20},
{Name: "r3.parquet", IDs: []int64{5}, Lower: 5, Upper: 6, SizeBytes: 3 << 20},
{Name: "r4.parquet", IDs: []int64{7}, Lower: 7, Upper: 8, SizeBytes: 6 << 20},
{Name: "r5.parquet", IDs: []int64{9}, Lower: 9, Upper: 10, SizeBytes: 6 << 20},
})
handler := NewHandler(nil)
config := Config{TargetFileSizeBytes: 10 << 20, MinInputFiles: 3, MaxCommitRetries: 3}
result, _, err := handler.compactDataFiles(context.Background(), client, setup.BucketName, setup.tablePath(), config, nil)
if err != nil {
t.Fatalf("compactDataFiles: %v", err)
}
if !strings.Contains(result, "compacted 3 files into 1") {
t.Fatalf("expected the ordered run to survive the failed repack, got %q", result)
}
}
@@ -0,0 +1,97 @@
package iceberg
import (
"bytes"
"context"
"fmt"
"testing"
"github.com/apache/iceberg-go"
"github.com/parquet-go/parquet-go"
)
// Compacted files get row groups from the table's
// write.parquet.row-group-limit instead of one row group per file.
func TestCompactDataFilesWritesSeveralRowGroups(t *testing.T) {
fs, client := startFakeFiler(t)
setup := tableSetup{BucketName: "tb", Namespace: "ns", TableName: "tbl"}
type rowsT = struct {
ID int64
Name string
}
makeRows := func(lo, hi int64) []rowsT {
rows := make([]rowsT, 0, hi-lo)
for i := lo; i < hi; i++ {
rows = append(rows, rowsT{i, fmt.Sprintf("r%d", i)})
}
return rows
}
populateTableWithDeleteFiles(t, fs, setup,
[]struct {
Name string
Rows []rowsT
}{
{"d1.parquet", makeRows(0, 1500)},
{"d2.parquet", makeRows(1500, 3000)},
},
nil, nil,
)
handler := NewHandler(nil)
config := Config{
TargetFileSizeBytes: 256 * 1024 * 1024,
MinInputFiles: 2,
MaxCommitRetries: 3,
RowGroupRowLimit: 2000,
}
if _, _, err := handler.compactDataFiles(context.Background(), client, setup.BucketName, setup.tablePath(), config, nil); err != nil {
t.Fatalf("compactDataFiles: %v", err)
}
state, err := loadCurrentMetadata(context.Background(), client, setup.BucketName, setup.tablePath())
if err != nil {
t.Fatalf("loadCurrentMetadata: %v", err)
}
df := compactedDataFile(t, client, setup)
data, err := loadFileByIcebergPath(context.Background(), client, setup.BucketName, state.DataPath, df.FilePath())
if err != nil {
t.Fatalf("load compacted file: %v", err)
}
f, err := parquet.OpenFile(bytes.NewReader(data), int64(len(data)))
if err != nil {
t.Fatalf("open compacted file: %v", err)
}
if got := len(f.Metadata().RowGroups); got != 2 {
t.Errorf("row groups = %d, want 2 (3000 rows at row-group-limit 2000)", got)
}
}
// A configured row-group limit is a cap, not a floor: small explicit limits
// are honored instead of being raised to the estimate floor.
func TestRowsPerRowGroupHonorsConfiguredLimit(t *testing.T) {
entries := make([]iceberg.ManifestEntry, 0, 2)
spec := *iceberg.UnpartitionedSpec
for _, name := range []string{"a.parquet", "b.parquet"} {
dfb, err := iceberg.NewDataFileBuilder(spec, iceberg.EntryContentData,
"s3://b/data/"+name, iceberg.ParquetFile, map[int]any{}, nil, nil, 300, 9000)
if err != nil {
t.Fatalf("build data file: %v", err)
}
snapID := int64(1)
entries = append(entries, iceberg.NewManifestEntry(iceberg.EntryStatusADDED, &snapID, nil, nil, dfb.Build()))
}
bin := compactionBin{Entries: entries, TotalSize: 18000}
if got := rowsPerRowGroup(bin, Config{RowGroupRowLimit: 128}); got != 128 {
t.Errorf("explicit limit: rowsPerRowGroup = %d, want 128", got)
}
// 30 bytes/row, byte cap 128MB: the estimate of ~4.4M rows exceeds the 2M
// default, so the default holds; the floor lifts a degenerate estimate.
if got := rowsPerRowGroup(bin, Config{RowGroupSizeBytes: 1024}); got != minRowGroupRows {
t.Errorf("floored estimate: rowsPerRowGroup = %d, want %d", got, minRowGroupRows)
}
if got := rowsPerRowGroup(bin, Config{}); got != defaultRowGroupRowLimit {
t.Errorf("default: rowsPerRowGroup = %d, want %d", got, defaultRowGroupRowLimit)
}
}
+231
View File
@@ -0,0 +1,231 @@
package iceberg
import (
"bytes"
"math"
"sort"
"strings"
"time"
"unicode/utf8"
"github.com/apache/iceberg-go"
"github.com/parquet-go/parquet-go"
"github.com/parquet-go/parquet-go/format"
)
// boundTruncateLength is the default metrics mode of the Iceberg reference
// implementation, truncate(16), applied to string and binary bounds.
const boundTruncateLength = 16
// columnStats is what a data file's manifest entry records about its columns,
// keyed by Iceberg field id. Readers plan scans from it.
type columnStats struct {
sizes, values, nulls map[int]int64
lower, upper map[int][]byte
splitOffsets []int64
}
func (s *columnStats) applyTo(b *iceberg.DataFileBuilder) {
b.ColumnSizes(s.sizes).ValueCounts(s.values).NullValueCounts(s.nulls).
LowerBoundValues(s.lower).UpperBoundValues(s.upper).SplitOffsets(s.splitOffsets)
}
// collectColumnStats reads back the footer of a Parquet file the compactor
// wrote. Counts and sizes are summed over the row groups and the bounds are
// the minimum and maximum of the row group statistics, encoded with the
// spec's single-value serialization. A column whose bounds cannot be read
// for every row group, or converted exactly, gets no bounds.
func collectColumnStats(data []byte, schema *iceberg.Schema) (*columnStats, error) {
file, err := parquet.OpenFile(bytes.NewReader(data), int64(len(data)),
parquet.SkipPageIndex(true), parquet.SkipBloomFilters(true))
if err != nil {
return nil, err
}
var leaves []*parquet.Column
var walk func(*parquet.Column)
walk = func(c *parquet.Column) {
if c.Leaf() {
leaves = append(leaves, c)
return
}
for _, child := range c.Columns() {
walk(child)
}
}
walk(file.Root())
type bounds struct {
field iceberg.NestedField
typ parquet.Type
min, max parquet.Value
seen bool
invalid bool
}
stats := &columnStats{sizes: map[int]int64{}, values: map[int]int64{}, nulls: map[int]int64{},
lower: map[int][]byte{}, upper: map[int][]byte{}}
agg := map[int]*bounds{}
meta := file.Metadata()
for rgIdx, rg := range file.RowGroups() {
stats.splitOffsets = append(stats.splitOffsets, rowGroupOffset(&meta.RowGroups[rgIdx]))
for col, cc := range rg.ColumnChunks() {
field, ok := icebergFieldOf(leaves[col], schema)
if !ok {
continue
}
stats.sizes[field.ID] += meta.RowGroups[rgIdx].Columns[col].MetaData.TotalCompressedSize
stats.values[field.ID] += cc.NumValues()
b := agg[field.ID]
if b == nil {
b = &bounds{field: field, typ: leaves[col].Type()}
agg[field.ID] = b
}
chunk, ok := cc.(*parquet.FileColumnChunk)
if !ok {
b.invalid = true
continue
}
stats.nulls[field.ID] += chunk.NullCount()
if chunk.NullCount() == chunk.NumValues() {
continue // an all-null chunk constrains nothing
}
lo, hi, ok := chunk.Bounds()
switch {
case !ok:
b.invalid = true
case !b.seen:
b.min, b.max, b.seen = lo.Clone(), hi.Clone(), true
default:
if b.typ.Compare(lo, b.min) < 0 {
b.min = lo.Clone()
}
if b.typ.Compare(hi, b.max) > 0 {
b.max = hi.Clone()
}
}
}
}
for id, b := range agg {
if b.invalid || !b.seen {
continue
}
lo, okLo := boundLiteral(b.min, b.typ, b.field.Type)
hi, okHi := boundLiteral(b.max, b.typ, b.field.Type)
if !okLo || !okHi {
continue
}
lo, hi, okHi = truncateBounds(lo, hi)
if raw, err := lo.MarshalBinary(); err == nil {
stats.lower[id] = raw
}
if raw, err := hi.MarshalBinary(); err == nil && okHi {
stats.upper[id] = raw
}
}
sort.Slice(stats.splitOffsets, func(i, j int) bool { return stats.splitOffsets[i] < stats.splitOffsets[j] })
return stats, nil
}
// icebergFieldOf finds the Iceberg field of a Parquet leaf column by the field
// id the writer stored, or by its dotted path in a file written without ids.
// Only primitive fields outside lists and maps carry statistics.
func icebergFieldOf(leaf *parquet.Column, schema *iceberg.Schema) (iceberg.NestedField, bool) {
var field iceberg.NestedField
var ok bool
if leaf.MaxRepetitionLevel() > 0 {
return field, false
}
if id := leaf.ID(); id > 0 {
field, ok = schema.FindFieldByID(id)
} else {
field, ok = schema.FindFieldByName(strings.Join(leaf.Path(), "."))
}
if !ok {
return field, false
}
_, primitive := field.Type.(iceberg.PrimitiveType)
return field, primitive
}
// rowGroupOffset is where a row group's bytes start: the split offset readers
// plan one task per row group from.
func rowGroupOffset(rg *format.RowGroup) int64 {
if rg.FileOffset > 0 || len(rg.Columns) == 0 {
return rg.FileOffset
}
first := rg.Columns[0].MetaData
if first.DictionaryPageOffset > 0 && first.DictionaryPageOffset < first.DataPageOffset {
return first.DictionaryPageOffset
}
return first.DataPageOffset
}
// boundLiteral converts a Parquet statistics value into a literal of the
// column's Iceberg type, when the conversion is exact. Other types (decimal,
// time, uuid, fixed, timestamps not in the type's unit) get no bound.
func boundLiteral(v parquet.Value, pt parquet.Type, t iceberg.Type) (iceberg.Literal, bool) {
if v.IsNull() {
return nil, false
}
switch t.(type) {
case iceberg.BooleanType:
return iceberg.BoolLiteral(v.Boolean()), v.Kind() == parquet.Boolean
case iceberg.Int32Type:
return iceberg.Int32Literal(v.Int32()), v.Kind() == parquet.Int32
case iceberg.Int64Type:
switch v.Kind() {
case parquet.Int64:
return iceberg.Int64Literal(v.Int64()), true
case parquet.Int32:
return iceberg.Int64Literal(int64(v.Int32())), true
}
case iceberg.Float32Type:
f := v.Float()
return iceberg.Float32Literal(f), v.Kind() == parquet.Float && !math.IsNaN(float64(f))
case iceberg.Float64Type:
f := v.Double()
return iceberg.Float64Literal(f), v.Kind() == parquet.Double && !math.IsNaN(f)
case iceberg.DateType:
return iceberg.DateLiteral(iceberg.Date(v.Int32())), v.Kind() == parquet.Int32
case iceberg.TimestampType, iceberg.TimestampTzType:
return iceberg.TimestampLiteral(iceberg.Timestamp(v.Int64())), v.Kind() == parquet.Int64 && timeUnit(pt) == time.Microsecond
case iceberg.TimestampNsType:
return iceberg.TimestampNsLiteral(iceberg.TimestampNano(v.Int64())), v.Kind() == parquet.Int64 && timeUnit(pt) == time.Nanosecond
case iceberg.StringType:
b := v.ByteArray()
return iceberg.StringLiteral(string(b)), v.Kind() == parquet.ByteArray && utf8.Valid(b)
case iceberg.BinaryType:
return iceberg.BinaryLiteral(bytes.Clone(v.ByteArray())), v.Kind() == parquet.ByteArray
}
return nil, false
}
func timeUnit(pt parquet.Type) time.Duration {
if pt == nil || pt.LogicalType() == nil {
return 0
}
if ts, ok := pt.LogicalType().Value.(*format.TimestampType); ok && ts.Unit.Value != nil {
return ts.Unit.Value.Duration()
}
return 0
}
// truncateBounds applies truncate(16) to string and binary bounds. A prefix is
// still a lower bound; a shortened upper bound is not, so a value longer than
// the limit keeps no upper bound (the reference implementation increments the
// last character instead).
func truncateBounds(lo, hi iceberg.Literal) (iceberg.Literal, iceberg.Literal, bool) {
switch l := lo.(type) {
case iceberg.StringLiteral:
if r := []rune(string(l)); len(r) > boundTruncateLength {
lo = iceberg.StringLiteral(string(r[:boundTruncateLength]))
}
return lo, hi, utf8.RuneCountInString(string(hi.(iceberg.StringLiteral))) <= boundTruncateLength
case iceberg.BinaryLiteral:
if len(l) > boundTruncateLength {
lo = l[:boundTruncateLength]
}
return lo, hi, len(hi.(iceberg.BinaryLiteral)) <= boundTruncateLength
}
return lo, hi, true
}
@@ -0,0 +1,148 @@
package iceberg
import (
"bytes"
"context"
"testing"
"github.com/apache/iceberg-go"
"github.com/parquet-go/parquet-go"
filer_pb "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
)
// compactedDataFile returns the data file the current snapshot's data
// manifest references after a compaction run.
func compactedDataFile(t *testing.T, client filer_pb.SeaweedFilerClient, setup tableSetup) iceberg.DataFile {
t.Helper()
found := liveDataFiles(t, client, setup)
if len(found) != 1 {
t.Fatalf("expected 1 live data file after compaction, got %d", len(found))
}
return found[0]
}
func liveDataFiles(t *testing.T, client filer_pb.SeaweedFilerClient, setup tableSetup) []iceberg.DataFile {
t.Helper()
state, err := loadCurrentMetadata(context.Background(), client, setup.BucketName, setup.tablePath())
if err != nil {
t.Fatalf("loadCurrentMetadata: %v", err)
}
manifests, err := loadCurrentManifests(context.Background(), client, setup.BucketName, state.DataPath, state.Metadata)
if err != nil {
t.Fatalf("loadCurrentManifests: %v", err)
}
var found []iceberg.DataFile
for _, mf := range manifests {
if mf.ManifestContent() != iceberg.ManifestContentData {
continue
}
manifestData, err := loadFileByIcebergPath(context.Background(), client, setup.BucketName, state.DataPath, mf.FilePath())
if err != nil {
t.Fatalf("load manifest: %v", err)
}
entries, err := iceberg.ReadManifest(mf, bytes.NewReader(manifestData), true)
if err != nil {
t.Fatalf("read manifest: %v", err)
}
for _, entry := range entries {
if entry.Status() != iceberg.EntryStatusDELETED {
found = append(found, entry.DataFile())
}
}
}
return found
}
func compactedFileBytes(t *testing.T, client filer_pb.SeaweedFilerClient, setup tableSetup, filePath string) []byte {
t.Helper()
state, err := loadCurrentMetadata(context.Background(), client, setup.BucketName, setup.tablePath())
if err != nil {
t.Fatalf("loadCurrentMetadata: %v", err)
}
data, err := loadFileByIcebergPath(context.Background(), client, setup.BucketName, state.DataPath, filePath)
if err != nil {
t.Fatalf("load compacted file: %v", err)
}
return data
}
func compactedIDs(t *testing.T, client filer_pb.SeaweedFilerClient, setup tableSetup, filePath string) []int64 {
t.Helper()
data := compactedFileBytes(t, client, setup, filePath)
f, err := parquet.OpenFile(bytes.NewReader(data), int64(len(data)))
if err != nil {
t.Fatalf("open compacted file: %v", err)
}
var ids []int64
for _, rg := range f.RowGroups() {
rows := rg.Rows()
buf := make([]parquet.Row, 8)
for {
n, err := rows.ReadRows(buf)
for _, row := range buf[:n] {
ids = append(ids, row[0].Int64())
}
if err != nil {
break
}
}
rows.Close()
}
return ids
}
// A compacted file's manifest entry carries the column metrics readers
// prune files by: counts, sizes, bounds and split offsets read back
// from the file's own footer.
func TestCompactDataFilesRecordsColumnStatistics(t *testing.T) {
fs, client := startFakeFiler(t)
setup := tableSetup{BucketName: "tb", Namespace: "ns", TableName: "tbl"}
populateTableWithDeleteFiles(t, fs, setup,
[]struct {
Name string
Rows []struct {
ID int64
Name string
}
}{
{"d1.parquet", []struct {
ID int64
Name string
}{{1, "a"}, {2, "b"}}},
{"d2.parquet", []struct {
ID int64
Name string
}{{3, "c"}}},
},
nil, nil,
)
handler := NewHandler(nil)
config := Config{TargetFileSizeBytes: 256 * 1024 * 1024, MinInputFiles: 2, MaxCommitRetries: 3}
if _, _, err := handler.compactDataFiles(context.Background(), client, setup.BucketName, setup.tablePath(), config, nil); err != nil {
t.Fatalf("compactDataFiles: %v", err)
}
df := compactedDataFile(t, client, setup)
if df.ValueCounts()[1] != 3 {
t.Errorf("value_counts[1] = %d, want 3", df.ValueCounts()[1])
}
if len(df.LowerBoundValues()) == 0 || len(df.UpperBoundValues()) == 0 {
t.Fatal("compacted file has no bounds")
}
lo, err := iceberg.LiteralFromBytes(iceberg.PrimitiveTypes.Int64, df.LowerBoundValues()[1])
if err != nil {
t.Fatalf("decode lower bound: %v", err)
}
hi, err := iceberg.LiteralFromBytes(iceberg.PrimitiveTypes.Int64, df.UpperBoundValues()[1])
if err != nil {
t.Fatalf("decode upper bound: %v", err)
}
if lo.(iceberg.Int64Literal) != 1 || hi.(iceberg.Int64Literal) != 3 {
t.Errorf("id bounds = %v..%v, want 1..3", lo, hi)
}
if len(df.SplitOffsets()) == 0 {
t.Error("compacted file has no split_offsets")
}
}
+4
View File
@@ -130,6 +130,10 @@ type Config struct {
SortMaxInputBytes int64
SortBufferRows int64
SortSpillDir string
// Row groups of a compacted file, from the table's write.parquet.*
// properties; zero means the default.
RowGroupSizeBytes int64
RowGroupRowLimit int64
}
// ParseConfig extracts an iceberg maintenance Config from plugin config values.
+1 -1
View File
@@ -477,7 +477,7 @@ func hasEligibleCompaction(
}
targetSize := compactionTargetSizeForPlan(config, rewritePlan)
bins := buildCompactionBins(candidateEntries, targetSize, minInputFiles)
bins := buildCompactionBins(candidateEntries, targetSize, minInputFiles, meta)
bins = filterCompactionBinsByPlan(bins, config, rewritePlan)
return len(bins) > 0, nil
}
+10 -10
View File
@@ -701,7 +701,7 @@ func TestBuildCompactionBins(t *testing.T) {
{path: "data/f3.parquet", size: 4096, partition: map[int]any{}},
})
bins := buildCompactionBins(entries, targetSize, minFiles)
bins := buildCompactionBins(entries, targetSize, minFiles, nil)
if len(bins) != 1 {
t.Fatalf("expected 1 bin, got %d", len(bins))
}
@@ -720,7 +720,7 @@ func TestBuildCompactionBinsFiltersLargeFiles(t *testing.T) {
{path: "data/large.parquet", size: 5000, partition: map[int]any{}},
})
bins := buildCompactionBins(entries, targetSize, minFiles)
bins := buildCompactionBins(entries, targetSize, minFiles, nil)
if len(bins) != 1 {
t.Fatalf("expected 1 bin, got %d", len(bins))
}
@@ -756,7 +756,7 @@ func TestBuildCompactionBinsLowercaseParquetFormat(t *testing.T) {
entries[0] = lowercaseFormatEntry{entries[0]}
entries[1] = lowercaseFormatEntry{entries[1]}
bins := buildCompactionBins(entries, targetSize, minFiles)
bins := buildCompactionBins(entries, targetSize, minFiles, nil)
if len(bins) != 1 {
t.Fatalf("expected 1 bin, got %d", len(bins))
}
@@ -774,7 +774,7 @@ func TestBuildCompactionBinsMinFilesThreshold(t *testing.T) {
{path: "data/f2.parquet", size: 2048, partition: map[int]any{}},
})
bins := buildCompactionBins(entries, targetSize, minFiles)
bins := buildCompactionBins(entries, targetSize, minFiles, nil)
if len(bins) != 0 {
t.Errorf("expected 0 bins (below min threshold), got %d", len(bins))
}
@@ -801,7 +801,7 @@ func TestBuildCompactionBinsMultiplePartitions(t *testing.T) {
{path: "data/b3.parquet", size: 4096, partition: partB, partitionSpec: &partitionSpec},
})
bins := buildCompactionBins(entries, targetSize, minFiles)
bins := buildCompactionBins(entries, targetSize, minFiles, nil)
if len(bins) != 2 {
t.Fatalf("expected 2 bins (one per partition), got %d", len(bins))
}
@@ -977,7 +977,7 @@ func TestBuildCompactionBinsMultipleSpecs(t *testing.T) {
{path: "data/s1-f2.parquet", size: 2048, partition: map[int]any{}, specID: 1},
}, partSpecs)
bins := buildCompactionBins(entries, targetSize, minFiles)
bins := buildCompactionBins(entries, targetSize, minFiles, nil)
if len(bins) != 2 {
t.Fatalf("expected 2 bins (one per spec), got %d", len(bins))
}
@@ -1011,7 +1011,7 @@ func TestBuildCompactionBinsSingleSpec(t *testing.T) {
{path: "data/f3.parquet", size: 4096, partition: map[int]any{}, specID: 0},
}, partSpecs)
bins := buildCompactionBins(entries, targetSize, minFiles)
bins := buildCompactionBins(entries, targetSize, minFiles, nil)
if len(bins) != 1 {
t.Fatalf("expected 1 bin, got %d", len(bins))
}
@@ -1303,7 +1303,7 @@ func TestMergeParquetFilesWithPositionDeletes(t *testing.T) {
merged, count, err := mergeParquetFiles(
context.Background(), client, "test-bucket", "ns/tbl",
entries, posDeletes, nil, nil,
entries, posDeletes, nil, nil, 0,
)
if err != nil {
t.Fatalf("mergeParquetFiles: %v", err)
@@ -1419,7 +1419,7 @@ func TestMergeParquetFilesWithEqualityDeletes(t *testing.T) {
merged, count, err := mergeParquetFiles(
context.Background(), client, "test-bucket", "ns/tbl",
entries, nil, eqGroups, schema,
entries, nil, eqGroups, schema, 0,
)
if err != nil {
t.Fatalf("mergeParquetFiles: %v", err)
@@ -1479,7 +1479,7 @@ func TestMergeParquetFilesDictionaryEncodedInput(t *testing.T) {
merged, count, err := mergeParquetFiles(
context.Background(), client, "test-bucket", "ns/tbl",
entries, nil, nil, nil,
entries, nil, nil, nil, 0,
)
if err != nil {
t.Fatalf("mergeParquetFiles: %v", err)
@@ -19,6 +19,8 @@ import (
const (
propTargetFileSize = "write.target-file-size-bytes"
propDeleteTargetFileSize = "write.delete.target-file-size-bytes"
propRowGroupSize = "write.parquet.row-group-size-bytes"
propRowGroupRowLimit = "write.parquet.row-group-limit"
propMaxSnapshotAgeMs = "history.expire.max-snapshot-age-ms"
propMinSnapshotsToKeep = "history.expire.min-snapshots-to-keep"
)
@@ -50,6 +52,12 @@ func applyTableProperties(cfg Config, props iceberg.Properties) Config {
if v, ok := propInt64(props, propDeleteTargetFileSize); ok {
cfg.DeleteTargetFileSizeBytes = v
}
if v, ok := propInt64(props, propRowGroupSize); ok {
cfg.RowGroupSizeBytes = v
}
if v, ok := propInt64(props, propRowGroupRowLimit); ok {
cfg.RowGroupRowLimit = v
}
if v, ok := propInt64(props, propMaxSnapshotAgeMs); ok {
cfg.SnapshotRetentionMs = v
}