Skip to content
Open
Show file tree
Hide file tree
Changes from all 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
20 changes: 14 additions & 6 deletions src/iceberg/manifest/manifest_group.cc
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,16 @@ namespace iceberg {

namespace {

// Selected with if/else rather than a conditional operator because older MSVC versions
// (e.g. 14.42) try to copy the move-only Result operands of the conditional operator.
Result<ManifestEntryStreamPtr> OpenEntriesStream(ManifestReader& reader,
bool ignore_deleted) {
if (ignore_deleted) {
return reader.LiveEntriesStream();
}
return reader.EntriesStream();
}

std::shared_ptr<Schema> DataFileFilterSchema() {
auto empty_partition_type = std::make_shared<StructType>(std::vector<SchemaField>{});
return std::make_shared<Schema>(std::vector<SchemaField>{
Expand Down Expand Up @@ -341,9 +351,8 @@ class ManifestGroup::FilePlanningStream final : public FileScanTaskStream {
}

ICEBERG_ASSIGN_OR_RAISE(auto reader, group_->MakeReader(manifest, columns_));
ICEBERG_ASSIGN_OR_RAISE(entry_stream_, group_->ignore_deleted_
? reader->LiveEntriesStream()
: reader->EntriesStream());
ICEBERG_ASSIGN_OR_RAISE(entry_stream_,
OpenEntriesStream(*reader, group_->ignore_deleted_));
current_spec_id_ = manifest.partition_spec_id;
return true;
}
Expand Down Expand Up @@ -376,9 +385,8 @@ class ManifestGroup::FilePlanningStream final : public FileScanTaskStream {
[this](const ManifestFile* manifest) -> Result<std::vector<TaggedStream>> {
ICEBERG_ASSIGN_OR_RAISE(auto reader,
group_->MakeReader(*manifest, columns_));
ICEBERG_ASSIGN_OR_RAISE(auto stream, group_->ignore_deleted_
? reader->LiveEntriesStream()
: reader->EntriesStream());
ICEBERG_ASSIGN_OR_RAISE(
auto stream, OpenEntriesStream(*reader, group_->ignore_deleted_));

std::vector<TaggedStream> tagged_streams;
tagged_streams.emplace_back(manifest->partition_spec_id, std::move(stream));
Expand Down
2 changes: 1 addition & 1 deletion src/iceberg/test/cache_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -95,7 +95,7 @@ TEST(MemoizeLruTest, Threads) {
group.SetExecutor(std::ref(executor));

for (int32_t thread_id = 0; thread_id < 4; ++thread_id) {
group.Submit([&] -> Status {
group.Submit([&]() -> Status {
for (int32_t i = 0; i < 100; ++i) {
const int32_t key = i % 8;
if (memoized(key) != key * 2) {
Expand Down
2 changes: 1 addition & 1 deletion src/iceberg/test/merging_snapshot_update_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -1028,7 +1028,7 @@ TEST_F(MergingSnapshotUpdateTest, WriteDeleteGroups) {
op->WriteManifestsWith(executor, 3);

constexpr size_t kFileCount = 15'000;
auto files = std::views::iota(0UZ, kFileCount) |
auto files = std::views::iota(size_t{0}, kFileCount) |
std::views::transform([this](size_t index) {
return MakeDeleteFile(std::format("/delete/group_{}.parquet", index),
static_cast<int64_t>(index % 2));
Expand Down
2 changes: 1 addition & 1 deletion src/iceberg/update/snapshot_update.cc
Original file line number Diff line number Diff line change
Expand Up @@ -103,7 +103,7 @@ Result<std::vector<ManifestFile>> WriteManifestGroups(OptionalExecutor executor,
// TODO(zehua): Replace the manual offset calculation with `std::views::chunk`
// once the supported libc++ provides it.
auto groups =
std::views::iota(0UZ, group_count) |
std::views::iota(size_t{0}, group_count) |
std::views::transform([files, group_size](size_t group_index) {
const size_t offset = group_index * group_size;
return files.subspan(offset, std::min(group_size, files.size() - offset));
Expand Down
138 changes: 78 additions & 60 deletions src/iceberg/util/executor_util_internal.h
Original file line number Diff line number Diff line change
Expand Up @@ -82,6 +82,19 @@ concept ParallelCollectible =
requires ParallelReducible<ParallelCollectValueT<InputRange, Task>, Options...>;
};

// Checked through a class template rather than a lambda in the requires-clause because
// older MSVC versions (e.g. 14.42) cannot expand the Args pack inside such a lambda.
template <typename Indices, typename ArgsTuple, auto... Options>
struct ParallelCollectibleArgs : std::false_type {};

template <std::size_t... I, typename... Args, auto... Options>
struct ParallelCollectibleArgs<std::index_sequence<I...>, std::tuple<Args...>, Options...>
: std::bool_constant<(
ParallelCollectible<typename ParallelCollectTraits<I, Args...>::input_type,
typename ParallelCollectTraits<I, Args...>::task_type,
Options...> &&
...)> {};

} // namespace internal

template <typename... Args>
Expand Down Expand Up @@ -157,82 +170,87 @@ struct ParallelReduce<std::tuple<Ts...>> {
}
};

// The helpers below replace lambdas that were expanded over the pair index pack, which
// older MSVC versions (e.g. 14.42) cannot compile.
namespace internal {

template <typename ArgsTuple, std::size_t... I>
auto MakeParallelCollectValues(ArgsTuple& args_tuple, std::index_sequence<I...>) {
return std::tuple{
std::vector<ParallelCollectValueT<std::tuple_element_t<I * 2, ArgsTuple>,
std::tuple_element_t<I * 2 + 1, ArgsTuple>>>(
std::ranges::size(std::get<I * 2>(args_tuple)))...};
}

template <auto... Options, typename ValuesTuple, std::size_t... I>
auto ReduceParallelCollectValues(ValuesTuple& values_tuple, std::index_sequence<I...>) {
if constexpr (sizeof...(I) == 1) {
return ParallelReduce<typename std::tuple_element_t<0, ValuesTuple>::value_type,
Options...>::Reduce(std::get<0>(values_tuple));
} else {
return std::tuple{
ParallelReduce<typename std::tuple_element_t<I, ValuesTuple>::value_type,
Options...>::Reduce(std::get<I>(values_tuple))...};
}
}

template <std::size_t I, typename Group, typename ArgsTuple, typename ValuesTuple>
void SubmitParallelCollectPair(Group& group, ArgsTuple& args_tuple,
ValuesTuple& values_tuple) {
using item_ref = std::ranges::range_reference_t<std::tuple_element_t<I * 2, ArgsTuple>>;

for (auto&& [item, value] :
std::views::zip(std::get<I * 2>(args_tuple), std::get<I>(values_tuple))) {
if constexpr (std::is_lvalue_reference_v<item_ref>) {
group.Submit([&]() -> Status {
ICEBERG_ASSIGN_OR_RAISE(value,
std::invoke(std::get<I * 2 + 1>(args_tuple), item));
return {};
});
} else {
group.Submit([&, item = std::move(item)]() mutable -> Status {
ICEBERG_ASSIGN_OR_RAISE(
value, std::invoke(std::get<I * 2 + 1>(args_tuple), std::move(item)));
return {};
});
}
}
}

template <typename Group, typename ArgsTuple, typename ValuesTuple, std::size_t... I>
void SubmitParallelCollectTasks(Group& group, ArgsTuple& args_tuple,
ValuesTuple& values_tuple, std::index_sequence<I...>) {
(SubmitParallelCollectPair<I>(group, args_tuple, values_tuple), ...);
}

} // namespace internal

template <auto... Options, typename... Args>
requires(sizeof...(Args) >= 2 && sizeof...(Args) % 2 == 0 &&
[]<std::size_t... I>(std::index_sequence<I...>) consteval {
return (internal::ParallelCollectible<
typename internal::ParallelCollectTraits<I, Args...>::input_type,
typename internal::ParallelCollectTraits<I, Args...>::task_type,
Options...> &&
...);
}(std::make_index_sequence<sizeof...(Args) / 2>{}))
requires(
sizeof...(Args) >= 2 && sizeof...(Args) % 2 == 0 &&
internal::ParallelCollectibleArgs<std::make_index_sequence<sizeof...(Args) / 2>,
std::tuple<Args...>, Options...>::value)
auto ParallelCollect(OptionalExecutor executor, Args&&... args) {
constexpr std::size_t pair_count = sizeof...(Args) / 2;
using indices = std::make_index_sequence<pair_count>;

auto args_tuple = std::forward_as_tuple(std::forward<Args>(args)...);
auto values_tuple = internal::MakeParallelCollectValues(args_tuple, indices{});

auto values_tuple = [&]<std::size_t... I>(std::index_sequence<I...>) {
return std::tuple{[&] {
using traits = internal::ParallelCollectTraits<I, Args...>;

return std::vector<typename traits::value_type>(
std::ranges::size(std::get<I * 2>(args_tuple)));
}()...};
}(indices{});

auto reduce_all = [&]<std::size_t... I>(std::index_sequence<I...>) {
auto reduce_one = [&]<std::size_t PairIndex> {
using traits = internal::ParallelCollectTraits<PairIndex, Args...>;
using value_type = typename traits::value_type;
return ParallelReduce<value_type, Options...>::Reduce(
std::get<PairIndex>(values_tuple));
};

if constexpr (pair_count == 1) {
return reduce_one.template operator()<0>();
} else {
return std::tuple{reduce_one.template operator()<I>()...};
}
};

using result_type = decltype(reduce_all(indices{}));
using result_type = decltype(internal::ReduceParallelCollectValues<Options...>(
values_tuple, indices{}));

TaskGroup group;
group.SetExecutor(executor);

[&]<std::size_t... I>(std::index_sequence<I...>) {
(
[&] {
using item_ref = std::ranges::range_reference_t<
typename internal::ParallelCollectTraits<I, Args...>::input_type>;

for (auto&& [item, value] :
std::views::zip(std::get<I * 2>(args_tuple), std::get<I>(values_tuple))) {
if constexpr (std::is_lvalue_reference_v<item_ref>) {
group.Submit([&]() -> Status {
ICEBERG_ASSIGN_OR_RAISE(
value, std::invoke(std::get<I * 2 + 1>(args_tuple), item));
return {};
});
} else {
group.Submit([&, item = std::move(item)]() mutable -> Status {
ICEBERG_ASSIGN_OR_RAISE(
value, std::invoke(std::get<I * 2 + 1>(args_tuple), std::move(item)));
return {};
});
}
}
}(),
...);
}(indices{});
internal::SubmitParallelCollectTasks(group, args_tuple, values_tuple, indices{});

auto status = std::move(group).Run();
if (!status.has_value()) {
return Result<result_type>(std::unexpected<Error>(status.error()));
}

return Result<result_type>(reduce_all(indices{}));
return Result<result_type>(
internal::ReduceParallelCollectValues<Options...>(values_tuple, indices{}));
}

} // namespace iceberg
Loading