Processing four terabytes of newspaper data with EC2 and Hadoop
Summary
Derek Gottfrid describes how the New York Times turned four terabytes of archive data into usable PDF files with Hadoop, EC2 and S3. The source data comprised about four terabytes of TIFF files and metadata. The job ran on a hundred Amazon EC2 instances.
Ideas
- Temporarily rented computing power replaces permanently oversized infrastructure.
- Hadoop distributes independent conversion tasks across many short-lived instances.
- S3 decouples input data, the computing cluster and finished results.
- A reproducible job processes decades of archive holdings in one run.
- Parallelisation shortens a multi-week single-machine task to about one day.
- Usage-based billing makes one-off large projects economically predictable.
Insights
- Elastic infrastructure changes not only costs but also feasible project sizes.
- Data locality and independent work packages determine how well a batch job scales.
- Object storage creates a stable boundary between transient computing power and permanent results.
- Rare peak loads justify rented capacity more than owning hardware.
Facts
- Processing took less than 24 hours.
- The result consisted of about 1.5 terabytes of PDF files.
Recommendations
- Split batch processing into repeatable, independent units.
- Store intermediate results outside transient compute nodes.
- Calculate transfer, storage and failed retries together.
References
Links to the original source and the Web Archive open in a new tab.