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
- Partition has raw files:
A.parquet, B.parquet, C.parquet (300 total rows)
- Hourly compaction runs → creates
metrics_compacted.parquet (300 rows)
- Upload succeeds
- Pod crashes before deletion completes
- Storage now contains:
A.parquet, B.parquet, C.parquet + metrics_compacted.parquet
- Result: 600 rows in storage (300 unique, 300 duplicated)
Daily Tier (Same Issue)
- Hourly-compacted files exist:
hour03_compacted.parquet, hour04_compacted.parquet
- Daily compaction → creates
metrics_daily.parquet
- Upload succeeds, deletion fails
- Storage has both
_compacted.parquet files AND _daily.parquet
- 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 |
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
Root Cause
1. Non-Atomic Compaction Operation
The compaction job in
job.go:183-259performs these steps sequentially: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:3. Re-compaction Creates Additional Duplicates
When the next compaction cycle runs,
ShouldCompactintier.go:205-210triggers re-compaction if source files still exist alongside compacted files:This creates a second compacted file from the same source files, compounding the duplication.
Reproduction Scenario
Hourly Tier
A.parquet,B.parquet,C.parquet(300 total rows)metrics_compacted.parquet(300 rows)A.parquet,B.parquet,C.parquet+metrics_compacted.parquetDaily Tier (Same Issue)
hour03_compacted.parquet,hour04_compacted.parquetmetrics_daily.parquet_compacted.parquetfiles AND_daily.parquetSuggested Solution: Manifest-Based Tracking
Implement a manifest file system to track compaction state and enable crash recovery.
Design
Compaction Flow (Modified)
Recovery Flow (On Startup / Each Cycle)
Manifest Storage
Store manifests in the same storage backend under a dedicated path:
State Transitions
Benefits