mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-09-30 11:45:42 +00:00
vacuum: keep disk-full read-only volumes reclaimable (#11519)
* storage/topology: keep disk-full read-only volumes vacuumable The vacuum sweep skipped every read-only replica, so a volume that went read-only because its disk filled could never reclaim its garbage — the exact situation compaction exists for. The volume server now reports disk_space_low in VacuumVolumeCheckResponse, and the sweep skips a read-only replica only when the flag is clear. An explicit volumeId vacuum is unaffected: it already bypassed the read-only rule. The field takes number 4: 2 and 3 are downstream-allocated for tombstone retention, keeping the wire merge clean. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * storage: measure vacuum free space against live bytes The pre-compaction space check required the current .dat + .idx size free, which includes the garbage being reclaimed — on a nearly full disk that estimate can never fit, so the volume stayed garbage-bound forever. Measure against the estimated compacted output instead: superblock plus live index entries plus live content bytes, with the existing ten percent buffer unchanged. Mirrors the same check in the Rust volume server. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * vacuum: count per-needle framing in the compacted-size estimate The live-bytes estimate covered each live needle's content and index entry but not its .dat framing (header, checksum, timestamp, padding — ~32 bytes on version 3). For small-needle volumes that is more than the 10% headroom, so a disk with space between the estimate and the real output still ran out mid-compaction. Rust side mirrors the same formula. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * storage: report disk_space_low only when it is the sole read-only cause Review feedback (ihnokim, greptile, devin): a volume read-only for low disk space AND an operator mark or I/O quarantine was still eligible for the automatic sweep, rewriting a copy meant to stay protected. The flag now reports only the benign sole-cause case in both servers. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * topology: fail closed when the read-only lookup misses in the sweep A heartbeat can drop the volume from the DataNode cache between the location-list copy and VacuumVolumeCheck; a lookup error previously skipped the read-only check entirely. Review feedback (coderabbit). 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>
This commit is contained in:
co-authored by
Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
parent
38ce95d960
commit
0978e7f833
@@ -168,6 +168,7 @@ message VacuumVolumeCheckRequest {
|
||||
}
|
||||
message VacuumVolumeCheckResponse {
|
||||
double garbage_ratio = 1;
|
||||
bool disk_space_low = 4; // the volume is read-only solely because its disk is low on space — a cause compaction itself reclaims
|
||||
}
|
||||
|
||||
message VacuumVolumeCompactRequest {
|
||||
|
||||
@@ -1127,14 +1127,24 @@ impl VolumeServer for VolumeGrpcService {
|
||||
) -> Result<Response<volume_server_pb::VacuumVolumeCheckResponse>, Status> {
|
||||
let vid = VolumeId(request.into_inner().volume_id);
|
||||
let store = self.state.store.read().unwrap();
|
||||
let garbage_ratio = match store.find_volume(vid) {
|
||||
Some((_, vol)) => vol.garbage_level(),
|
||||
let (garbage_ratio, disk_space_low) = match store.find_volume(vid) {
|
||||
Some((_, vol)) => {
|
||||
// disk_space_low only counts when it is the sole read-only
|
||||
// cause — an operator mark or I/O quarantine still shields
|
||||
// the volume.
|
||||
let (_, no_write_or_delete, no_write_can_delete, is_low) = vol.read_only_reasons();
|
||||
(
|
||||
vol.garbage_level(),
|
||||
is_low && !no_write_or_delete && !no_write_can_delete,
|
||||
)
|
||||
}
|
||||
None => {
|
||||
return Err(crate::storage::volume::VolumeError::VolumeNotFound(vid).into());
|
||||
}
|
||||
};
|
||||
Ok(Response::new(volume_server_pb::VacuumVolumeCheckResponse {
|
||||
garbage_ratio,
|
||||
disk_space_low,
|
||||
}))
|
||||
}
|
||||
|
||||
|
||||
@@ -14,9 +14,9 @@ use crate::pb::master_pb;
|
||||
use crate::storage::disk_location::DiskLocation;
|
||||
use crate::storage::erasure_coding::ec_shard::{EcVolumeShard, MAX_SHARD_COUNT, ShardId};
|
||||
use crate::storage::erasure_coding::ec_volume::{EcVolume, is_usable_ecx_file};
|
||||
use crate::storage::needle::needle::Needle;
|
||||
use crate::storage::needle::needle::{Needle, get_actual_size};
|
||||
use crate::storage::needle_map::NeedleMapKind;
|
||||
use crate::storage::super_block::ReplicaPlacement;
|
||||
use crate::storage::super_block::{ReplicaPlacement, SUPER_BLOCK_SIZE};
|
||||
use crate::storage::types::*;
|
||||
use crate::storage::volume::{CompactionJob, VifVolumeInfo, VolumeError, VolumeSpec};
|
||||
|
||||
@@ -1539,13 +1539,21 @@ impl Store {
|
||||
preallocate: u64,
|
||||
) -> Result<Option<CompactionJob>, VolumeError> {
|
||||
// Required space matches Go's CompactVolume check: the larger of the
|
||||
// requested preallocation and the estimated volume size.
|
||||
// requested preallocation and the estimated compacted size — the live
|
||||
// needles, not the .dat the garbage already occupies, so a full disk
|
||||
// can still be reclaimed.
|
||||
let (loc_idx, space_needed) = {
|
||||
let (loc_idx, v) = self
|
||||
.find_volume(vid)
|
||||
.ok_or(VolumeError::VolumeNotFound(vid))?;
|
||||
let estimated = v.dat_file_size().unwrap_or(0) + v.idx_file_size();
|
||||
(loc_idx, std::cmp::max(preallocate, estimated))
|
||||
let live_count = (v.file_count() - v.deleted_count()).max(0) as u64;
|
||||
let live_bytes = v.content_size().saturating_sub(v.deleted_size());
|
||||
let per_needle = (get_actual_size(Size(0), v.version())
|
||||
+ NEEDLE_PADDING_SIZE as i64
|
||||
+ NEEDLE_MAP_ENTRY_SIZE as i64) as u64;
|
||||
let estimated = SUPER_BLOCK_SIZE as u64 + live_count * per_needle + live_bytes;
|
||||
let space_needed = std::cmp::max(preallocate, estimated);
|
||||
(loc_idx, space_needed + space_needed / 10)
|
||||
};
|
||||
|
||||
let dir = self.locations[loc_idx].directory.clone();
|
||||
|
||||
@@ -2690,6 +2690,19 @@ impl Volume {
|
||||
|| self.location_disk_space_low.load(Ordering::Relaxed)
|
||||
}
|
||||
|
||||
/// Mirrors Go's ReadOnlyReasons: `no_write_or_delete` already covers the
|
||||
/// io_unavailable quarantine.
|
||||
pub fn read_only_reasons(&self) -> (bool, bool, bool, bool) {
|
||||
let no_write_or_delete = self.no_write_or_delete || self.io_unavailable.is_some();
|
||||
let disk_space_low = self.location_disk_space_low.load(Ordering::Relaxed);
|
||||
(
|
||||
no_write_or_delete || self.no_write_can_delete || disk_space_low,
|
||||
no_write_or_delete,
|
||||
self.no_write_can_delete,
|
||||
disk_space_low,
|
||||
)
|
||||
}
|
||||
|
||||
/// The reason the volume refuses all I/O, when a failed recovery left the
|
||||
/// .dat/index pair unverified. Mirrors Go's unavailableError.
|
||||
pub fn unavailable_error(&self) -> Option<VolumeError> {
|
||||
|
||||
@@ -168,6 +168,7 @@ message VacuumVolumeCheckRequest {
|
||||
}
|
||||
message VacuumVolumeCheckResponse {
|
||||
double garbage_ratio = 1;
|
||||
bool disk_space_low = 4; // the volume is read-only solely because its disk is low on space — a cause compaction itself reclaims
|
||||
}
|
||||
|
||||
message VacuumVolumeCompactRequest {
|
||||
|
||||
@@ -436,6 +436,7 @@ func (x *VacuumVolumeCheckRequest) GetVolumeId() uint32 {
|
||||
type VacuumVolumeCheckResponse struct {
|
||||
state protoimpl.MessageState `protogen:"open.v1"`
|
||||
GarbageRatio float64 `protobuf:"fixed64,1,opt,name=garbage_ratio,json=garbageRatio,proto3" json:"garbage_ratio,omitempty"`
|
||||
DiskSpaceLow bool `protobuf:"varint,4,opt,name=disk_space_low,json=diskSpaceLow,proto3" json:"disk_space_low,omitempty"` // the volume is read-only solely because its disk is low on space — a cause compaction itself reclaims
|
||||
unknownFields protoimpl.UnknownFields
|
||||
sizeCache protoimpl.SizeCache
|
||||
}
|
||||
@@ -477,6 +478,13 @@ func (x *VacuumVolumeCheckResponse) GetGarbageRatio() float64 {
|
||||
return 0
|
||||
}
|
||||
|
||||
func (x *VacuumVolumeCheckResponse) GetDiskSpaceLow() bool {
|
||||
if x != nil {
|
||||
return x.DiskSpaceLow
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
type VacuumVolumeCompactRequest struct {
|
||||
state protoimpl.MessageState `protogen:"open.v1"`
|
||||
VolumeId uint32 `protobuf:"varint,1,opt,name=volume_id,json=volumeId,proto3" json:"volume_id,omitempty"`
|
||||
@@ -7276,9 +7284,10 @@ const file_volume_server_proto_rawDesc = "" +
|
||||
"\aversion\x18\x05 \x01(\rR\aversion\"\a\n" +
|
||||
"\x05Empty\"7\n" +
|
||||
"\x18VacuumVolumeCheckRequest\x12\x1b\n" +
|
||||
"\tvolume_id\x18\x01 \x01(\rR\bvolumeId\"@\n" +
|
||||
"\tvolume_id\x18\x01 \x01(\rR\bvolumeId\"f\n" +
|
||||
"\x19VacuumVolumeCheckResponse\x12#\n" +
|
||||
"\rgarbage_ratio\x18\x01 \x01(\x01R\fgarbageRatio\"[\n" +
|
||||
"\rgarbage_ratio\x18\x01 \x01(\x01R\fgarbageRatio\x12$\n" +
|
||||
"\x0edisk_space_low\x18\x04 \x01(\bR\fdiskSpaceLow\"[\n" +
|
||||
"\x1aVacuumVolumeCompactRequest\x12\x1b\n" +
|
||||
"\tvolume_id\x18\x01 \x01(\rR\bvolumeId\x12 \n" +
|
||||
"\vpreallocate\x18\x02 \x01(\x03R\vpreallocate\"f\n" +
|
||||
|
||||
@@ -22,9 +22,10 @@ func (vs *VolumeServer) VacuumVolumeCheck(ctx context.Context, req *volume_serve
|
||||
|
||||
resp := &volume_server_pb.VacuumVolumeCheckResponse{}
|
||||
|
||||
garbageRatio, err := vs.store.CheckCompactVolume(needle.VolumeId(req.VolumeId))
|
||||
garbageRatio, diskSpaceLow, err := vs.store.CheckCompactVolume(needle.VolumeId(req.VolumeId))
|
||||
|
||||
resp.GarbageRatio = garbageRatio
|
||||
resp.DiskSpaceLow = diskSpaceLow
|
||||
|
||||
if err != nil {
|
||||
glog.V(3).Infof("check volume %d: %v", req.VolumeId, err)
|
||||
|
||||
@@ -7,16 +7,21 @@ import (
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/glog"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/needle"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/super_block"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/types"
|
||||
)
|
||||
|
||||
var ErrInsufficientSpace = fmt.Errorf("insufficient free space")
|
||||
|
||||
func (s *Store) CheckCompactVolume(volumeId needle.VolumeId) (float64, error) {
|
||||
func (s *Store) CheckCompactVolume(volumeId needle.VolumeId) (garbageRatio float64, diskSpaceLow bool, err error) {
|
||||
if v := s.findVolume(volumeId); v != nil {
|
||||
glog.V(3).Infof("volume %d garbage level: %f", volumeId, v.garbageLevel())
|
||||
return v.garbageLevel(), nil
|
||||
// diskSpaceLow only counts when it is the sole read-only cause — an
|
||||
// operator mark or I/O quarantine still shields the volume.
|
||||
_, noWriteOrDelete, noWriteCanDelete, isLow := v.ReadOnlyReasons()
|
||||
return v.garbageLevel(), isLow && !noWriteOrDelete && !noWriteCanDelete, nil
|
||||
}
|
||||
return 0, fmt.Errorf("volume id %d is not found during check compact: %w", volumeId, ErrVolumeNotFound)
|
||||
return 0, false, fmt.Errorf("volume id %d is not found during check compact: %w", volumeId, ErrVolumeNotFound)
|
||||
}
|
||||
|
||||
func (s *Store) CompactVolume(vid needle.VolumeId, preallocate int64, compactionBytePerSecond int64, progressFn ProgressFunc) error {
|
||||
@@ -56,18 +61,39 @@ func (s *Store) CommitCleanupVolume(vid needle.VolumeId) error {
|
||||
return fmt.Errorf("volume id %d is not found during cleaning up: %w", vid, ErrVolumeNotFound)
|
||||
}
|
||||
|
||||
// estimatedCompactedSize is what compaction writes: a superblock, the live
|
||||
// needles with their on-disk framing, and an index with live entries only.
|
||||
// Deleted bytes do not carry over, so a mostly-garbage volume needs far less
|
||||
// space than it occupies.
|
||||
func estimatedCompactedSize(v *Volume) int64 {
|
||||
liveCount := v.FileCount()
|
||||
if deleted := v.DeletedCount(); deleted < liveCount {
|
||||
liveCount -= deleted
|
||||
} else {
|
||||
liveCount = 0
|
||||
}
|
||||
liveBytes := v.ContentSize()
|
||||
if deleted := v.DeletedSize(); deleted < liveBytes {
|
||||
liveBytes -= deleted
|
||||
} else {
|
||||
liveBytes = 0
|
||||
}
|
||||
perNeedle := needle.GetActualSize(0, v.Version()) + types.NeedlePaddingSize + types.NeedleMapEntrySize
|
||||
return super_block.SuperBlockSize + int64(liveCount)*perNeedle + int64(liveBytes)
|
||||
}
|
||||
|
||||
func ensureCompactVolumeSpace(v *Volume, preallocate int64) error {
|
||||
// Get current volume size for space calculation
|
||||
volumeSize, indexSize, _ := v.FileStat()
|
||||
|
||||
// Calculate space needed for compaction:
|
||||
// 1. Space for the new compacted volume (approximately same as current volume size)
|
||||
// 2. Use the larger of preallocate or estimated volume size
|
||||
estimatedCompactSize := int64(volumeSize + indexSize)
|
||||
// The compacted output holds live needles only, so measure against the
|
||||
// estimated compacted size — otherwise a disk full of garbage can never
|
||||
// reclaim itself.
|
||||
estimatedCompactSize := estimatedCompactedSize(v)
|
||||
spaceNeeded := preallocate
|
||||
if estimatedCompactSize > preallocate {
|
||||
spaceNeeded = estimatedCompactSize
|
||||
}
|
||||
spaceNeeded += spaceNeeded / 10
|
||||
|
||||
diskStatus := stats.NewDiskStatus(v.dir)
|
||||
if int64(diskStatus.Free) < spaceNeeded {
|
||||
|
||||
@@ -2,6 +2,12 @@ package storage
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"github.com/stretchr/testify/require"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/needle"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/super_block"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/types"
|
||||
)
|
||||
|
||||
func TestSpaceCalculation(t *testing.T) {
|
||||
@@ -49,3 +55,104 @@ func TestSpaceCalculation(t *testing.T) {
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// Compaction writes live needles only, so the space check must be measured
|
||||
// against the live size, not the .dat the garbage occupies — a full disk
|
||||
// needs the estimate to shrink or it can never reclaim.
|
||||
func TestEstimatedCompactedSizeCountsLiveNeedles(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
|
||||
v, err := NewVolume(dir, dir, "", 1, NeedleMapInMemory, &super_block.ReplicaPlacement{}, &needle.TTL{}, 0, needle.GetCurrentVersion(), 0, 0)
|
||||
if err != nil {
|
||||
t.Fatalf("volume creation: %v", err)
|
||||
}
|
||||
defer v.Close()
|
||||
|
||||
const count = 20
|
||||
for i := 1; i <= count; i++ {
|
||||
if _, _, _, err := v.writeNeedle2(newRandomNeedle(uint64(i)), true, false, false); err != nil {
|
||||
t.Fatalf("write needle %d: %v", i, err)
|
||||
}
|
||||
}
|
||||
datSize, _, _ := v.FileStat()
|
||||
|
||||
fullEstimate := estimatedCompactedSize(v)
|
||||
if fullEstimate <= super_block.SuperBlockSize {
|
||||
t.Fatalf("estimate for all-live volume = %d, want > superblock", fullEstimate)
|
||||
}
|
||||
|
||||
for i := 1; i < count; i++ {
|
||||
if _, err := v.doDeleteRequest(newEmptyNeedle(uint64(i))); err != nil {
|
||||
t.Fatalf("delete needle %d: %v", i, err)
|
||||
}
|
||||
}
|
||||
|
||||
estimate := estimatedCompactedSize(v)
|
||||
if estimate >= int64(datSize) {
|
||||
t.Fatalf("estimate %d not below .dat size %d with 19/20 needles deleted", estimate, datSize)
|
||||
}
|
||||
if estimate <= super_block.SuperBlockSize {
|
||||
t.Fatalf("estimate %d lost the one live needle", estimate)
|
||||
}
|
||||
live := int64(v.FileCount()-v.DeletedCount())*types.NeedleMapEntrySize + super_block.SuperBlockSize
|
||||
if estimate < live {
|
||||
t.Fatalf("estimate %d below superblock + live index entries %d", estimate, live)
|
||||
}
|
||||
|
||||
if _, err := v.doDeleteRequest(newEmptyNeedle(uint64(count))); err != nil {
|
||||
t.Fatalf("delete last needle: %v", err)
|
||||
}
|
||||
if estimate := estimatedCompactedSize(v); estimate != super_block.SuperBlockSize {
|
||||
t.Fatalf("all-deleted estimate = %d, want superblock only (%d)", estimate, super_block.SuperBlockSize)
|
||||
}
|
||||
}
|
||||
|
||||
// The estimate must cover what compaction writes on disk: each live needle's
|
||||
// content plus its header, checksum, timestamp and padding. An all-live
|
||||
// volume's compacted .dat is byte-for-byte its current one, so the estimate
|
||||
// may not fall below the current file.
|
||||
func TestEstimatedCompactedSizeCoversNeedleFraming(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
|
||||
v, err := NewVolume(dir, dir, "", 1, NeedleMapInMemory, &super_block.ReplicaPlacement{}, &needle.TTL{}, 0, needle.GetCurrentVersion(), 0, 0)
|
||||
if err != nil {
|
||||
t.Fatalf("volume creation: %v", err)
|
||||
}
|
||||
defer v.Close()
|
||||
|
||||
for i := 1; i <= 100; i++ {
|
||||
if _, _, _, err := v.writeNeedle2(newRandomNeedle(uint64(i)), true, false, false); err != nil {
|
||||
t.Fatalf("write needle %d: %v", i, err)
|
||||
}
|
||||
}
|
||||
datSize, _, _ := v.FileStat()
|
||||
|
||||
if estimate := estimatedCompactedSize(v); estimate < int64(datSize) {
|
||||
t.Fatalf("estimate %d below .dat size %d for an all-live volume: missing per-needle framing", estimate, datSize)
|
||||
}
|
||||
}
|
||||
|
||||
// disk_space_low is only reported when low space is the sole read-only cause,
|
||||
// so a volume also marked read-only by an operator or quarantined by failed
|
||||
// I/O stays out of the sweep.
|
||||
func TestCheckCompactVolumeDiskLowSoleCauseOnly(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
store := newSingleDirStore(t, dir)
|
||||
defer store.Close()
|
||||
const vid = needle.VolumeId(7)
|
||||
require.NoError(t, store.AddVolume(vid, "", NeedleMapInMemory, "000", "", 0, needle.GetCurrentVersion(), 0, types.HardDriveType, 0))
|
||||
|
||||
_, low, err := store.CheckCompactVolume(vid)
|
||||
require.NoError(t, err)
|
||||
require.False(t, low)
|
||||
|
||||
store.Locations[0].isDiskSpaceLow.Store(true)
|
||||
_, low, err = store.CheckCompactVolume(vid)
|
||||
require.NoError(t, err)
|
||||
require.True(t, low)
|
||||
|
||||
require.NoError(t, store.MarkVolumeReadonly(vid, false, false))
|
||||
_, low, err = store.CheckCompactVolume(vid)
|
||||
require.NoError(t, err)
|
||||
require.False(t, low)
|
||||
}
|
||||
|
||||
@@ -23,11 +23,11 @@ import (
|
||||
)
|
||||
|
||||
func (t *Topology) batchVacuumVolumeCheck(grpcDialOption grpc.DialOption, vid needle.VolumeId,
|
||||
locationlist *VolumeLocationList, garbageThreshold float64) (*VolumeLocationList, bool) {
|
||||
locationlist *VolumeLocationList, garbageThreshold float64, skipReadOnly bool) (*VolumeLocationList, bool) {
|
||||
ch := make(chan int, locationlist.Length())
|
||||
errCount := int32(0)
|
||||
for index, dn := range locationlist.list {
|
||||
go func(index int, url pb.ServerAddress, vid needle.VolumeId) {
|
||||
go func(index int, dn *DataNode, url pb.ServerAddress, vid needle.VolumeId) {
|
||||
err := operation.WithVolumeServerClient(false, url, grpcDialOption, func(volumeServerClient volume_server_pb.VolumeServerClient) error {
|
||||
resp, err := volumeServerClient.VacuumVolumeCheck(context.Background(), &volume_server_pb.VacuumVolumeCheckRequest{
|
||||
VolumeId: uint32(vid),
|
||||
@@ -37,6 +37,21 @@ func (t *Topology) batchVacuumVolumeCheck(grpcDialOption grpc.DialOption, vid ne
|
||||
ch <- -1
|
||||
return err
|
||||
}
|
||||
// A sweep skips a read-only copy unless low disk space is its
|
||||
// only read-only cause — that is the copy compaction exists for.
|
||||
if skipReadOnly {
|
||||
v, lookErr := dn.GetVolumesById(vid)
|
||||
if lookErr != nil {
|
||||
atomic.AddInt32(&errCount, 1)
|
||||
ch <- -1
|
||||
return lookErr
|
||||
}
|
||||
if v.ReadOnly && !resp.DiskSpaceLow {
|
||||
glog.V(0).Infof("skip vacuuming read-only volume %d on %s", vid, url)
|
||||
ch <- -1
|
||||
return nil
|
||||
}
|
||||
}
|
||||
if resp.GarbageRatio >= garbageThreshold {
|
||||
ch <- index
|
||||
} else {
|
||||
@@ -47,7 +62,7 @@ func (t *Topology) batchVacuumVolumeCheck(grpcDialOption grpc.DialOption, vid ne
|
||||
if err != nil {
|
||||
glog.V(0).Infof("Checking vacuuming %d on %s: %v", vid, url, err)
|
||||
}
|
||||
}(index, dn.ServerAddress(), vid)
|
||||
}(index, dn, dn.ServerAddress(), vid)
|
||||
}
|
||||
vacuumLocationList := NewVolumeLocationList()
|
||||
|
||||
@@ -356,18 +371,16 @@ func (t *Topology) vacuumOneVolumeLayout(grpcDialOption grpc.DialOption, volumeL
|
||||
}
|
||||
|
||||
// skipReadOnly is set by the background scan and all-volumes sweep, where a
|
||||
// read-only flag usually means an unhealthy disk. An explicit volumeId clears
|
||||
// it so a benignly read-only (full/oversized) volume can be reclaimed.
|
||||
// read-only flag usually means an unhealthy disk. Even then a copy that is
|
||||
// read-only because its disk is low on space stays eligible — compaction is
|
||||
// how the space comes back. An explicit volumeId clears the rule entirely.
|
||||
func (t *Topology) vacuumOneVolumeId(grpcDialOption grpc.DialOption, volumeLayout *VolumeLayout, c *Collection, garbageThreshold float64, locationList *VolumeLocationList, vid needle.VolumeId, preallocate int64, skipReadOnly bool) {
|
||||
volumeLayout.accessLock.RLock()
|
||||
isReadOnly := volumeLayout.vid2location[vid].AnyReadOnly()
|
||||
isEnoughCopies := volumeLayout.enoughCopies(vid)
|
||||
volumeLayout.accessLock.RUnlock()
|
||||
|
||||
if isReadOnly {
|
||||
if skipReadOnly {
|
||||
return
|
||||
}
|
||||
if isReadOnly && !skipReadOnly {
|
||||
glog.V(0).Infof("vacuuming read-only volume %d on explicit request", vid)
|
||||
}
|
||||
if !isEnoughCopies {
|
||||
@@ -377,7 +390,7 @@ func (t *Topology) vacuumOneVolumeId(grpcDialOption grpc.DialOption, volumeLayou
|
||||
|
||||
glog.V(1).Infof("check vacuum on collection:%s volume:%d", c.Name, vid)
|
||||
if vacuumLocationList, needVacuum := t.batchVacuumVolumeCheck(
|
||||
grpcDialOption, vid, locationList, garbageThreshold); needVacuum {
|
||||
grpcDialOption, vid, locationList, garbageThreshold, skipReadOnly); needVacuum {
|
||||
if t.batchVacuumVolumeCompact(grpcDialOption, volumeLayout, vid, vacuumLocationList, preallocate) {
|
||||
t.batchVacuumVolumeCommit(grpcDialOption, volumeLayout, vid, vacuumLocationList, locationList)
|
||||
} else {
|
||||
|
||||
@@ -64,7 +64,7 @@ func (f *fakeVolumeDeleteServer) VolumeDelete(ctx context.Context, req *volume_s
|
||||
return &volume_server_pb.VolumeDeleteResponse{}, nil
|
||||
}
|
||||
|
||||
func startFakeVolumeServer(t *testing.T, vs *fakeVolumeDeleteServer) (grpcPort int, dialOption grpc.DialOption) {
|
||||
func startFakeVolumeServer(t *testing.T, vs volume_server_pb.VolumeServerServer) (grpcPort int, dialOption grpc.DialOption) {
|
||||
t.Helper()
|
||||
lis, err := net.Listen("tcp", "127.0.0.1:0")
|
||||
if err != nil {
|
||||
@@ -253,3 +253,110 @@ func TestDeleteEmptyVolumesKeepsVidWhenCopyDeleteFails(t *testing.T) {
|
||||
t.Fatalf("VolumeDelete calls = %d, want 1 (only the reachable copy)", len(fake.deletes))
|
||||
}
|
||||
}
|
||||
|
||||
type fakeVacuumServer struct {
|
||||
volume_server_pb.UnimplementedVolumeServerServer
|
||||
mu sync.Mutex
|
||||
checks map[uint32]*volume_server_pb.VacuumVolumeCheckResponse
|
||||
committed []uint32
|
||||
}
|
||||
|
||||
func (f *fakeVacuumServer) VacuumVolumeCheck(_ context.Context, req *volume_server_pb.VacuumVolumeCheckRequest) (*volume_server_pb.VacuumVolumeCheckResponse, error) {
|
||||
resp, ok := f.checks[req.VolumeId]
|
||||
if !ok {
|
||||
return nil, fmt.Errorf("volume %d not found", req.VolumeId)
|
||||
}
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
func (f *fakeVacuumServer) VacuumVolumeCompact(_ *volume_server_pb.VacuumVolumeCompactRequest, _ volume_server_pb.VolumeServer_VacuumVolumeCompactServer) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (f *fakeVacuumServer) VacuumVolumeCommit(_ context.Context, req *volume_server_pb.VacuumVolumeCommitRequest) (*volume_server_pb.VacuumVolumeCommitResponse, error) {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
f.committed = append(f.committed, req.VolumeId)
|
||||
return &volume_server_pb.VacuumVolumeCommitResponse{IsReadOnly: true}, nil
|
||||
}
|
||||
|
||||
func (f *fakeVacuumServer) VacuumVolumeCleanup(_ context.Context, _ *volume_server_pb.VacuumVolumeCleanupRequest) (*volume_server_pb.VacuumVolumeCleanupResponse, error) {
|
||||
return &volume_server_pb.VacuumVolumeCleanupResponse{}, nil
|
||||
}
|
||||
|
||||
// A sweep keeps a read-only replica eligible only when the volume server
|
||||
// reports disk_space_low; other read-only causes stay skipped unless the
|
||||
// request names the volume explicitly.
|
||||
func TestVacuumReadOnlyDiskLowVolume(t *testing.T) {
|
||||
fake := &fakeVacuumServer{
|
||||
checks: map[uint32]*volume_server_pb.VacuumVolumeCheckResponse{
|
||||
1: {GarbageRatio: 0.9, DiskSpaceLow: true},
|
||||
2: {GarbageRatio: 0.9},
|
||||
3: {GarbageRatio: 0.9},
|
||||
4: {GarbageRatio: 0.1, DiskSpaceLow: true},
|
||||
5: {GarbageRatio: 0.9},
|
||||
},
|
||||
}
|
||||
grpcPort, dialOption := startFakeVolumeServer(t, fake)
|
||||
|
||||
topo := NewTopology("weedfs", sequence.NewMemorySequencer(), 32*1024, 5, false)
|
||||
dn := topo.GetOrCreateDataCenter("dc1").GetOrCreateRack("rack1").
|
||||
GetOrCreateDataNode("127.0.0.1", 8080, grpcPort, "127.0.0.1", "dn", map[string]uint32{"": 10})
|
||||
|
||||
newVolumeInfo := func(vid needle.VolumeId, readOnly bool) storage.VolumeInfo {
|
||||
return storage.VolumeInfo{
|
||||
Id: vid,
|
||||
Size: 1 << 20,
|
||||
Collection: "c",
|
||||
ReadOnly: readOnly,
|
||||
ModifiedAtSecond: time.Now().Unix(),
|
||||
Version: needle.GetCurrentVersion(),
|
||||
ReplicaPlacement: &super_block.ReplicaPlacement{},
|
||||
Ttl: needle.EMPTY_TTL,
|
||||
}
|
||||
}
|
||||
vl := topo.GetVolumeLayout("c", &super_block.ReplicaPlacement{}, needle.EMPTY_TTL, types.ToDiskType(""))
|
||||
volumes := []storage.VolumeInfo{
|
||||
newVolumeInfo(1, true),
|
||||
newVolumeInfo(2, true),
|
||||
newVolumeInfo(3, false),
|
||||
newVolumeInfo(4, true),
|
||||
newVolumeInfo(5, true),
|
||||
}
|
||||
dn.UpdateVolumes(volumes)
|
||||
for _, v := range volumes {
|
||||
topo.RegisterVolumeLayout(v, dn)
|
||||
}
|
||||
c := NewCollection("c", topo.volumeSizeLimit, false)
|
||||
|
||||
vacuum := func(vid needle.VolumeId, skipReadOnly bool) {
|
||||
vl.accessLock.RLock()
|
||||
ll := vl.vid2location[vid].Copy()
|
||||
vl.accessLock.RUnlock()
|
||||
topo.vacuumOneVolumeId(dialOption, vl, c, 0.3, ll, vid, 0, skipReadOnly)
|
||||
}
|
||||
|
||||
vacuum(1, true) // read-only but disk_low: compacted
|
||||
vacuum(2, true) // read-only otherwise: skipped by sweep
|
||||
vacuum(3, true) // writable: compacted
|
||||
vacuum(4, true) // disk_low but below threshold: skipped
|
||||
vacuum(5, false) // read-only, explicit request: compacted
|
||||
|
||||
fake.mu.Lock()
|
||||
defer fake.mu.Unlock()
|
||||
want := map[uint32]bool{1: true, 3: true, 5: true}
|
||||
got := map[uint32]bool{}
|
||||
for _, vid := range fake.committed {
|
||||
got[vid] = true
|
||||
}
|
||||
for vid := range want {
|
||||
if !got[vid] {
|
||||
t.Errorf("volume %d was not vacuum-committed", vid)
|
||||
}
|
||||
}
|
||||
for _, vid := range fake.committed {
|
||||
if !want[vid] {
|
||||
t.Errorf("volume %d vacuum-committed unexpectedly", vid)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user