Skip to content

Commit b757e36

Browse files
wangshao1wangshaoyi
andauthored
fix: some streams errors such as pkpatternmatchdel etc (OpenAtomFoundation#2726)
* fix pkpatternmatchdel error --------- Co-authored-by: wangshaoyi <wangshaoyi@360.cn>
1 parent 0a87260 commit b757e36

5 files changed

Lines changed: 548 additions & 527 deletions

File tree

src/storage/src/base_filter.h

Lines changed: 21 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@
1616
#include "src/base_value_format.h"
1717
#include "src/base_meta_value_format.h"
1818
#include "src/lists_meta_value_format.h"
19+
#include "src/pika_stream_meta_value.h"
1920
#include "src/strings_value_format.h"
2021
#include "src/zsets_data_key_format.h"
2122
#include "src/debug.h"
@@ -36,11 +37,12 @@ class BaseMetaFilter : public rocksdb::CompactionFilter {
3637
* The field designs of the remaining zset,set,hash and stream in meta-value
3738
* are the same, so the same filtering strategy is used
3839
*/
40+
ParsedBaseKey parsed_key(key);
3941
auto type = static_cast<enum DataType>(static_cast<uint8_t>(value[0]));
4042
DEBUG("==========================START==========================");
4143
if (type == DataType::kStrings) {
4244
ParsedStringsValue parsed_strings_value(value);
43-
DEBUG("[StringsFilter] key: {}, value = {}, timestamp: {}, cur_time: {}", key.ToString().c_str(),
45+
DEBUG("[string type] key: %s, value = %s, timestamp: %llu, cur_time: %llu", parsed_key.Key().ToString().c_str(),
4446
parsed_strings_value.UserValue().ToString().c_str(), parsed_strings_value.Etime(), cur_time);
4547
if (parsed_strings_value.Etime() != 0 && parsed_strings_value.Etime() < cur_time) {
4648
DEBUG("Drop[Stale]");
@@ -49,9 +51,17 @@ class BaseMetaFilter : public rocksdb::CompactionFilter {
4951
DEBUG("Reserve");
5052
return false;
5153
}
54+
} else if (type == DataType::kStreams) {
55+
ParsedStreamMetaValue parsed_stream_meta_value(value);
56+
DEBUG("[stream meta type], key: %s, entries_added = %llu, first_id: %s, last_id: %s, version: %llu",
57+
parsed_key.Key().ToString().c_str(), parsed_stream_meta_value.entries_added(),
58+
parsed_stream_meta_value.first_id().ToString().c_str(),
59+
parsed_stream_meta_value.last_id().ToString().c_str(),
60+
parsed_stream_meta_value.version());
61+
return false;
5262
} else if (type == DataType::kLists) {
5363
ParsedListsMetaValue parsed_lists_meta_value(value);
54-
DEBUG("[ListMetaFilter], key: {}, count = {}, timestamp: {}, cur_time: {}, version: {}", key.ToString().c_str(),
64+
DEBUG("[list meta type], key: %s, count = %d, timestamp: %llu, cur_time: %llu, version: %llu", parsed_key.Key().ToString().c_str(),
5565
parsed_lists_meta_value.Count(), parsed_lists_meta_value.Etime(), cur_time,
5666
parsed_lists_meta_value.Version());
5767

@@ -68,8 +78,9 @@ class BaseMetaFilter : public rocksdb::CompactionFilter {
6878
return false;
6979
} else {
7080
ParsedBaseMetaValue parsed_base_meta_value(value);
71-
DEBUG("[MetaFilter] key: {}, count = {}, timestamp: {}, cur_time: {}, version: {}", key.ToString().c_str(),
72-
parsed_base_meta_value.Count(), parsed_base_meta_value.Etime(), cur_time, parsed_base_meta_value.Version());
81+
DEBUG("[%s meta type] key: %s, count = %d, timestamp: %llu, cur_time: %llu, version: %llu",
82+
DataTypeToString(type), parsed_key.Key().ToString().c_str(), parsed_base_meta_value.Count(),
83+
parsed_base_meta_value.Etime(), cur_time, parsed_base_meta_value.Version());
7384

7485
if (parsed_base_meta_value.Etime() != 0 && parsed_base_meta_value.Etime() < cur_time &&
7586
parsed_base_meta_value.Version() < cur_time) {
@@ -143,7 +154,12 @@ class BaseDataFilter : public rocksdb::CompactionFilter {
143154
auto type = static_cast<enum DataType>(static_cast<uint8_t>(meta_value[0]));
144155
if (type != type_) {
145156
return true;
146-
} else if (type == DataType::kHashes || type == DataType::kSets || type == DataType::kStreams || type == DataType::kZSets) {
157+
} else if (type == DataType::kStreams) {
158+
ParsedStreamMetaValue parsed_stream_meta_value(meta_value);
159+
meta_not_found_ = false;
160+
cur_meta_version_ = parsed_stream_meta_value.version();
161+
cur_meta_etime_ = 0; // stream do not support ttl
162+
} else if (type == DataType::kHashes || type == DataType::kSets || type == DataType::kZSets) {
147163
ParsedBaseMetaValue parsed_base_meta_value(&meta_value);
148164
meta_not_found_ = false;
149165
cur_meta_version_ = parsed_base_meta_value.Version();

src/storage/src/pika_stream_meta_value.h

Lines changed: 9 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -82,7 +82,8 @@ class StreamMetaValue {
8282
value_ = std::move(value);
8383
assert(value_.size() == kDefaultStreamValueLength);
8484
if (value_.size() != kDefaultStreamValueLength) {
85-
LOG(ERROR) << "Invalid stream meta value length: ";
85+
LOG(ERROR) << "Invalid stream meta value length: " << value_.size()
86+
<< " expected: " << kDefaultStreamValueLength;
8687
return;
8788
}
8889
char* pos = &value_[0];
@@ -215,7 +216,8 @@ class ParsedStreamMetaValue {
215216
ParsedStreamMetaValue(const Slice& value) {
216217
assert(value.size() == kDefaultStreamValueLength);
217218
if (value.size() != kDefaultStreamValueLength) {
218-
LOG(ERROR) << "Invalid stream meta value length: ";
219+
LOG(ERROR) << "Invalid stream meta value length: " << value.size()
220+
<< " expected: " << kDefaultStreamValueLength;
219221
return;
220222
}
221223
char* pos = const_cast<char*>(value.data());
@@ -294,7 +296,7 @@ class StreamCGroupMetaValue {
294296
uint64_t needed = kDefaultStreamCGroupValueLength;
295297
assert(value_.size() == 0);
296298
if (value_.size() != 0) {
297-
LOG(FATAL) << "Init on a existed stream cgroup meta value!";
299+
LOG(ERROR) << "Init on a existed stream cgroup meta value!";
298300
return;
299301
}
300302
value_.resize(needed);
@@ -314,7 +316,8 @@ class StreamCGroupMetaValue {
314316
value_ = std::move(value);
315317
assert(value_.size() == kDefaultStreamCGroupValueLength);
316318
if (value_.size() != kDefaultStreamCGroupValueLength) {
317-
LOG(FATAL) << "Invalid stream cgroup meta value length: ";
319+
LOG(ERROR) << "Invalid stream cgroup meta value length: " << value_.size()
320+
<< " expected: " << kDefaultStreamValueLength;
318321
return;
319322
}
320323
if (value_.size() == kDefaultStreamCGroupValueLength) {
@@ -373,7 +376,7 @@ class StreamConsumerMetaValue {
373376
value_ = std::move(value);
374377
assert(value_.size() == kDefaultStreamConsumerValueLength);
375378
if (value_.size() != kDefaultStreamConsumerValueLength) {
376-
LOG(FATAL) << "Invalid stream consumer meta value length: " << value_.size()
379+
LOG(ERROR) << "Invalid stream consumer meta value length: " << value_.size()
377380
<< " expected: " << kDefaultStreamConsumerValueLength;
378381
return;
379382
}
@@ -391,7 +394,7 @@ class StreamConsumerMetaValue {
391394
pel_ = pel;
392395
assert(value_.size() == 0);
393396
if (value_.size() != 0) {
394-
LOG(FATAL) << "Invalid stream consumer meta value length: " << value_.size() << " expected: 0";
397+
LOG(ERROR) << "Invalid stream consumer meta value length: " << value_.size() << " expected: 0";
395398
return;
396399
}
397400
uint64_t needed = kDefaultStreamConsumerValueLength;

src/storage/src/redis_strings.cc

Lines changed: 3 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1703,19 +1703,19 @@ rocksdb::Status Redis::PKPatternMatchDel(const std::string& pattern, int32_t* re
17031703
rocksdb::WriteBatch batch;
17041704
rocksdb::Iterator* iter = db_->NewIterator(iterator_options, handles_[kMetaCF]);
17051705
iter->SeekToFirst();
1706-
key = iter->key().ToString();
17071706
while (iter->Valid()) {
17081707
auto meta_type = static_cast<enum DataType>(static_cast<uint8_t>(iter->value()[0]));
17091708
ParsedBaseMetaKey parsed_meta_key(iter->key().ToString());
1709+
key = iter->key().ToString();
1710+
meta_value = iter->value().ToString();
1711+
17101712
if (meta_type == DataType::kStrings) {
1711-
meta_value = iter->value().ToString();
17121713
ParsedStringsValue parsed_strings_value(&meta_value);
17131714
if (!parsed_strings_value.IsStale() &&
17141715
(StringMatch(pattern.data(), pattern.size(), parsed_meta_key.Key().data(), parsed_meta_key.Key().size(), 0) != 0)) {
17151716
batch.Delete(key);
17161717
}
17171718
} else if (meta_type == DataType::kLists) {
1718-
meta_value = iter->value().ToString();
17191719
ParsedListsMetaValue parsed_lists_meta_value(&meta_value);
17201720
if (!parsed_lists_meta_value.IsStale() && (parsed_lists_meta_value.Count() != 0U) &&
17211721
(StringMatch(pattern.data(), pattern.size(), parsed_meta_key.Key().data(), parsed_meta_key.Key().size(), 0) !=
@@ -1732,7 +1732,6 @@ rocksdb::Status Redis::PKPatternMatchDel(const std::string& pattern, int32_t* re
17321732
batch.Put(handles_[kMetaCF], key, stream_meta_value.value());
17331733
}
17341734
} else {
1735-
meta_value = iter->value().ToString();
17361735
ParsedBaseMetaValue parsed_meta_value(&meta_value);
17371736
if (!parsed_meta_value.IsStale() && (parsed_meta_value.Count() != 0) &&
17381737
(StringMatch(pattern.data(), pattern.size(), parsed_meta_key.Key().data(), parsed_meta_key.Key().size(), 0) !=

src/storage/src/storage.cc

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1401,11 +1401,14 @@ Status Storage::PKRScanRange(const DataType& data_type, const Slice& key_start,
14011401

14021402
Status Storage::PKPatternMatchDel(const DataType& data_type, const std::string& pattern, int32_t* ret) {
14031403
Status s;
1404+
*ret = 0;
14041405
for (const auto& inst : insts_) {
1405-
s = inst->PKPatternMatchDel(pattern, ret);
1406+
int32_t tmp_ret = 0;
1407+
s = inst->PKPatternMatchDel(pattern, &tmp_ret);
14061408
if (!s.ok()) {
14071409
return s;
14081410
}
1411+
*ret += tmp_ret;
14091412
}
14101413
return s;
14111414
}

0 commit comments

Comments
 (0)