概述
在RocksDB 3.0中加入了Column Family特性,加入这个特性之后,每一个KV对都会关联一个Column Family,其中默认的Column Family是 “default”. Column Family主要是提供给RocksDB一个逻辑的分区.从实现上来看不同的Column Family共享WAL,而都有自己的Memtable和SST.这就意味着我们可以很 快速已经方便的设置不同的属性给不同的Column Family以及快速删除对应的Column Family.
主要API
首先是创建Column Family,这里注意我们可以通过两种方式来创建Column Family,一种是在Open DB的时候通过传递需要创建的Column Family,一种是当DB创建并打开之后, 通过直接的CreateColumnFamily来创建Column Family.
1. DB::Open(const DBOptions& db_options, const std::string& name, const std::vector<ColumnFamilyDescriptor>& column_families, std::vector<ColumnFamilyHandle*>* handles, DB** dbptr);
2. DB::CreateColumnFamily(const ColumnFamilyOptions& options, const std::string& column_family_name, ColumnFamilyHandle** handle);
这里可以看到不管是哪一种方式最终都会返回一个ColumnFamilyHandle给调用者来使用.
然后就是删除Column Family的方式,这里很简单就是传递之前创建的ColumnFamilyHandle给RocksDB,然后用以删除.
1. DropColumnFamily(ColumnFamilyHandle* column_family);
实现
所有的Column Family都是通过一个叫做ColumnFamilySet的结构来管理的,而每一个Column Family都是一个ColumnFamilyData.
先来看ColumnFamilySet,这里可以看到它有两个数据结构来管理Column Family,分别是map(column_family_data_)以及一个双向链表(dummy_cfd_). 其中map用来保存Column Family名字和对应的id以及ColumnFamilyData的映射. 这里要注意在RocksDB内部是将没一个ColumnFamily的名字表示为一个uint32类型的ID(max_column_family_).也就是这个ID是一个简单的递增的数值.
1. class ColumnFamilySet {
2. public:
3. // ColumnFamilySet supports iteration
4. public:
5. .................................
7. ColumnFamilyData* CreateColumnFamily(const std::string& name, uint32_t id,
8. Version* dummy_version,
9. const ColumnFamilyOptions& options);
10. iterator begin() { return iterator(dummy_cfd_->next_); }
11. iterator end() { return iterator(dummy_cfd_); }
12. ...............................
13. private:
14. friend class ColumnFamilyData;
15. // helper function that gets called from cfd destructor
16. // REQUIRES: DB mutex held
17. void RemoveColumnFamily(ColumnFamilyData* cfd);
19. // column_families_ and column_family_data_ need to be protected:
20. // * when mutating both conditions have to be satisfied:
21. // 1. DB mutex locked
22. // 2. thread currently in single-threaded write thread
23. // * when reading, at least one condition needs to be satisfied:
24. // 1. DB mutex locked
25. // 2. accessed from a single-threaded write thread
26. std::unordered_map<std::string, uint32_t> column_families_;
27. std::unordered_map<uint32_t, ColumnFamilyData*> column_family_data_;
29. uint32_t max_column_family_;
30. ColumnFamilyData* dummy_cfd_;
31. // We don't hold the refcount here, since default column family always exists
32. // We are also not responsible for cleaning up default_cfd_cache_. This is
33. // just a cache that makes common case (accessing default column family)
34. // faster
35. ColumnFamilyData* default_cfd_cache_;
37. ..................................
38. };
然后来看ColumnFamilyData,这个数据结构就是用来表示一个ColumnFamily,保存了对应的信息,我们可以看到有ID/name以及当前ColumnFamily对应的所有的version(dummy_versions_). 其中这里的next_/prev_就是在ColumnFamilySet中用来表示所有ColumnFamily的双向链表.
1. class ColumnFamilyData {
2. public:
3. ~ColumnFamilyData();
5. // thread-safe
6. uint32_t GetID() const { return id_; }
7. // thread-safe
8. const std::string& GetName() const { return name_; }
10. // Ref() can only be called from a context where the caller can guarantee
11. // that ColumnFamilyData is alive (while holding a non-zero ref already,
12. // holding a DB mutex, or as the leader in a write batch group).
13. void Ref() { refs_.fetch_add(1, std::memory_order_relaxed); }
15. // Unref decreases the reference count, but does not handle deletion
16. // when the count goes to 0. If this method returns true then the
17. // caller should delete the instance immediately, or later, by calling
18. // FreeDeadColumnFamilies(). Unref() can only be called while holding
19. // a DB mutex, or during single-threaded recovery.
20. bool Unref() {
21. int old_refs = refs_.fetch_sub(1, std::memory_order_relaxed);
22. assert(old_refs > 0);
23. return old_refs == 1;
24. }
25. ..............................
27. private:
28. friend class ColumnFamilySet;
29. ColumnFamilyData(uint32_t id, const std::string& name,
30. Version* dummy_versions, Cache* table_cache,
31. WriteBufferManager* write_buffer_manager,
32. const ColumnFamilyOptions& options,
33. const ImmutableDBOptions& db_options,
34. const EnvOptions& env_options,
35. ColumnFamilySet* column_family_set);
37. uint32_t id_;
38. const std::string name_;
39. Version* dummy_versions_; // Head of circular doubly-linked list of versions.
40. Version* current_; // == dummy_versions->prev_
41. ......................................................
43. // Thread's local copy of SuperVersion pointer
44. // This needs to be destructed before mutex_
45. std::unique_ptr<ThreadLocalPtr> local_sv_;
47. // pointers for a circular linked list. we use it to support iterations over
48. // all column families that are alive (note: dropped column families can also
49. // be alive as long as client holds a reference)
50. ColumnFamilyData* next_;
51. ColumnFamilyData* prev_;
52. ...................................
54. ColumnFamilySet* column_family_set_;
55. ..................................
56. };
然后就是返回给调用者的ColumnFamilyHandleImpl结构,这个结构主要是封装了ColumnFamilyData.
1. // ColumnFamilyHandleImpl is the class that clients use to access different
2. // column families. It has non-trivial destructor, which gets called when client
3. // is done using the column family
4. class ColumnFamilyHandleImpl : public ColumnFamilyHandle {
5. public:
6. // create while holding the mutex
7. ColumnFamilyHandleImpl(
8. ColumnFamilyData* cfd, DBImpl* db, InstrumentedMutex* mutex);
9. // destroy without mutex
10. virtual ~ColumnFamilyHandleImpl();
11. virtual ColumnFamilyData* cfd() const { return cfd_; }
12. ......................................
14. private:
15. ColumnFamilyData* cfd_;
16. DBImpl* db_;
17. InstrumentedMutex* mutex_;
18. };
接下来我们就来从ColumnFamily的创建以及删除来分析ColumnFamily的实现.我们从DBImpl::CreateColumnFamilyImpl开始.在这个函数 中首先就是通过调用GetNextColumnFamilyID来得到当前创建的ColumnFamily对应的ID(自增).然后再调用LogAndApply来对ColumnFamily 进行对应的操作.最后再返回封装好的ColumnFamilyHandle给调用者.
1. Status DBImpl::CreateColumnFamilyImpl(const ColumnFamilyOptions& cf_options,
2. const std::string& column_family_name,
3. ColumnFamilyHandle** handle) {
4. .......................................
6. {
7. ...................................
8. VersionEdit edit;
9. edit.AddColumnFamily(column_family_name);
10. uint32_t new_id = versions_->GetColumnFamilySet()->GetNextColumnFamilyID();
11. edit.SetColumnFamily(new_id);
12. edit.SetLogNumber(logfile_number_);
13. edit.SetComparatorName(cf_options.comparator->Name());
15. // LogAndApply will both write the creation in MANIFEST and create
16. // ColumnFamilyData object
17. { // write thread
18. WriteThread::Writer w;
19. write_thread_.EnterUnbatched(&w, &mutex_);
20. // LogAndApply will both write the creation in MANIFEST and create
21. // ColumnFamilyData object
22. s = versions_->LogAndApply(nullptr, MutableCFOptions(cf_options), &edit,
23. &mutex_, directories_.GetDbDir(), false,
24. &cf_options);
25. write_thread_.ExitUnbatched(&w);
26. }
27. if (s.ok()) {
28. ........................................
29. *handle = new ColumnFamilyHandleImpl(cfd, this, &mutex_);
30. ROCKS_LOG_INFO(immutable_db_options_.info_log,
31. "Created column family [%s] (ID %u)",
32. column_family_name.c_str(), (unsigned)cfd->GetID());
33. }
34. .............................................
35. } // InstrumentedMutexLock l(&mutex_)
37. .................................
38. return s;
39. }
最终会在LogAndApply调用ColumnFamilySet的CreateColumnFamily函数(通过VersionSet::CreateColumnFamily),这个函数我们可看到主要做了下面三件事情
- 创建ColumnFamilyData对象
- 将新的创建好的CFD加入到双向链表
- 对应的Map数据结构更新数据
```
- // under a DB mutex AND write thread
- ColumnFamilyData* ColumnFamilySet::CreateColumnFamily(
- const std::string& name, uint32_t id, Version* dummy_versions,
- const ColumnFamilyOptions& options) {
- assert(column_families_.find(name) == column_families_.end());
- ColumnFamilyData* new_cfd = new ColumnFamilyData(
- id, name, dummy_versions, table_cache_, write_buffer_manager_, options,
- *db_options_, env_options_, this);
- column_families_.insert({name, id});
- column_family_data_.insert({id, new_cfd});
- max_column_family_ = std::max(max_column_family_, id);
- // add to linked list
- new_cfd->next_ = dummy_cfd_;
- auto prev = dummy_cfd_->prev_;
- new_cfd->prev_ = prev;
- prev->next_ = new_cfd;
- dummy_cfd_->prev_ = new_cfd;
- if (id == 0) {
- default_cfd_cache_ = new_cfd;
- }
- return new_cfd;
- } ```
然后来看如何删除ColumnFamily,这里所有的删除最终都会调用ColumnFamilySet::RemoveColumnFamily函数,这个函数是是从两个Map中删除对应的ColumnFamily. 这里或许我们要问了,为什么管理的双向链表不需要删除呢。这里原因是这样的,由于ColumnFamilyData是通过引用计数管理的,因此只有当所有的引用计数都清零之后, 才需要真正的函数ColumnFamilyData(也就是会从双向链表中删除数据).
```
- // under a DB mutex AND from a write thread
- void ColumnFamilySet::RemoveColumnFamily(ColumnFamilyData* cfd) {
- auto cfd_iter = column_family_data_.find(cfd->GetID());
- assert(cfd_iter != column_family_data_.end());
- column_family_data_.erase(cfd_iter);
- column_families_.erase(cfd->GetName());
- } ```
因此我们来看ColumnFamilyData的析构函数.可以看到析构函数中会从双向链表中删除对应的数据,以及处理对应的Version(corrent_).
1. // DB mutex held
2. ColumnFamilyData::~ColumnFamilyData() {
3. assert(refs_.load(std::memory_order_relaxed) == 0);
4. // remove from linked list
5. auto prev = prev_;
6. auto next = next_;
7. prev->next_ = next;
8. next->prev_ = prev;
10. if (!dropped_ && column_family_set_ != nullptr) {
11. // If it's dropped, it's already removed from column family set
12. // If column_family_set_ == nullptr, this is dummy CFD and not in
13. // ColumnFamilySet
14. column_family_set_->RemoveColumnFamily(this);
15. }
17. if (current_ != nullptr) {
18. current_->Unref();
19. }
20. ..............................
21. }
最后我们来看一下在磁盘上ColumnFamily是如何保存的,首先需要明确的是ColumnFamily是保存在MANIFEST文件中的,信息的保存比较简单(之前的文章有介绍), 和MANIFEST中其他的信息没什么区别,因此这里我们主要来看数据的读取以及初始化,这里所有的操作都是包含在VersionSet::Recover中,我们来看这个函数.
函数主要的逻辑就是读取MANIFEST然后来再来将磁盘上读取的ColumnFamily的信息初始化(初始化ColumnFamilySet结构),可以看到这里相当于将之前的create/drop 的操作全部回放一遍,也就是会调用CreateColumnFamily/DropColumnFamily来将磁盘的信息初始化到内存.
1. while (reader.ReadRecord(&record, &scratch) && s.ok()) {
2. VersionEdit edit;
3. s = edit.DecodeFrom(record);
4. if (!s.ok()) {
5. break;
6. }
8. // Not found means that user didn't supply that column
9. // family option AND we encountered column family add
10. // record. Once we encounter column family drop record,
11. // we will delete the column family from
12. // column_families_not_found.
13. bool cf_in_not_found =
14. column_families_not_found.find(edit.column_family_) !=
15. column_families_not_found.end();
16. // in builders means that user supplied that column family
17. // option AND that we encountered column family add record
18. bool cf_in_builders =
19. builders.find(edit.column_family_) != builders.end();
21. // they can't both be true
22. assert(!(cf_in_not_found && cf_in_builders));
24. ColumnFamilyData* cfd = nullptr;
26. if (edit.is_column_family_add_) {
27. if (cf_in_builders || cf_in_not_found) {
28. s = Status::Corruption(
29. "Manifest adding the same column family twice");
30. break;
31. }
32. auto cf_options = cf_name_to_options.find(edit.column_family_name_);
33. if (cf_options == cf_name_to_options.end()) {
34. column_families_not_found.insert(
35. {edit.column_family_, edit.column_family_name_});
36. } else {
37. cfd = CreateColumnFamily(cf_options->second, &edit);
38. cfd->set_initialized();
39. builders.insert(
40. {edit.column_family_, new BaseReferencedVersionBuilder(cfd)});
41. }
42. } else if (edit.is_column_family_drop_) {
43. if (cf_in_builders) {
44. auto builder = builders.find(edit.column_family_);
45. assert(builder != builders.end());
46. delete builder->second;
47. builders.erase(builder);
48. cfd = column_family_set_->GetColumnFamily(edit.column_family_);
49. if (cfd->Unref()) {
50. delete cfd;
51. cfd = nullptr;
52. } else {
53. // who else can have reference to cfd!?
54. assert(false);
55. }
56. } else if (cf_in_not_found) {
57. column_families_not_found.erase(edit.column_family_);
58. } else {
59. s = Status::Corruption(
60. "Manifest - dropping non-existing column family");
61. break;
62. }
63. } else if (!cf_in_not_found) {
64. if (!cf_in_builders) {
65. s = Status::Corruption(
66. "Manifest record referencing unknown column family");
67. break;
68. }
70. cfd = column_family_set_->GetColumnFamily(edit.column_family_);
71. // this should never happen since cf_in_builders is true
72. assert(cfd != nullptr);
74. // if it is not column family add or column family drop,
75. // then it's a file add/delete, which should be forwarded
76. // to builder
77. auto builder = builders.find(edit.column_family_);
78. assert(builder != builders.end());
79. builder->second->version_builder()->Apply(&edit);
80. }
82. if (cfd != nullptr) {
83. if (edit.has_log_number_) {
84. if (cfd->GetLogNumber() > edit.log_number_) {
85. ROCKS_LOG_WARN(
86. db_options_->info_log,
87. "MANIFEST corruption detected, but ignored - Log numbers in "
88. "records NOT monotonically increasing");
89. } else {
90. cfd->SetLogNumber(edit.log_number_);
91. have_log_number = true;
92. }
93. }
94. if (edit.has_comparator_ &&
95. edit.comparator_ != cfd->user_comparator()->Name()) {
96. s = Status::InvalidArgument(
97. cfd->user_comparator()->Name(),
98. "does not match existing comparator " + edit.comparator_);
99. break;
100. }
101. }
103. if (edit.has_prev_log_number_) {
104. previous_log_number = edit.prev_log_number_;
105. have_prev_log_number = true;
106. }
108. if (edit.has_next_file_number_) {
109. next_file = edit.next_file_number_;
110. have_next_file = true;
111. }
113. if (edit.has_max_column_family_) {
114. max_column_family = edit.max_column_family_;
115. }
117. if (edit.has_last_sequence_) {
118. last_sequence = edit.last_sequence_;
119. have_last_sequence = true;
120. }
121. }
