MCPcopy Create free account
hub / github.com/apache/incubator-pegasus / write_batch_put_ctx

Method write_batch_put_ctx

src/server/rocksdb_wrapper.cpp:112–164  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

110}
111
112int rocksdb_wrapper::write_batch_put_ctx(const db_write_context &ctx,
113 dsn::string_view raw_key,
114 dsn::string_view value,
115 uint32_t expire_sec)
116{
117 FAIL_POINT_INJECT_F("db_write_batch_put",
118 [](dsn::string_view) -> int { return FAIL_DB_WRITE_BATCH_PUT; });
119
120 uint64_t new_timetag = ctx.remote_timetag;
121 if (!ctx.is_duplicated_write()) { // local write
122 new_timetag = generate_timetag(ctx.timestamp, get_cluster_id_if_exists(), false);
123 }
124
125 if (ctx.verify_timetag && // needs read-before-write
126 _pegasus_data_version >= 1 && // data version 0 doesn't support timetag.
127 !raw_key.empty()) { // not an empty write
128
129 db_get_context get_ctx;
130 int err = get(raw_key, &get_ctx);
131 if (dsn_unlikely(err != rocksdb::Status::kOk)) {
132 return err;
133 }
134 // if record exists and is not expired.
135 if (get_ctx.found && !get_ctx.expired) {
136 uint64_t local_timetag =
137 pegasus_extract_timetag(_pegasus_data_version, get_ctx.raw_value);
138
139 if (local_timetag >= new_timetag) {
140 // ignore this stale update with lower timetag,
141 // and write an empty record instead
142 raw_key = value = dsn::string_view();
143 }
144 }
145 }
146
147 rocksdb::Slice skey = utils::to_rocksdb_slice(raw_key);
148 rocksdb::SliceParts skey_parts(&skey, 1);
149 rocksdb::SliceParts svalue = _value_generator->generate_value(
150 _pegasus_data_version, value, db_expire_ts(expire_sec), new_timetag);
151 rocksdb::Status s = _write_batch->Put(skey_parts, svalue);
152 if (dsn_unlikely(!s.ok())) {
153 ::dsn::blob hash_key, sort_key;
154 pegasus_restore_key(::dsn::blob(raw_key.data(), 0, raw_key.size()), hash_key, sort_key);
155 LOG_ERROR_ROCKSDB("WriteBatchPut",
156 s.ToString(),
157 "decree: {}, hash_key: {}, sort_key: {}, expire_ts: {}",
158 ctx.decree,
159 utils::c_escape_string(hash_key),
160 utils::c_escape_string(sort_key),
161 expire_sec);
162 }
163 return s.code();
164}
165
166int rocksdb_wrapper::write(int64_t decree)
167{

Callers 3

multi_putMethod · 0.80
batch_putMethod · 0.80
single_setMethod · 0.80

Calls 15

generate_timetagFunction · 0.85
get_cluster_id_if_existsFunction · 0.85
pegasus_extract_timetagFunction · 0.85
to_rocksdb_sliceFunction · 0.85
pegasus_restore_keyFunction · 0.85
is_duplicated_writeMethod · 0.80
PutMethod · 0.80
okMethod · 0.80
blobClass · 0.70
string_viewClass · 0.50
c_escape_stringFunction · 0.50
emptyMethod · 0.45

Tested by 1

single_setMethod · 0.64