From 7ff684c193dd6c0fa71b8ea9b0d0141e17adc3dd Mon Sep 17 00:00:00 2001 From: Ben Johnson Date: Fri, 31 May 2019 17:29:11 -0600 Subject: [PATCH] Allow partial translate file reads. This commit fixes an issue where translation `LogEntry` must be read in its entirety, however, large entries can exceed the buffer size. This has been changed so that partial entries reads are allowed. The `LogEntry.ReadFrom()` may still generate large byte slices during reads of large individual fields or keys. --- http/handler.go | 14 ---------- translate.go | 67 ++++++++++++----------------------------------- translate_test.go | 39 +++++++++++++++++++++++++++ 3 files changed, 56 insertions(+), 64 deletions(-) diff --git a/http/handler.go b/http/handler.go index 8e932271a..aea500582 100644 --- a/http/handler.go +++ b/http/handler.go @@ -1485,10 +1485,6 @@ type defaultClusterMessageResponse struct{} // translateStoreBufferSize is the buffer size used for streaming data. const translateStoreBufferSize = 1 << 16 // 64k -// translateStoreBufferSizeMax is the maximum size that the buffer is allowed -// to grow before raising an error. -const translateStoreBufferSizeMax = 1 << 22 // 4Mb - func (h *Handler) handleGetTranslateData(w http.ResponseWriter, r *http.Request) { q := r.URL.Query() offset, _ := strconv.ParseInt(q.Get("offset"), 10, 64) @@ -1517,16 +1513,6 @@ func (h *Handler) handleGetTranslateData(w http.ResponseWriter, r *http.Request) n, err := rdr.Read(buf) if err == io.EOF { return - } else if err == pilosa.ErrTranslateReadTargetUndersized { - // Increase the buffer size and try to read again. - useBufferSize *= 2 - // Prevent the buffer from growing without bound. - if useBufferSize > translateStoreBufferSizeMax { - h.logger.Printf("http: translate store buffer exceeded max size: %s", err) - return - } - buf = make([]byte, useBufferSize) - continue } else if err != nil { h.logger.Printf("http: translate store read error: %s", err) return diff --git a/translate.go b/translate.go index 94eb57186..1b6eb6244 100644 --- a/translate.go +++ b/translate.go @@ -45,11 +45,10 @@ const ( // Translate store errors. var ( - ErrTranslateStoreClosed = errors.New("pilosa: translate store closed") - ErrTranslateStoreReaderClosed = errors.New("pilosa: translate store reader closed") - ErrReplicationNotSupported = errors.New("pilosa: replication not supported") - ErrTranslateStoreReadOnly = errors.New("pilosa: translate store could not find or create key, translate store read only") - ErrTranslateReadTargetUndersized = errors.New("pilosa: translate read target is undersized") + ErrTranslateStoreClosed = errors.New("pilosa: translate store closed") + ErrTranslateStoreReaderClosed = errors.New("pilosa: translate store reader closed") + ErrReplicationNotSupported = errors.New("pilosa: replication not supported") + ErrTranslateStoreReadOnly = errors.New("pilosa: translate store could not find or create key, translate store read only") ) // TranslateStore is the storage for translation string-to-uint64 values. @@ -746,48 +745,38 @@ func (e *LogEntry) ReadFrom(r io.Reader) (_ int64, err error) { return int64(uVarintSize(e.Length)), err } - // Slurp entire entry and replace reader. - buf := make([]byte, e.Length) - n, err := io.ReadFull(r, buf) - n64 := int64(n + uVarintSize(e.Length)) - if err != nil { - return n64, err - } - bufr := bytes.NewReader(buf) - br, r = bufr, bufr - // Read the entry type. if err := binary.Read(r, binary.BigEndian, &e.Type); err != nil { - return n64, err + return 0, err } // Read index name. if sz, err := binary.ReadUvarint(br); err != nil { - return n64, err + return 0, err } else if sz == 0 { e.Index = nil } else { e.Index = make([]byte, sz) if _, err := io.ReadFull(r, e.Index); err != nil { - return n64, err + return 0, err } } // Read field name. if sz, err := binary.ReadUvarint(br); err != nil { - return n64, err + return 0, err } else if sz == 0 { e.Field = nil } else { e.Field = make([]byte, sz) if _, err := io.ReadFull(r, e.Field); err != nil { - return n64, err + return 0, err } } // Read key count. if n, err := binary.ReadUvarint(br); err != nil { - return n64, err + return 0, err } else if n == 0 { e.IDs, e.Keys = nil, nil } else { @@ -798,20 +787,21 @@ func (e *LogEntry) ReadFrom(r io.Reader) (_ int64, err error) { for i := range e.Keys { // Read identifier. if e.IDs[i], err = binary.ReadUvarint(br); err != nil { - return n64, err + return 0, err } // Read key. if sz, err := binary.ReadUvarint(br); err != nil { - return n64, err + return 0, err } else if sz > 0 { e.Keys[i] = make([]byte, sz) if _, err := io.ReadFull(r, e.Keys[i]); err != nil { - return n64, err + return 0, err } } } - return n64, nil + + return int64(uVarintSize(e.Length)) + int64(e.Length), nil } // WriteTo serializes a LogEntry to w. @@ -875,22 +865,6 @@ func (e *LogEntry) WriteTo(w io.Writer) (_ int64, err error) { return int64(sz) + n, err } -// validLogEntriesLen returns the maximum length of p that contains valid entries. -func validLogEntriesLen(p []byte) (n int) { - r := bytes.NewReader(p) - for { - if sz, err := binary.ReadUvarint(r); err != nil { - return n - } else if off, err := r.Seek(int64(sz), io.SeekCurrent); err != nil { - return n - } else if off > int64(len(p)) { - return n - } else { - n = int(off) - } - } -} - type fieldKey struct { index string field string @@ -1126,7 +1100,7 @@ func (r *translateFileReader) Read(p []byte) (n int, err error) { } } -// read writes the bytes for zero or more valid entries to p. +// read reads up to len(p) bytes into p. func (r *translateFileReader) read(p []byte) (n int, err error) { sz := r.store.size() @@ -1137,20 +1111,13 @@ func (r *translateFileReader) read(p []byte) (n int, err error) { return 0, nil } - if max := sz - r.offset; max > int64(len(p)) { - // If p is not large enough to hold a single entry, - // return an error so the client can increase the - // size of p and try again. - return 0, ErrTranslateReadTargetUndersized - } else if int64(len(p)) > max { + if max := sz - r.offset; int64(len(p)) > max { // Shorten buffer to maximum read size. p = p[:max] } // Read data from file at offset. - // Limit the number of bytes read to only whole entries. n, err = r.file.ReadAt(p, r.offset) - n = validLogEntriesLen(p[:n]) r.offset += int64(n) return n, err } diff --git a/translate_test.go b/translate_test.go index 4fa977e62..344a820d5 100644 --- a/translate_test.go +++ b/translate_test.go @@ -370,6 +370,45 @@ func TestTranslateFile_Reader(t *testing.T) { t.Fatal(diff) } }) + + t.Run("TinyBuffer", func(t *testing.T) { + stringKeys := make([]string, 1024) + byteSliceKeys := make([][]byte, len(stringKeys)) + ids := make([]uint64, len(stringKeys)) + for i := range stringKeys { + stringKeys[i] = fmt.Sprintf("KEY%d", i) + byteSliceKeys[i] = []byte(stringKeys[i]) + ids[i] = uint64(i + 1) + } + + s := MustOpenTranslateFile() + defer s.MustClose() + if _, err := s.TranslateColumnsToUint64("IDX0", stringKeys); err != nil { + t.Fatal(err) + } + + // Obtain the reader and use the smallest possible buffer for bufio. + rc, err := s.Reader(context.Background(), 0) + if err != nil { + t.Fatal(err) + } + brc := bufio.NewReaderSize(rc, 16) + defer rc.Close() + + // Record should be able to be read using multiple reads. + var entry pilosa.LogEntry + if _, err := entry.ReadFrom(brc); err != nil { + t.Fatal(err) + } else if diff := cmp.Diff(entry, pilosa.LogEntry{ + Type: pilosa.LogEntryTypeInsertColumn, + Index: []byte("IDX0"), + IDs: ids, + Keys: byteSliceKeys, + Length: 9012, + }); diff != "" { + t.Fatal(diff) + } + }) } func TestPrintTranslateFile(t *testing.T) {