Skip to content

Bug: Compaction causes data duplication when deletion fails or pod restarts #157

Description

@khalid244

Description

Summary

The compaction process can cause data duplication when the deletion of source files fails after the compacted output file has been successfully uploaded. This occurs because the compaction operation is not atomic - there's a window between uploading the compacted file and deleting the source files where a failure (network error, pod restart, etc.) leaves both the original and compacted files in storage.

Observed Impact

I faced this issue when 400k rows of data became 800k rows in one day.

Impact

  • Data duplication in query results: Queries reading all parquet files will return duplicate rows
  • Storage waste: Both original and compacted files consume storage
  • Compounding duplication: Subsequent compaction cycles may re-compact the same files, creating additional duplicates

Root Cause

1. Non-Atomic Compaction Operation

The compaction job in job.go:183-259 performs these steps sequentially:

func (j *Job) Run(ctx context.Context) error {
    // 1. Download files
    // 2. Compact with DuckDB
    // 3. Upload compacted file  ← SUCCESS
    // 4. Delete source files    ← FAILURE POINT (pod crash, network error, etc.)
    
    if err := j.deleteOldFiles(ctx); err != nil {
        j.logger.Warn().Err(err).Msg("Failed to delete some old files")
        // Job still completes successfully!
    }
    return j.complete()
}

If step 4 fails after step 3 succeeds, both original and compacted files exist in storage.

2. Deletion Failure Doesn't Fail the Job

In job.go:250-253, deletion errors are logged as warnings but the job completes successfully:

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
}

3. Re-compaction Creates Additional Duplicates

When the next compaction cycle runs, ShouldCompact in tier.go:205-210 triggers re-compaction if source files still exist alongside compacted files:

// Case 2: Has compacted files, but many new uncompacted files accumulated
if len(compactedFiles) > 0 && len(uncompactedFiles) >= t.MinFiles {
    return true  // Triggers re-compaction of same source files!
}

This creates a second compacted file from the same source files, compounding the duplication.


Reproduction Scenario

Hourly Tier

  1. Partition has raw files: A.parquet, B.parquet, C.parquet (300 total rows)
  2. Hourly compaction runs → creates metrics_compacted.parquet (300 rows)
  3. Upload succeeds
  4. Pod crashes before deletion completes
  5. Storage now contains: A.parquet, B.parquet, C.parquet + metrics_compacted.parquet
  6. Result: 600 rows in storage (300 unique, 300 duplicated)

Daily Tier (Same Issue)

  1. Hourly-compacted files exist: hour03_compacted.parquet, hour04_compacted.parquet
  2. Daily compaction → creates metrics_daily.parquet
  3. Upload succeeds, deletion fails
  4. Storage has both _compacted.parquet files AND _daily.parquet
  5. Result: Data duplicated across tier outputs

Suggested Solution: Manifest-Based Tracking

Implement a manifest file system to track compaction state and enable crash recovery.

Design

Compaction Flow (Modified)

1. Download input files
2. Compact → create output file locally
3. Get output file size from local file
4. Write manifest file: {
     "output": "metrics_20250115_compacted.parquet",
     "output_size": 52428800,
     "inputs": ["A.parquet", "B.parquet", "C.parquet"],
     "created_at": "2025-01-15T10:00:00Z",
     "status": "pending"
   }
5. Upload output file
6. Delete input files
7. Delete manifest file (compaction complete)

Recovery Flow (On Startup / Each Cycle)

1. Scan for orphaned manifest files in _compaction_state/
2. For each manifest:
   a. Check if output file exists (use S3 HeadObject or equivalent)
   b. If output file missing → delete manifest, retry compaction
   c. If output file exists but size != manifest.output_size → delete partial output, delete manifest, retry compaction
   d. If output file exists and size matches → retry deletion of input files, then delete manifest
3. When finding candidates, exclude files listed in any manifest

Manifest Storage

Store manifests in the same storage backend under a dedicated path:

_compaction_state/
  hourly/
    {partition_path}_{timestamp}.json
  daily/
    {partition_path}_{timestamp}.json

State Transitions

Normal flow:
  (no manifest) → [write manifest] → pending → [upload] → [delete inputs] → [delete manifest] → (no manifest)

Crash recovery:
  pending + output missing        → delete manifest → retry compaction
  pending + output size mismatch  → delete partial output → delete manifest → retry compaction
  pending + output size matches   → delete inputs → delete manifest → done

Benefits

Benefit Description
Crash recovery Manifests persist across pod restarts, enabling deletion retry
Upload validation File size check detects partial/corrupted uploads
Idempotency Files listed in manifests are skipped in candidate scans
No re-compaction Source files won't be compacted again while manifest exists
Audit trail Manifests document what was compacted and when
Debuggability Easy to inspect compaction state via manifest files

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions