Skip to content

Better HDFS Support - #1556

Merged
echeipesh merged 21 commits into
locationtech:masterfrom
echeipesh:feature/better-hdfs-tile-reader
Jul 5, 2016
Merged

echeipesh merged 21 commits into
locationtech:masterfrom
echeipesh:feature/better-hdfs-tile-reader

Conversation

@echeipesh

@echeipesh echeipesh commented Jun 21, 2016 •

Copy link
Copy Markdown
Contributor

This PR improves HDFS layer writing and random value reading support.

HadoopValuesReader now opens up MapFile.Readers for each map file in the layer. These readers cache the index of available keys so they are able to provide quick lookups. This replaces and improves previous method of using FileInputFormat to query for a single record.

HadoopRDDWriter has multiple improvements:

  • Accomplishes it's task with a single shuffle step
    • This is possible because instead of estimating the blocks required by counting we simply roll over to a new file when the record being written is about to surpass the block boundary.
    • Instead of using a groupBy we sort our records on shuffle which uses IO index to partition the records
  • Writes happen through mapping the partition iterator, this is possible because the partition comes pre-sorted and greatly reduces memory pressure as records can be garbage collected after they are written.

Current testing shows:

  • HDFS ingest speeds exceed those seen in Accumulo
  • tile fetching latency is visually acceptable
  • that jobs using the GroupedCumulativeIterator approach are able to complete in settings where .groupBy/sort/write jobs are killed for memory violations.

Future improvements that on the mind:

  • Save bloom filter index for layer MapFiles to reduce memory requirements for HadoopValueReader and to produce quicker lookup in most cases.
  • In HadoopValueReader when fetching a record, cache all the tiles that share that index before filtering down to a single tile. These records are very likely to be asked for next and should be stored in LRU cache.
  • Devise a method to compare KeyIndex instances such that if the RDD to be saved is already partitioned by compatible index (where all key boundaries are shared) to avoid the extra shuffle.

MapFile.Writer.valueClass(classOf[BytesWritable]),
MapFile.Writer.compression(SequenceFile.CompressionType.NONE))
writer.setIndexInterval(1)
writer.setIndexInterval(32)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

how it would effect query time? (just curious)

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Layer query time would not be effected at all by this. When we're reading off ranges we're already seeking through the file, so not having as many points would have minimal impact if any. This is going to have more of an impact on random access through value reader. Spinning up a cluster to figure that part out. For reference the default interval value is 128

@echeipesh echeipesh changed the title [WIP] Better HDFS Support Better HDFS Support Jul 5, 2016
/**
* When record being written would exceed the block size of the current MapFile
* opens a new file to continue writing. This allows to split partition into block-sized
* chunks without foreknowledge of how bit it is.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

bit => big

@lossyrob

lossyrob commented Jul 5, 2016

Copy link
Copy Markdown
Member

+1 after comments addressed

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

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants