Skip to content

Commit

Permalink
Fix GenericAvroSerializer and address comments
Browse files Browse the repository at this point in the history
  • Loading branch information
zsxwing committed Dec 2, 2015
1 parent bfb3360 commit a5d965c
Show file tree
Hide file tree
Showing 2 changed files with 6 additions and 4 deletions.
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,7 @@

package org.apache.spark.serializer

import java.io.ByteArrayOutputStream
import java.io.{ByteArrayInputStream, ByteArrayOutputStream}
import java.nio.ByteBuffer

import scala.collection.mutable
Expand All @@ -31,7 +31,6 @@ import org.apache.commons.io.IOUtils

import org.apache.spark.{SparkException, SparkEnv}
import org.apache.spark.io.CompressionCodec
import org.apache.spark.util.ByteBufferInputStream

/**
* Custom serializer used for generic Avro records. If the user registers the schemas
Expand Down Expand Up @@ -82,7 +81,10 @@ private[serializer] class GenericAvroSerializer(schemas: Map[Long, String])
* seen values so to limit the number of times that decompression has to be done.
*/
def decompress(schemaBytes: ByteBuffer): Schema = decompressCache.getOrElseUpdate(schemaBytes, {
val bis = new ByteBufferInputStream(schemaBytes)
val bis = new ByteArrayInputStream(
schemaBytes.array(),
schemaBytes.arrayOffset() + schemaBytes.position(),
schemaBytes.remaining())
val bytes = IOUtils.toByteArray(codec.compressedInputStream(bis))
new Schema.Parser().parse(new String(bytes, "UTF-8"))
})
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -332,7 +332,7 @@ private void decodeBinaryBatch(int col, int num) throws IOException {
for (int n = 0; n < num; ++n) {
if (columnReaders[col].next()) {
ByteBuffer bytes = columnReaders[col].nextBinary().toByteBuffer();
int len = bytes.limit() - bytes.position();
int len = bytes.remaining();
if (originalTypes[col] == OriginalType.UTF8) {
UTF8String str =
UTF8String.fromBytes(bytes.array(), bytes.arrayOffset() + bytes.position(), len);
Expand Down

0 comments on commit a5d965c

Please sign in to comment.