write_batch.cc 6.4 KB
Newer Older
J
jorlow@chromium.org 已提交
1 2 3 4 5 6 7 8 9
// Copyright (c) 2011 The LevelDB Authors. All rights reserved.
// Use of this source code is governed by a BSD-style license that can be
// found in the LICENSE file. See the AUTHORS file for names of contributors.
//
// WriteBatch::rep_ :=
//    sequence: fixed64
//    count: fixed32
//    data: record[count]
// record :=
10 11
//    kTypeValue varstring varstring
//    kTypeMerge varstring varstring
J
jorlow@chromium.org 已提交
12 13 14 15 16
//    kTypeDeletion varstring
// varstring :=
//    len: varint32
//    data: uint8[len]

17
#include "rocksdb/write_batch.h"
J
jorlow@chromium.org 已提交
18

19 20
#include "rocksdb/options.h"
#include "rocksdb/statistics.h"
J
jorlow@chromium.org 已提交
21
#include "db/dbformat.h"
22
#include "db/db_impl.h"
J
jorlow@chromium.org 已提交
23
#include "db/memtable.h"
24
#include "db/snapshot.h"
J
jorlow@chromium.org 已提交
25 26
#include "db/write_batch_internal.h"
#include "util/coding.h"
27
#include <stdexcept>
J
jorlow@chromium.org 已提交
28 29 30

namespace leveldb {

31 32 33
// WriteBatch header has an 8-byte sequence number followed by a 4-byte count.
static const size_t kHeader = 12;

J
jorlow@chromium.org 已提交
34 35 36 37 38 39
WriteBatch::WriteBatch() {
  Clear();
}

WriteBatch::~WriteBatch() { }

40 41
WriteBatch::Handler::~Handler() { }

42 43 44 45
void WriteBatch::Handler::Merge(const Slice& key, const Slice& value) {
  throw std::runtime_error("Handler::Merge not implemented!");
}

J
Jim Paton 已提交
46 47 48 49 50
void WriteBatch::Handler::LogData(const Slice& blob) {
  // If the user has not specified something to do with blobs, then we ignore
  // them.
}

51 52 53 54
bool WriteBatch::Handler::Continue() {
  return true;
}

J
jorlow@chromium.org 已提交
55 56
void WriteBatch::Clear() {
  rep_.clear();
57
  rep_.resize(kHeader);
J
jorlow@chromium.org 已提交
58 59
}

H
Haobo Xu 已提交
60 61 62 63
int WriteBatch::Count() const {
  return WriteBatchInternal::Count(this);
}

64 65
Status WriteBatch::Iterate(Handler* handler) const {
  Slice input(rep_);
66
  if (input.size() < kHeader) {
67 68 69
    return Status::Corruption("malformed WriteBatch (too small)");
  }

70
  input.remove_prefix(kHeader);
J
Jim Paton 已提交
71
  Slice key, value, blob;
72
  int found = 0;
73
  while (!input.empty() && handler->Continue()) {
74 75 76 77 78 79 80
    char tag = input[0];
    input.remove_prefix(1);
    switch (tag) {
      case kTypeValue:
        if (GetLengthPrefixedSlice(&input, &key) &&
            GetLengthPrefixedSlice(&input, &value)) {
          handler->Put(key, value);
J
Jim Paton 已提交
81
          found++;
82 83 84 85 86 87 88
        } else {
          return Status::Corruption("bad WriteBatch Put");
        }
        break;
      case kTypeDeletion:
        if (GetLengthPrefixedSlice(&input, &key)) {
          handler->Delete(key);
J
Jim Paton 已提交
89
          found++;
90 91 92 93
        } else {
          return Status::Corruption("bad WriteBatch Delete");
        }
        break;
94 95 96 97
      case kTypeMerge:
        if (GetLengthPrefixedSlice(&input, &key) &&
            GetLengthPrefixedSlice(&input, &value)) {
          handler->Merge(key, value);
J
Jim Paton 已提交
98
          found++;
99 100 101 102
        } else {
          return Status::Corruption("bad WriteBatch Merge");
        }
        break;
J
Jim Paton 已提交
103 104 105 106 107 108 109
      case kTypeLogData:
        if (GetLengthPrefixedSlice(&input, &blob)) {
          handler->LogData(blob);
        } else {
          return Status::Corruption("bad WriteBatch Blob");
        }
        break;
110 111 112 113 114 115 116 117 118 119 120
      default:
        return Status::Corruption("unknown WriteBatch tag");
    }
  }
  if (found != WriteBatchInternal::Count(this)) {
    return Status::Corruption("WriteBatch has wrong count");
  } else {
    return Status::OK();
  }
}

J
jorlow@chromium.org 已提交
121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149
int WriteBatchInternal::Count(const WriteBatch* b) {
  return DecodeFixed32(b->rep_.data() + 8);
}

void WriteBatchInternal::SetCount(WriteBatch* b, int n) {
  EncodeFixed32(&b->rep_[8], n);
}

SequenceNumber WriteBatchInternal::Sequence(const WriteBatch* b) {
  return SequenceNumber(DecodeFixed64(b->rep_.data()));
}

void WriteBatchInternal::SetSequence(WriteBatch* b, SequenceNumber seq) {
  EncodeFixed64(&b->rep_[0], seq);
}

void WriteBatch::Put(const Slice& key, const Slice& value) {
  WriteBatchInternal::SetCount(this, WriteBatchInternal::Count(this) + 1);
  rep_.push_back(static_cast<char>(kTypeValue));
  PutLengthPrefixedSlice(&rep_, key);
  PutLengthPrefixedSlice(&rep_, value);
}

void WriteBatch::Delete(const Slice& key) {
  WriteBatchInternal::SetCount(this, WriteBatchInternal::Count(this) + 1);
  rep_.push_back(static_cast<char>(kTypeDeletion));
  PutLengthPrefixedSlice(&rep_, key);
}

150 151 152 153 154 155 156
void WriteBatch::Merge(const Slice& key, const Slice& value) {
  WriteBatchInternal::SetCount(this, WriteBatchInternal::Count(this) + 1);
  rep_.push_back(static_cast<char>(kTypeMerge));
  PutLengthPrefixedSlice(&rep_, key);
  PutLengthPrefixedSlice(&rep_, value);
}

J
Jim Paton 已提交
157 158 159 160
void WriteBatch::PutLogData(const Slice& blob) {
  rep_.push_back(static_cast<char>(kTypeLogData));
  PutLengthPrefixedSlice(&rep_, blob);
}
161

162 163 164 165 166
namespace {
class MemTableInserter : public WriteBatch::Handler {
 public:
  SequenceNumber sequence_;
  MemTable* mem_;
167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183
  const Options* options_;
  DBImpl* db_;
  const bool filter_deletes_;

