@@ -48,6 +48,8 @@ std::shared_ptr<rocksdb::Cache> pegasus_server_impl::_s_block_cache;
4848::dsn::task_ptr pegasus_server_impl::_update_server_rdb_stat;
4949::dsn::perf_counter_wrapper pegasus_server_impl::_pfc_rdb_block_cache_mem_usage;
5050const std::string pegasus_server_impl::COMPRESSION_HEADER = " per_level:" ;
51+ const std::string pegasus_server_impl::DATA_COLUMN_FAMILY_NAME = " default" ;
52+ const std::string pegasus_server_impl::META_COLUMN_FAMILY_NAME = " pegasus_meta_cf" ;
5153
5254pegasus_server_impl::pegasus_server_impl (dsn::replication::replica *r)
5355 : dsn::apps::rrdb_service(r),
@@ -452,8 +454,7 @@ pegasus_server_impl::~pegasus_server_impl()
452454{
453455 if (_is_open) {
454456 dassert (_db != nullptr , " " );
455- delete _db;
456- _db = nullptr ;
457+ release_db ();
457458 }
458459}
459460
@@ -1647,6 +1648,7 @@ ::dsn::error_code pegasus_server_impl::start(int argc, char **argv)
16471648 // 2, we can parse restore info from app env, which is stored in argv
16481649 // 3, restore_dir is exist
16491650 //
1651+ bool db_exist = true ;
16501652 auto path = ::dsn::utils::filesystem::path_combine (data_dir (), " rdb" );
16511653 if (::dsn::utils::filesystem::path_exists (path)) {
16521654 // only case 1
@@ -1662,6 +1664,7 @@ ::dsn::error_code pegasus_server_impl::start(int argc, char **argv)
16621664 replica_name ());
16631665 return ::dsn::ERR_FILE_OPERATION_FAILED ;
16641666 } else {
1667+ db_exist = false ;
16651668 dinfo (" %s: open a new db, path = %s" , replica_name (), path.c_str ());
16661669 }
16671670 } else {
@@ -1688,6 +1691,7 @@ ::dsn::error_code pegasus_server_impl::start(int argc, char **argv)
16881691 restore_dir.c_str ());
16891692 return ::dsn::ERR_FILE_OPERATION_FAILED ;
16901693 } else {
1694+ db_exist = false ;
16911695 dwarn (
16921696 " %s: try to restore and restore_dir(%s) isn't exist, but we don't force "
16931697 " it, the role of this replica must not primary, so we open a new db on the "
@@ -1702,16 +1706,36 @@ ::dsn::error_code pegasus_server_impl::start(int argc, char **argv)
17021706
17031707 ddebug (" %s: start to open rocksDB's rdb(%s)" , replica_name (), path.c_str ());
17041708
1705- auto status = rocksdb::DB::Open (rocksdb::Options (_db_opts, _data_cf_opts), path, &_db);
1709+ bool need_open_with_meta_cf = false ;
1710+ // Check meta CF only when db exist.
1711+ if (db_exist && check_meta_cf (path, &need_open_with_meta_cf) != ::dsn::ERR_OK ) {
1712+ derror_replica (" check meta column family failed" );
1713+ return ::dsn::ERR_LOCAL_APP_FAILURE ;
1714+ }
1715+ std::vector<rocksdb::ColumnFamilyDescriptor> column_families (
1716+ {{DATA_COLUMN_FAMILY_NAME , _data_cf_opts}});
1717+ if (need_open_with_meta_cf) {
1718+ column_families.emplace_back (rocksdb::ColumnFamilyDescriptor (
1719+ META_COLUMN_FAMILY_NAME , rocksdb::ColumnFamilyOptions ()));
1720+ }
1721+ std::vector<rocksdb::ColumnFamilyHandle *> handles_opened;
1722+ auto status = rocksdb::DB::Open (_db_opts, path, column_families, &handles_opened, &_db);
17061723 if (status.ok ()) {
1724+ dcheck_eq_replica (column_families.size (), handles_opened.size ());
1725+ dcheck_eq_replica (handles_opened[0 ]->GetName (), DATA_COLUMN_FAMILY_NAME );
1726+ _data_cf = handles_opened[0 ];
1727+ if (handles_opened.size () == 2 ) {
1728+ dcheck_eq_replica (handles_opened[1 ]->GetName (), META_COLUMN_FAMILY_NAME );
1729+ _meta_cf = handles_opened[1 ];
1730+ }
1731+
17071732 _last_committed_decree = _db->GetLastFlushedDecree ();
17081733 _pegasus_data_version = _db->GetPegasusDataVersion ();
17091734 if (_pegasus_data_version > PEGASUS_DATA_VERSION_MAX ) {
17101735 derror (" %s: open app failed, unsupported data version %" PRIu32,
17111736 replica_name (),
17121737 _pegasus_data_version);
1713- delete _db;
1714- _db = nullptr ;
1738+ release_db ();
17151739 return ::dsn::ERR_LOCAL_APP_FAILURE ;
17161740 }
17171741
@@ -1736,8 +1760,7 @@ ::dsn::error_code pegasus_server_impl::start(int argc, char **argv)
17361760 auto err = async_checkpoint (false );
17371761 if (err != ::dsn::ERR_OK ) {
17381762 derror (" %s: create checkpoint failed, error = %s" , replica_name (), err.to_string ());
1739- delete _db;
1740- _db = nullptr ;
1763+ release_db ();
17411764 return err;
17421765 }
17431766 dassert (last_flushed == last_durable_decree (),
@@ -1822,8 +1845,7 @@ ::dsn::error_code pegasus_server_impl::stop(bool clear_state)
18221845 _context_cache.clear ();
18231846
18241847 _is_open = false ;
1825- delete _db;
1826- _db = nullptr ;
1848+ release_db ();
18271849
18281850 std::deque<int64_t > reserved_checkpoints;
18291851 {
@@ -2894,5 +2916,41 @@ void pegasus_server_impl::set_partition_version(int32_t partition_version)
28942916 // TODO(heyuchen): set filter _partition_version in further pr
28952917}
28962918
2919+ ::dsn::error_code pegasus_server_impl::check_meta_cf (const std::string &path,
2920+ bool *need_open_with_meta_cf)
2921+ {
2922+ *need_open_with_meta_cf = false ;
2923+ std::vector<std::string> column_families;
2924+ auto s = rocksdb::DB::ListColumnFamilies (rocksdb::DBOptions (), path, &column_families);
2925+ if (!s.ok ()) {
2926+ derror_replica (" rocksdb::DB::ListColumnFamilies failed, error = {}" , s.ToString ());
2927+ return ::dsn::ERR_LOCAL_APP_FAILURE ;
2928+ }
2929+
2930+ for (const auto &column_family : column_families) {
2931+ if (column_family == DATA_COLUMN_FAMILY_NAME ) {
2932+ continue ;
2933+ }
2934+ if (column_family == META_COLUMN_FAMILY_NAME ) {
2935+ *need_open_with_meta_cf = true ;
2936+ continue ;
2937+ }
2938+ dassert_replica (false , " Column family '{}' should not present" );
2939+ }
2940+ return ::dsn::ERR_OK ;
2941+ }
2942+
2943+ void pegasus_server_impl::release_db ()
2944+ {
2945+ _db->DestroyColumnFamilyHandle (_data_cf);
2946+ _data_cf = nullptr ;
2947+ if (_meta_cf != nullptr ) {
2948+ _db->DestroyColumnFamilyHandle (_meta_cf);
2949+ }
2950+ _meta_cf = nullptr ;
2951+ delete _db;
2952+ _db = nullptr ;
2953+ }
2954+
28972955} // namespace server
28982956} // namespace pegasus
0 commit comments