Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
102 changes: 73 additions & 29 deletions internal/compaction/job.go
Original file line number Diff line number Diff line change
Expand Up @@ -120,25 +120,28 @@ type Job struct {
DurationSeconds float64

// Internal
logger zerolog.Logger
mu sync.Mutex
db *sql.DB // Shared DuckDB connection
compactedFiles []string // Files that were actually compacted (valid files only)
logger zerolog.Logger
mu sync.Mutex
db *sql.DB // Shared DuckDB connection
compactedFiles []string // Files that were actually compacted (valid files only)
manifestManager *ManifestManager
manifestPath string // Path to the manifest file for this job
}

// JobConfig holds configuration for creating a compaction job
type JobConfig struct {
Measurement string
PartitionPath string
Files []string
StorageBackend storage.Backend
Database string
TargetSizeMB int
Tier string
TempDirectory string // Base temp directory for compaction files (default: ./data/compaction)
SortKeys []string // Sort keys for this measurement (for ORDER BY in compaction)
Logger zerolog.Logger
DB *sql.DB // Shared DuckDB connection (avoids memory retention from temp connections)
Measurement string
PartitionPath string
Files []string
StorageBackend storage.Backend
Database string
TargetSizeMB int
Tier string
TempDirectory string // Base temp directory for compaction files (default: ./data/compaction)
SortKeys []string // Sort keys for this measurement (for ORDER BY in compaction)
Logger zerolog.Logger
DB *sql.DB // Shared DuckDB connection (avoids memory retention from temp connections)
ManifestManager *ManifestManager // Manifest manager for crash recovery (optional, recommended)
}

// NewJob creates a new compaction job
Expand All @@ -164,19 +167,20 @@ func NewJob(cfg *JobConfig) *Job {
}

return &Job{
Measurement: cfg.Measurement,
PartitionPath: cfg.PartitionPath,
Files: cfg.Files,
StorageBackend: cfg.StorageBackend,
Database: cfg.Database,
TargetSizeMB: cfg.TargetSizeMB,
Tier: cfg.Tier,
TempDirectory: tempDir,
SortKeys: sortKeys,
JobID: jobID,
Status: JobStatusPending,
logger: cfg.Logger.With().Str("job_id", jobID).Logger(),
db: cfg.DB,
Measurement: cfg.Measurement,
PartitionPath: cfg.PartitionPath,
Files: cfg.Files,
StorageBackend: cfg.StorageBackend,
Database: cfg.Database,
TargetSizeMB: cfg.TargetSizeMB,
Tier: cfg.Tier,
TempDirectory: tempDir,
SortKeys: sortKeys,
JobID: jobID,
Status: JobStatusPending,
logger: cfg.Logger.With().Str("job_id", jobID).Logger(),
db: cfg.DB,
manifestManager: cfg.ManifestManager,
}
}

