Skip to content
Merged
Show file tree
Hide file tree
Changes from 4 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 0 additions & 9 deletions tiledb/api/c_api/array/array_api_internal.h
Original file line number Diff line number Diff line change
Expand Up @@ -98,15 +98,6 @@ struct tiledb_array_handle_t
array_->delete_array(uri);
}

void delete_fragments(
tiledb::sm::ContextResources& resources,
const tiledb::sm::URI& uri,
uint64_t ts_start,
uint64_t ts_end,
std::optional<tiledb::sm::ArrayDirectory> array_dir = std::nullopt) {
array_->delete_fragments(resources, uri, ts_start, ts_end, array_dir);
}

void delete_fragments(
const tiledb::sm::URI& uri,
uint64_t timestamp_start,
Expand Down
34 changes: 13 additions & 21 deletions tiledb/sm/array/array.cc
Original file line number Diff line number Diff line change
Expand Up @@ -658,42 +658,33 @@ Status Array::close() {
}

void Array::delete_fragments(
ContextResources& resources,
const URI& uri,
uint64_t timestamp_start,
uint64_t timestamp_end,
std::optional<ArrayDirectory> array_dir) {
// Get the fragment URIs to be deleted
if (array_dir == std::nullopt) {
array_dir = ArrayDirectory(resources, uri, timestamp_start, timestamp_end);
}
auto filtered_fragment_uris = array_dir->filtered_fragment_uris(true);
ContextResources& resources, ArrayDirectory& array_dir) {
auto filtered_fragment_uris = array_dir.filtered_fragment_uris(true);
const auto& fragment_uris = filtered_fragment_uris.fragment_uris();

// Retrieve commit uris to delete and ignore
std::vector<URI> commit_uris_to_delete;
std::vector<URI> commit_uris_to_ignore;
for (auto& fragment : fragment_uris) {
auto commit_uri = array_dir->get_commit_uri(fragment.uri_);
auto commit_uri = array_dir.get_commit_uri(fragment.uri_);
commit_uris_to_delete.emplace_back(commit_uri);
if (array_dir->consolidated_commit_uris_set().count(commit_uri.c_str()) !=
0) {
if (array_dir.consolidated_commit_uris_set().count(commit_uri) != 0) {
commit_uris_to_ignore.emplace_back(commit_uri);
}
}

// Write ignore file
if (commit_uris_to_ignore.size() != 0) {
array_dir->write_commit_ignore_file(commit_uris_to_ignore);
array_dir.write_commit_ignore_file(commit_uris_to_ignore);
}

// Delete fragments and commits
auto vfs = &(resources.vfs());
auto& vfs = resources.vfs();
throw_if_not_ok(parallel_for(
&resources.compute_tp(), 0, fragment_uris.size(), [&](size_t i) {
vfs->remove_dir(fragment_uris[i].uri_);
if (vfs->is_file(commit_uris_to_delete[i])) {
vfs->remove_file(commit_uris_to_delete[i]);
vfs.remove_dir(fragment_uris[i].uri_);
if (vfs.is_file(commit_uris_to_delete[i])) {
vfs.remove_file(commit_uris_to_delete[i]);
}
return Status::Ok();
}));
Expand All @@ -714,7 +705,9 @@ void Array::delete_fragments(
rest_client->post_delete_fragments_to_rest(
uri, this, timestamp_start, timestamp_end);
} else {
Array::delete_fragments(resources_, uri, timestamp_start, timestamp_end);
auto array_dir =
ArrayDirectory(resources_, uri, timestamp_start, timestamp_end);
Array::delete_fragments(resources_, array_dir);
}
}

Expand All @@ -724,8 +717,7 @@ void Array::delete_array(ContextResources& resources, const URI& uri) {
ArrayDirectory(resources, uri, 0, std::numeric_limits<uint64_t>::max());

// Delete fragments and commits
Array::delete_fragments(
resources, uri, 0, std::numeric_limits<uint64_t>::max(), array_dir);
Array::delete_fragments(resources, array_dir);

// Delete array metadata, fragment metadata and array schema files
// Note: metadata files may not be present, try to delete anyway
Expand Down
23 changes: 5 additions & 18 deletions tiledb/sm/array/array.h
Original file line number Diff line number Diff line change
Expand Up @@ -100,7 +100,7 @@ class OpenedArray {
uint64_t timestamp_end_opened_at,
bool is_remote)
: resources_(resources)
, array_dir_(ArrayDirectory(resources, array_uri))
, array_dir_(std::in_place, resources, array_uri)
, array_schema_latest_(nullptr)
, metadata_(memory_tracker)
, metadata_loaded_(false)
Expand All @@ -124,8 +124,8 @@ class OpenedArray {
}

/** Sets the array directory. */
inline void set_array_directory(const ArrayDirectory&& dir) {
array_dir_ = dir;
inline void set_array_directory(ArrayDirectory&& dir) {
array_dir_.emplace(std::move(dir));
}

/** Returns the latest array schema. */
Expand Down Expand Up @@ -320,7 +320,7 @@ class Array {
}

/** Set the array directory. */
inline void set_array_directory(const ArrayDirectory&& dir) {
inline void set_array_directory(ArrayDirectory&& dir) {
opened_array_->set_array_directory(std::move(dir));
}

Expand Down Expand Up @@ -463,23 +463,10 @@ class Array {
* between the provided timestamps.
*
* @param resources The context resources.
* @param uri The uri of the Array whose fragments are to be deleted.
* @param timestamp_start The start timestamp at which to delete fragments.
* @param timestamp_end The end timestamp at which to delete fragments.
* @param array_dir An optional ArrayDirectory from which to delete fragments.
*
* @section Maturity Notes
* This is legacy code, ported from StorageManager during its removal process.
* Its existence supports the non-static `delete_fragments` API below,
* performing the actual deletion of fragments. This function is slated for
* removal and should be directly integrated into the function below.
*/
static void delete_fragments(
ContextResources& resources,
const URI& uri,
uint64_t timestamp_start,
uint64_t timstamp_end,
std::optional<ArrayDirectory> array_dir = std::nullopt);
ContextResources& resources, ArrayDirectory& array_dir);

/**
* Handles local and remote deletion of fragments between the provided
Expand Down
35 changes: 15 additions & 20 deletions tiledb/sm/array/array_directory.cc
Original file line number Diff line number Diff line change
Expand Up @@ -253,8 +253,8 @@ const std::vector<URI>& ArrayDirectory::commit_uris_to_vacuum() const {
return commit_uris_to_vacuum_;
}

const std::unordered_set<std::string>&
ArrayDirectory::consolidated_commit_uris_set() const {
const std::unordered_set<URI>& ArrayDirectory::consolidated_commit_uris_set()
const {
return consolidated_commit_uris_set_;
}

Expand Down Expand Up @@ -309,7 +309,7 @@ void ArrayDirectory::delete_fragments_list(
for (auto& timestamped_uri : uris) {
auto commit_uri = get_commit_uri(timestamped_uri);
commit_uris_to_delete.emplace_back(commit_uri);
if (consolidated_commit_uris_set().count(commit_uri.c_str()) != 0) {
if (consolidated_commit_uris_set().count(commit_uri) != 0) {
commit_uris_to_ignore.emplace_back(commit_uri);
}
}
Expand Down Expand Up @@ -483,7 +483,7 @@ ArrayDirectory::filtered_fragment_uris(const bool full_overlap_only) const {
if (mode_ == ArrayDirectoryMode::VACUUM_FRAGMENTS) {
for (auto& uri : fragment_uris_to_vacuum.value()) {
auto commit_uri = get_commit_uri(uri);
if (consolidated_commit_uris_set_.count(commit_uri.c_str()) == 0) {
if (consolidated_commit_uris_set_.count(commit_uri) == 0) {
commit_uris_to_vacuum.emplace_back(commit_uri);
} else {
commit_uris_to_ignore.emplace_back(commit_uri);
Expand Down Expand Up @@ -672,8 +672,7 @@ ArrayDirectory::load_commits_dir_uris_v12_or_higher(
for (size_t i = 0; i < commits_dir_uris.size(); ++i) {
if (stdx::string::ends_with(
commits_dir_uris[i].to_string(), constants::write_file_suffix)) {
if (consolidated_commit_uris_set_.count(commits_dir_uris[i].c_str()) ==
0) {
if (consolidated_commit_uris_set_.count(commits_dir_uris[i]) == 0) {
auto name = commits_dir_uris[i].last_path_part();
name =
name.substr(0, name.size() - constants::write_file_suffix.size());
Expand All @@ -694,8 +693,7 @@ ArrayDirectory::load_commits_dir_uris_v12_or_higher(

// Add the delete tile location if it overlaps the open start/end times
if (timestamps_overlap(timestamp_range, false)) {
if (consolidated_commit_uris_set_.count(commits_dir_uris[i].c_str()) ==
0) {
if (consolidated_commit_uris_set_.count(commits_dir_uris[i]) == 0) {
const auto base_uri_size = uri_.to_string().size();
delete_and_update_tiles_location_.emplace_back(
commits_dir_uris[i],
Expand All @@ -716,10 +714,7 @@ ArrayDirectory::list_fragment_metadata_dir_uris_v12_or_higher() {
return ls(uri_.join_path(constants::array_fragment_meta_dir_name));
}

tuple<
Status,
optional<std::vector<URI>>,
optional<std::unordered_set<std::string>>>
tuple<Status, optional<std::vector<URI>>, optional<std::unordered_set<URI>>>
ArrayDirectory::load_consolidated_commit_uris(
const std::vector<URI>& commits_dir_uris) {
auto timer_se = stats_->start_timer("load_consolidated_commit_uris");
Expand All @@ -745,7 +740,7 @@ ArrayDirectory::load_consolidated_commit_uris(

// Load all commit URIs. This is done in serial for now as it can be optimized
// by vacuuming.
std::unordered_set<std::string> uris_set;
std::unordered_set<URI> uris_set;
std::vector<std::pair<URI, std::string>> meta_files;
for (uint64_t i = 0; i < commits_dir_uris.size(); i++) {
auto& uri = commits_dir_uris[i];
Expand All @@ -763,7 +758,7 @@ ArrayDirectory::load_consolidated_commit_uris(
std::stringstream ss(names);
for (std::string condition_marker; std::getline(ss, condition_marker);) {
if (ignore_set.count(condition_marker) == 0) {
uris_set.emplace(uri_.to_string() + condition_marker);
uris_set.emplace(uri_.append_string(condition_marker));
}

// If we have a delete, process the condition tile
Expand Down Expand Up @@ -807,7 +802,7 @@ ArrayDirectory::load_consolidated_commit_uris(
uint64_t count = 0;
bool all_in_set = true;
for (std::string uri_str; std::getline(ss, uri_str);) {
if (uris_set.count(uri_.to_string() + uri_str) > 0) {
if (uris_set.count(uri_.append_string(uri_str)) > 0) {
count++;
} else {
all_in_set = false;
Expand Down Expand Up @@ -879,7 +874,7 @@ void ArrayDirectory::load_commits_uris_to_consolidate(
const std::vector<URI>& array_dir_uris,
const std::vector<URI>& commits_dir_uris,
const std::vector<URI>& consolidated_uris,
const std::unordered_set<std::string>& consolidated_uris_set) {
const std::unordered_set<URI>& consolidated_uris_set) {
// Make a set of existing commit URIs.
std::unordered_set<std::string> uris_set;
for (auto& uri : array_dir_uris) {
Expand All @@ -902,7 +897,7 @@ void ArrayDirectory::load_commits_uris_to_consolidate(
// Add the ok file URIs not already in the list.
for (auto& uri : array_dir_uris) {
if (stdx::string::ends_with(uri.to_string(), constants::ok_file_suffix)) {
if (consolidated_uris_set.count(uri.c_str()) == 0) {
if (consolidated_uris_set.count(uri) == 0) {
commit_uris_to_consolidate_.emplace_back(uri);
}
}
Expand All @@ -914,7 +909,7 @@ void ArrayDirectory::load_commits_uris_to_consolidate(
uri.to_string(), constants::write_file_suffix) ||
stdx::string::ends_with(
uri.to_string(), constants::delete_file_suffix)) {
if (consolidated_uris_set.count(uri.c_str()) == 0) {
if (consolidated_uris_set.count(uri) == 0) {
commit_uris_to_consolidate_.emplace_back(uri);
}
}
Expand Down Expand Up @@ -1258,7 +1253,7 @@ bool ArrayDirectory::is_vacuum_file(const URI& uri) const {
Status ArrayDirectory::is_fragment(
const URI& uri,
const std::unordered_set<std::string>& ok_uris_set,
const std::unordered_set<std::string>& consolidated_uris_set,
const std::unordered_set<URI>& consolidated_uris_set,
int* is_fragment) const {
// If the fragment ID does not have a name, ignore it.
if (!FragmentID::has_fragment_name(uri)) {
Expand Down Expand Up @@ -1292,7 +1287,7 @@ Status ArrayDirectory::is_fragment(

// Check set membership in consolidated uris
if (consolidated_uris_set.count(
uri.to_string() + constants::ok_file_suffix) != 0) {
uri.append_string(constants::ok_file_suffix)) != 0) {
*is_fragment = 1;
return Status::Ok();
}
Expand Down
19 changes: 10 additions & 9 deletions tiledb/sm/array/array_directory.h
Original file line number Diff line number Diff line change
Expand Up @@ -311,6 +311,10 @@ class ArrayDirectory {
uint64_t timestamp_end,
ArrayDirectoryMode mode = ArrayDirectoryMode::READ);

DISABLE_COPY_AND_COPY_ASSIGN(ArrayDirectory);
Comment thread
teo-tsirpanis marked this conversation as resolved.

ArrayDirectory(ArrayDirectory&&) = default;

/** Destructor. */
~ArrayDirectory() = default;

Expand Down Expand Up @@ -436,7 +440,7 @@ class ArrayDirectory {
const std::vector<URI>& commit_uris_to_vacuum() const;

/** Returns the consolidated commit URI set. */
const std::unordered_set<std::string>& consolidated_commit_uris_set() const;
const std::unordered_set<URI>& consolidated_commit_uris_set() const;

/** Returns the URIs of the consolidated commit files to vacuum. */
const std::vector<URI>& consolidated_commits_uris_to_vacuum() const;
Expand Down Expand Up @@ -511,7 +515,7 @@ class ArrayDirectory {
}

/** Accessor to consolidated_commit_uris_set_ */
inline std::unordered_set<std::string>& consolidated_commit_uris_set() {
inline std::unordered_set<URI>& consolidated_commit_uris_set() {
return consolidated_commit_uris_set_;
}

Expand Down Expand Up @@ -589,7 +593,7 @@ class ArrayDirectory {
std::vector<URI> unfiltered_fragment_uris_;

/** Consolidated commit URI set. */
std::unordered_set<std::string> consolidated_commit_uris_set_;
std::unordered_set<URI> consolidated_commit_uris_set_;

/** The URIs of all the array schema files. */
std::vector<URI> array_schema_uris_;
Expand Down Expand Up @@ -710,18 +714,15 @@ class ArrayDirectory {
const std::vector<URI>& array_dir_uris,
const std::vector<URI>& commits_dir_uris,
const std::vector<URI>& consolidated_uris,
const std::unordered_set<std::string>& consolidated_uris_set);
const std::unordered_set<URI>& consolidated_uris_set);

/**
* Loads the consolidated commit URI from the commit directory and the files
* to vacuum for the commits directory.
*
* @return Status, consolidated uris, set of all consolidated uris.
*/
tuple<
Status,
optional<std::vector<URI>>,
optional<std::unordered_set<std::string>>>
tuple<Status, optional<std::vector<URI>>, optional<std::unordered_set<URI>>>
load_consolidated_commit_uris(const std::vector<URI>& commits_dir_uris);

/** Loads the array metadata URIs. */
Expand Down Expand Up @@ -821,7 +822,7 @@ class ArrayDirectory {
Status is_fragment(
const URI& uri,
const std::unordered_set<std::string>& ok_uris_set,
const std::unordered_set<std::string>& consolidated_uris_set,
const std::unordered_set<URI>& consolidated_uris_set,
int32_t* is_fragment) const;

/**
Expand Down
2 changes: 1 addition & 1 deletion tiledb/sm/consolidator/consolidator.cc
Original file line number Diff line number Diff line change
Expand Up @@ -252,7 +252,7 @@ void Consolidator::fragments_consolidate(

void Consolidator::write_consolidated_commits_file(
format_version_t write_version,
ArrayDirectory array_dir,
const ArrayDirectory& array_dir,
const std::vector<URI>& commit_uris,
ContextResources& resources) {
// Compute the file name.
Expand Down
2 changes: 1 addition & 1 deletion tiledb/sm/consolidator/consolidator.h
Original file line number Diff line number Diff line change
Expand Up @@ -215,7 +215,7 @@ class Consolidator {
*/
static void write_consolidated_commits_file(
format_version_t write_version,
ArrayDirectory array_dir,
const ArrayDirectory& array_dir,
const std::vector<URI>& commit_uris,
ContextResources& resources);

Expand Down
Loading
Loading