bk99.de entertain the web since 1997

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

Read the original article

Search the Web Archive