  MemTableInserter(SequenceNumber sequence, MemTable* mem, const Options* opts,
                   DB* db, const bool filter_deletes)
    : sequence_(sequence),
      mem_(mem),
      options_(opts),
      db_(reinterpret_cast<DBImpl*>(db)),
      filter_deletes_(filter_deletes) {
    assert(mem_);
    if (filter_deletes_) {
      assert(options_);
      assert(db_);
    }
  }
184 185 186 187

  virtual void Put(const Slice& key, const Slice& value) {
    mem_->Add(sequence_, kTypeValue, key, value);
    sequence_++;
J
jorlow@chromium.org 已提交
188
  }
189 190 191 192
  virtual void Merge(const Slice& key, const Slice& value) {
    mem_->Add(sequence_, kTypeMerge, key, value);
    sequence_++;
  }
193
  virtual void Delete(const Slice& key) {
194 195 196 197 198 199 200 201 202 203
    if (filter_deletes_) {
      SnapshotImpl read_from_snapshot;
      read_from_snapshot.number_ = sequence_;
      ReadOptions ropts;
      ropts.snapshot = &read_from_snapshot;
      std::string value;
      if (!db_->KeyMayExist(ropts, key, &value)) {
        RecordTick(options_->statistics, NUMBER_FILTERED_DELETES);
        return;
      }
204
    }
205 206
    mem_->Add(sequence_, kTypeDeletion, key, Slice());
    sequence_++;
J
jorlow@chromium.org 已提交
207
  }
208
};
H
Hans Wennborg 已提交
209
}  // namespace
210

211 212 213 214 215
Status WriteBatchInternal::InsertInto(const WriteBatch* b, MemTable* mem,
                                      const Options* opts, DB* db,
                                      const bool filter_deletes) {
  MemTableInserter inserter(WriteBatchInternal::Sequence(b), mem, opts, db,
                            filter_deletes);
216
  return b->Iterate(&inserter);
J
jorlow@chromium.org 已提交
217 218 219
}

void WriteBatchInternal::SetContents(WriteBatch* b, const Slice& contents) {
220
  assert(contents.size() >= kHeader);
J
jorlow@chromium.org 已提交
221 222 223
  b->rep_.assign(contents.data(), contents.size());
}

224 225 226 227 228 229
void WriteBatchInternal::Append(WriteBatch* dst, const WriteBatch* src) {
  SetCount(dst, Count(dst) + Count(src));
  assert(src->rep_.size() >= kHeader);
  dst->rep_.append(src->rep_.data() + kHeader, src->rep_.size() - kHeader);
}

H
Hans Wennborg 已提交
230
}  // namespace leveldb