Expand Down Expand Up @@ -242,14 +246,54 @@ func (j *Job) Run(ctx context.Context) error {

// Upload compacted file
compactedKey := filepath.Join(j.PartitionPath, filepath.Base(compactedFile))

// Write manifest BEFORE upload to enable crash recovery
// If we crash after upload but before deletion, the manifest allows recovery
if j.manifestManager != nil {
manifest := &Manifest{
OutputPath: compactedKey,
OutputSize: j.BytesAfter,
InputFiles: j.compactedFiles,
Database: j.Database,
Measurement: j.Measurement,
PartitionPath: j.PartitionPath,
Tier: j.Tier,
Status: ManifestStatusPending,
CreatedAt: time.Now().UTC(),
JobID: j.JobID,
}

manifestPath, err := j.manifestManager.WriteManifest(ctx, manifest)
if err != nil {
return j.fail(fmt.Errorf("failed to write manifest: %w", err))
}
j.manifestPath = manifestPath
j.logger.Debug().Str("manifest", manifestPath).Msg("Wrote compaction manifest")
}

if err := j.uploadFile(ctx, compactedFile, compactedKey); err != nil {
// Upload failed - delete manifest since output doesn't exist
if j.manifestManager != nil && j.manifestPath != "" {
if delErr := j.manifestManager.DeleteManifest(ctx, j.manifestPath); delErr != nil {
j.logger.Warn().Err(delErr).Msg("Failed to delete manifest after upload failure")
}
}
return j.fail(fmt.Errorf("failed to upload compacted file: %w", err))
}

// Delete old files from storage
if err := j.deleteOldFiles(ctx); err != nil {
j.logger.Warn().Err(err).Msg("Failed to delete some old files")
// Don't fail the job for deletion errors
// Don't fail the job - manifest will enable recovery on next cycle
// The manifest remains so recovery can retry deletion
} else {
// Deletion succeeded - delete the manifest
if j.manifestManager != nil && j.manifestPath != "" {
if delErr := j.manifestManager.DeleteManifest(ctx, j.manifestPath); delErr != nil {
j.logger.Warn().Err(delErr).Msg("Failed to delete manifest after successful deletion")
// Non-fatal - manifest will be cleaned up during recovery
}
}
}

// Cleanup empty directories (best-effort, local storage only)
Expand Down
106 changes: 89 additions & 17 deletions internal/compaction/manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,8 +18,9 @@ var ErrCycleAlreadyRunning = errors.New("compaction cycle already running")

// Manager orchestrates compaction jobs across all measurements
type Manager struct {
StorageBackend storage.Backend
LockManager *LockManager
StorageBackend storage.Backend
LockManager *LockManager
ManifestManager *ManifestManager

// Configuration
MinAgeHours int
Expand All @@ -44,10 +45,11 @@ type Manager struct {
cycleID atomic.Int64

// Metrics
TotalJobsCompleted int
TotalJobsFailed int
TotalFilesCompacted int
TotalBytesSaved int64
TotalJobsCompleted int
TotalJobsFailed int
TotalFilesCompacted int
TotalBytesSaved int64
TotalManifestsRecov int // Number of manifests recovered

logger zerolog.Logger
mu sync.Mutex
Expand Down Expand Up @@ -99,9 +101,12 @@ func NewManager(cfg *ManagerConfig) *Manager {
defaultSortKeys = []string{"time"} // Default to time-only sorting
}

logger := cfg.Logger.With().Str("component", "compaction-manager").Logger()

m := &Manager{
StorageBackend: cfg.StorageBackend,
LockManager: cfg.LockManager,
ManifestManager: NewManifestManager(cfg.StorageBackend, logger),
MinAgeHours: cfg.MinAgeHours,
MinFiles: cfg.MinFiles,
TargetSizeMB: cfg.TargetSizeMB,
Expand All @@ -112,7 +117,7 @@ func NewManager(cfg *ManagerConfig) *Manager {
DefaultSortKeys: defaultSortKeys,
Tiers: cfg.Tiers,
jobHistory: make([]map[string]interface{}, 0),
logger: cfg.Logger.With().Str("component", "compaction-manager").Logger(),
logger: logger,
}

// Log tier information
Expand Down Expand Up @@ -442,6 +447,21 @@ func (m *Manager) RunCompactionCycleForTiers(ctx context.Context, tierNames []st
Strs("tiers", tierNames).
Msg("Starting compaction cycle for specific tiers")

// Run manifest recovery before starting new compactions
// This ensures interrupted compactions from previous cycles are completed
if m.ManifestManager != nil {
recovered, err := m.ManifestManager.RecoverOrphanedManifests(ctx)
if err != nil {
m.logger.Warn().Err(err).Msg("Manifest recovery encountered errors")
}
if recovered > 0 {
m.mu.Lock()
m.TotalManifestsRecov += recovered
m.mu.Unlock()
m.logger.Info().Int("recovered", recovered).Msg("Recovered orphaned compaction manifests")
}
}

// Build tier filter map for quick lookup
tierFilter := make(map[string]bool)
for _, name := range tierNames {
Expand Down Expand Up @@ -524,13 +544,22 @@ func (m *Manager) RunCompactionCycleForTiers(ctx context.Context, tierNames []st

// Process this measurement's candidates immediately
for _, candidate := range candidates {
// Filter out files that are tracked by manifests (pending compaction)
filteredCandidate, shouldProcess := m.filterCandidateFiles(ctx, candidate)
if !shouldProcess {
m.logger.Debug().
Str("partition", candidate.PartitionPath).
Msg("Skipping candidate: all files are tracked by manifests")
continue
}

// Split large candidates into batches to prevent DuckDB segfaults
// when processing too many files in a single read_parquet() call
batches := SplitCandidateIntoBatches(candidate)
batches := SplitCandidateIntoBatches(filteredCandidate)
if len(batches) > 1 {
m.logger.Info().
Str("partition", candidate.PartitionPath).
Int("total_files", len(candidate.Files)).
Str("partition", filteredCandidate.PartitionPath).
Int("total_files", len(filteredCandidate.Files)).
Int("batches", len(batches)).
Msg("Splitting large candidate into batches")
}
Expand Down Expand Up @@ -728,19 +757,62 @@ func (m *Manager) GetSortKeys(measurement string) []string {
return m.DefaultSortKeys
}

// filterCandidateFiles removes files that are tracked by manifests from a candidate.
// Returns the filtered candidate and whether it should still be processed.
func (m *Manager) filterCandidateFiles(ctx context.Context, candidate Candidate) (Candidate, bool) {
if m.ManifestManager == nil {
return candidate, len(candidate.Files) > 0
}

filesInManifests, err := m.ManifestManager.GetFilesInManifests(ctx)
if err != nil {
m.logger.Warn().Err(err).Msg("Failed to get files in manifests, proceeding without filtering")
return candidate, len(candidate.Files) > 0
}

if len(filesInManifests) == 0 {
return candidate, len(candidate.Files) > 0
}

// Filter out files that are in manifests
filteredFiles := make([]string, 0, len(candidate.Files))
var skipped int
for _, f := range candidate.Files {
if _, inManifest := filesInManifests[f]; inManifest {
skipped++
continue
}
filteredFiles = append(filteredFiles, f)
}

if skipped > 0 {
m.logger.Debug().
Str("partition", candidate.PartitionPath).
Int("skipped", skipped).
Int("remaining", len(filteredFiles)).
Msg("Filtered files tracked by manifests")
}

candidate.Files = filteredFiles
candidate.FileCount = len(filteredFiles)

return candidate, len(filteredFiles) > 0
}

// Stats returns compaction statistics
func (m *Manager) Stats() map[string]interface{} {
m.mu.Lock()
defer m.mu.Unlock()

stats := map[string]interface{}{
"total_jobs_completed": m.TotalJobsCompleted,
"total_jobs_failed": m.TotalJobsFailed,
"total_files_compacted": m.TotalFilesCompacted,
"total_bytes_saved": m.TotalBytesSaved,
"total_bytes_saved_mb": float64(m.TotalBytesSaved) / 1024 / 1024,
"cycle_running": m.cycleRunning.Load(),
"current_cycle_id": m.cycleID.Load(),
"total_jobs_completed": m.TotalJobsCompleted,
"total_jobs_failed": m.TotalJobsFailed,
"total_files_compacted": m.TotalFilesCompacted,
"total_bytes_saved": m.TotalBytesSaved,
"total_bytes_saved_mb": float64(m.TotalBytesSaved) / 1024 / 1024,
"total_manifests_recover": m.TotalManifestsRecov,
"cycle_running": m.cycleRunning.Load(),
"current_cycle_id": m.cycleID.Load(),
}

// Add recent jobs (last 10)
Expand Down
Loading