Skip to content
Draft
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
6 changes: 3 additions & 3 deletions cpp/include/cudf/detail/sorting.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,7 @@ std::unique_ptr<column> sorted_order(table_view const& input,
std::vector<order> const& column_order,
std::vector<null_order> const& null_precedence,
cuda::stream_ref stream,
rmm::device_async_resource_ref mr);
cudf::memory_resources mr);

/**
* @copydoc cudf::stable_sorted_order
Expand All @@ -38,7 +38,7 @@ std::unique_ptr<column> stable_sorted_order(table_view const& input,
std::vector<order> const& column_order,
std::vector<null_order> const& null_precedence,
cuda::stream_ref stream,
rmm::device_async_resource_ref mr);
cudf::memory_resources mr);

/**
* @copydoc cudf::sort_by_key
Expand All @@ -64,7 +64,7 @@ std::unique_ptr<column> rank(column_view const& input,
null_order null_precedence,
bool percentage,
cuda::stream_ref stream,
rmm::device_async_resource_ref mr);
cudf::memory_resources mr);

/**
* @copydoc cudf::stable_sort_by_key
Expand Down
25 changes: 12 additions & 13 deletions cpp/include/cudf/sorting.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,7 @@ namespace CUDF_EXPORT cudf {
* for each column. Size must be equal to `input.num_columns()` or empty.
* If empty, all columns will be sorted in `null_order::BEFORE`.
* @param stream CUDA stream used for device memory operations and kernel launches
* @param mr Device memory resource used to allocate the returned column's device memory
* @param mr Memory resources used for temporary allocations and the returned column
* @return A non-nullable column of elements containing the permuted row indices of
* `input` if it were sorted
*/
Expand All @@ -46,7 +46,7 @@ std::unique_ptr<column> sorted_order(
std::vector<order> const& column_order = {},
std::vector<null_order> const& null_precedence = {},
cuda::stream_ref stream = cudf::get_default_stream(),
rmm::device_async_resource_ref mr = cudf::get_current_device_resource_ref());
cudf::memory_resources mr = cudf::get_current_device_resource_ref());

/**
* @brief Computes the row indices that would produce `input` in a stable
Expand All @@ -61,7 +61,7 @@ std::unique_ptr<column> stable_sorted_order(
std::vector<order> const& column_order = {},
std::vector<null_order> const& null_precedence = {},
cuda::stream_ref stream = cudf::get_default_stream(),
rmm::device_async_resource_ref mr = cudf::get_current_device_resource_ref());
cudf::memory_resources mr = cudf::get_current_device_resource_ref());

/**
* @brief Checks whether the rows of a `table` are sorted in a lexicographical
Expand Down Expand Up @@ -216,18 +216,17 @@ std::unique_ptr<table> stable_sort_by_key(
* @param null_precedence The desired order of null rows compared to other elements
* @param percentage Flag to convert ranks to percentage in range (0,1]
* @param stream CUDA stream used for device memory operations and kernel launches
* @param mr Device memory resource used to allocate the returned column's device memory
* @param mr Memory resources used for temporary allocations and the returned column
* @return A column of containing the rank of the each element of the column of `input`
*/
std::unique_ptr<column> rank(
column_view const& input,
rank_method method,
order column_order,
null_policy null_handling,
null_order null_precedence,
bool percentage,
cuda::stream_ref stream = cudf::get_default_stream(),
rmm::device_async_resource_ref mr = cudf::get_current_device_resource_ref());
std::unique_ptr<column> rank(column_view const& input,
rank_method method,
order column_order,
null_policy null_handling,
null_order null_precedence,
bool percentage,
cuda::stream_ref stream = cudf::get_default_stream(),
cudf::memory_resources mr = cudf::get_current_device_resource_ref());

/**
* @brief Returns sorted order after sorting each segment in the table.
Expand Down
95 changes: 56 additions & 39 deletions cpp/src/sort/rank.cu
Original file line number Diff line number Diff line change
Expand Up @@ -57,17 +57,17 @@ struct unique_functor {
// Assign rank from 1 to n unique values. Equal values get same rank value.
rmm::device_uvector<size_type> sorted_dense_rank(column_view input_col,
column_view sorted_order_view,
cuda::stream_ref stream)
cuda::stream_ref stream,
rmm::device_async_resource_ref temp_mr)
{
auto const t_input = table_view{{input_col}};
auto const temp_mr = cudf::get_current_device_resource_ref();
auto const comparator = cudf::detail::row::equality::self_comparator{t_input, stream, temp_mr};

auto const sorted_index_order = cuda::make_permutation_iterator(
sorted_order_view.begin<size_type>(), cuda::counting_iterator<size_type>{0});

auto const input_size = input_col.size();
rmm::device_uvector<size_type> dense_rank_sorted(input_size, stream);
rmm::device_uvector<size_type> dense_rank_sorted(input_size, stream, temp_mr);

auto const comparator_helper = [&](auto const device_comparator) {
thrust::transform(rmm::exec_policy_nosync(stream, temp_mr),
Expand Down Expand Up @@ -120,13 +120,14 @@ void tie_break_ranks_transform(cudf::device_span<size_type const> dense_rank_sor
outputIterator rank_iter,
TieBreaker tie_breaker,
Transformer transformer,
cuda::stream_ref stream)
cuda::stream_ref stream,
rmm::device_async_resource_ref temp_mr)
{
auto const input_size = sorted_order_view.size();
// algorithm: reduce_by_key(dense_rank, 1, n, reduction_tie_breaker)
// reduction_tie_breaker = min, max, min_count
rmm::device_uvector<TieType> tie_sorted(sorted_order_view.size(), stream);
thrust::reduce_by_key(rmm::exec_policy_nosync(stream, cudf::get_current_device_resource_ref()),
rmm::device_uvector<TieType> tie_sorted(sorted_order_view.size(), stream, temp_mr);
thrust::reduce_by_key(rmm::exec_policy_nosync(stream, temp_mr),
dense_rank_sorted.begin(),
dense_rank_sorted.end(),
tie_iter,
Expand All @@ -142,7 +143,7 @@ void tie_break_ranks_transform(cudf::device_span<size_type const> dense_rank_sor
[tied_rank = tie_sorted.begin(), transformer] __device__(auto dense_pos) {
return transformer(tied_rank[dense_pos - 1]);
}));
thrust::scatter(rmm::exec_policy_nosync(stream, cudf::get_current_device_resource_ref()),
thrust::scatter(rmm::exec_policy_nosync(stream, temp_mr),
sorted_tied_rank,
sorted_tied_rank + input_size,
sorted_order_view.begin<size_type>(),
Expand All @@ -152,10 +153,11 @@ void tie_break_ranks_transform(cudf::device_span<size_type const> dense_rank_sor
template <typename outputType>
void rank_first(column_view sorted_order_view,
mutable_column_view rank_mutable_view,
cuda::stream_ref stream)
cuda::stream_ref stream,
rmm::device_async_resource_ref temp_mr)
{
// stable sort order ranking (no ties)
thrust::scatter(rmm::exec_policy_nosync(stream, cudf::get_current_device_resource_ref()),
thrust::scatter(rmm::exec_policy_nosync(stream, temp_mr),
cuda::counting_iterator<size_type>{1},
cuda::counting_iterator<size_type>{rank_mutable_view.size() + 1},
sorted_order_view.begin<size_type>(),
Expand All @@ -166,10 +168,11 @@ template <typename outputType>
void rank_dense(cudf::device_span<size_type const> dense_rank_sorted,
column_view sorted_order_view,
mutable_column_view rank_mutable_view,
cuda::stream_ref stream)
cuda::stream_ref stream,
rmm::device_async_resource_ref temp_mr)
{
// All equal values have same rank and rank always increases by 1 between groups
thrust::scatter(rmm::exec_policy_nosync(stream, cudf::get_current_device_resource_ref()),
thrust::scatter(rmm::exec_policy_nosync(stream, temp_mr),
dense_rank_sorted.begin(),
dense_rank_sorted.end(),
sorted_order_view.begin<size_type>(),
Expand All @@ -180,7 +183,8 @@ template <typename outputType>
void rank_min(cudf::device_span<size_type const> group_keys,
column_view sorted_order_view,
mutable_column_view rank_mutable_view,
cuda::stream_ref stream)
cuda::stream_ref stream,
rmm::device_async_resource_ref temp_mr)
{
// min of first in the group
// All equal values have min of ranks among them.
Expand All @@ -191,14 +195,16 @@ void rank_min(cudf::device_span<size_type const> group_keys,
rank_mutable_view.begin<outputType>(),
cuda::minimum{},
cuda::std::identity{},
stream);
stream,
temp_mr);
}

template <typename outputType>
void rank_max(cudf::device_span<size_type const> group_keys,
column_view sorted_order_view,
mutable_column_view rank_mutable_view,
cuda::stream_ref stream)
cuda::stream_ref stream,
rmm::device_async_resource_ref temp_mr)
{
// max of first in the group
// All equal values have max of ranks among them.
Expand All @@ -209,7 +215,8 @@ void rank_max(cudf::device_span<size_type const> group_keys,
rank_mutable_view.begin<outputType>(),
cuda::maximum{},
cuda::std::identity{},
stream);
stream,
temp_mr);
}

// Returns index, count
Expand All @@ -221,7 +228,8 @@ struct index_counter {
void rank_average(cudf::device_span<size_type const> group_keys,
column_view sorted_order_view,
mutable_column_view rank_mutable_view,
cuda::stream_ref stream)
cuda::stream_ref stream,
rmm::device_async_resource_ref temp_mr)
{
// k, k+1, .. k+n-1
// average = (n*k+ n*(n-1)/2)/n
Expand All @@ -245,7 +253,8 @@ void rank_average(cudf::device_span<size_type const> group_keys,
return static_cast<double>(cuda::std::get<0>(minrank_count)) +
(static_cast<double>(cuda::std::get<1>(minrank_count)) - 1) / 2.0;
}),
stream);
stream,
temp_mr);
}

} // anonymous namespace
Expand All @@ -257,77 +266,85 @@ std::unique_ptr<column> rank(column_view const& input,
null_order null_precedence,
bool percentage,
cuda::stream_ref stream,
rmm::device_async_resource_ref mr)
cudf::memory_resources mr)
{
auto const output_mr = mr.get_output_mr();
auto const temp_mr = mr.get_temporary_mr();
data_type const output_type = (percentage or method == rank_method::AVERAGE)
? data_type(type_id::FLOAT64)
: data_type(type_to_id<size_type>());
std::unique_ptr<column> rank_column = [&null_handling, &output_type, &input, &stream, &mr] {
std::unique_ptr<column> rank_column = [&null_handling, &output_type, &input, &stream, output_mr] {
// na_option=keep assign NA to NA values
if (null_handling == null_policy::EXCLUDE)
return make_numeric_column(output_type,
input.size(),
detail::copy_bitmask(input, stream, mr),
detail::copy_bitmask(input, stream, output_mr),
input.null_count(),
stream,
mr);
output_mr);
else
return make_numeric_column(output_type, input.size(), mask_state::UNALLOCATED, stream, mr);
return make_numeric_column(
output_type, input.size(), mask_state::UNALLOCATED, stream, output_mr);
}();
auto rank_mutable_view = rank_column->mutable_view();

std::unique_ptr<column> sorted_order =
(method == rank_method::FIRST)
? detail::stable_sorted_order(
table_view{{input}}, {column_order}, {null_precedence}, stream, mr)
: detail::sorted_order(table_view{{input}}, {column_order}, {null_precedence}, stream, mr);
table_view{{input}}, {column_order}, {null_precedence}, stream, {temp_mr, temp_mr})
: detail::sorted_order(
table_view{{input}}, {column_order}, {null_precedence}, stream, {temp_mr, temp_mr});
column_view sorted_order_view = sorted_order->view();

// dense: All equal values have same rank and rank always increases by 1 between groups
// acts as key for min, max, average to denote equal value groups
rmm::device_uvector<size_type> const dense_rank_sorted =
[&method, &input, &sorted_order_view, &stream] {
[&method, &input, &sorted_order_view, &stream, temp_mr] {
if (method != rank_method::FIRST)
return sorted_dense_rank(input, sorted_order_view, stream);
return sorted_dense_rank(input, sorted_order_view, stream, temp_mr);
else
return rmm::device_uvector<size_type>(0, stream);
return rmm::device_uvector<size_type>(0, stream, temp_mr);
}();

if (output_type.id() == type_id::FLOAT64) {
switch (method) {
case rank_method::FIRST:
rank_first<double>(sorted_order_view, rank_mutable_view, stream);
rank_first<double>(sorted_order_view, rank_mutable_view, stream, temp_mr);
break;
case rank_method::DENSE:
rank_dense<double>(dense_rank_sorted, sorted_order_view, rank_mutable_view, stream);
rank_dense<double>(
dense_rank_sorted, sorted_order_view, rank_mutable_view, stream, temp_mr);
break;
case rank_method::MIN:
rank_min<double>(dense_rank_sorted, sorted_order_view, rank_mutable_view, stream);
rank_min<double>(dense_rank_sorted, sorted_order_view, rank_mutable_view, stream, temp_mr);
break;
case rank_method::MAX:
rank_max<double>(dense_rank_sorted, sorted_order_view, rank_mutable_view, stream);
rank_max<double>(dense_rank_sorted, sorted_order_view, rank_mutable_view, stream, temp_mr);
break;
case rank_method::AVERAGE:
rank_average(dense_rank_sorted, sorted_order_view, rank_mutable_view, stream);
rank_average(dense_rank_sorted, sorted_order_view, rank_mutable_view, stream, temp_mr);
break;
default: CUDF_FAIL("Unexpected rank_method for rank()");
}
} else {
switch (method) {
case rank_method::FIRST:
rank_first<size_type>(sorted_order_view, rank_mutable_view, stream);
rank_first<size_type>(sorted_order_view, rank_mutable_view, stream, temp_mr);
break;
case rank_method::DENSE:
rank_dense<size_type>(dense_rank_sorted, sorted_order_view, rank_mutable_view, stream);
rank_dense<size_type>(
dense_rank_sorted, sorted_order_view, rank_mutable_view, stream, temp_mr);
break;
case rank_method::MIN:
rank_min<size_type>(dense_rank_sorted, sorted_order_view, rank_mutable_view, stream);
rank_min<size_type>(
dense_rank_sorted, sorted_order_view, rank_mutable_view, stream, temp_mr);
break;
case rank_method::MAX:
rank_max<size_type>(dense_rank_sorted, sorted_order_view, rank_mutable_view, stream);
rank_max<size_type>(
dense_rank_sorted, sorted_order_view, rank_mutable_view, stream, temp_mr);
break;
case rank_method::AVERAGE:
rank_average(dense_rank_sorted, sorted_order_view, rank_mutable_view, stream);
rank_average(dense_rank_sorted, sorted_order_view, rank_mutable_view, stream, temp_mr);
break;
default: CUDF_FAIL("Unexpected rank_method for rank()");
}
Expand All @@ -341,7 +358,7 @@ std::unique_ptr<column> rank(column_view const& input,
auto drs = dense_rank_sorted.data();
bool const is_dense = (method == rank_method::DENSE);
thrust::transform(
rmm::exec_policy_nosync(stream, cudf::get_current_device_resource_ref()),
rmm::exec_policy_nosync(stream, temp_mr),
rank_iter,
rank_iter + input.size(),
rank_iter,
Expand All @@ -360,7 +377,7 @@ std::unique_ptr<column> rank(column_view const& input,
null_order null_precedence,
bool percentage,
cuda::stream_ref stream,
rmm::device_async_resource_ref mr)
cudf::memory_resources mr)
{
CUDF_FUNC_RANGE();
return detail::rank(
Expand Down
4 changes: 2 additions & 2 deletions cpp/src/sort/sort.cu
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@ std::unique_ptr<column> sorted_order(table_view const& input,
std::vector<order> const& column_order,
std::vector<null_order> const& null_precedence,
cuda::stream_ref stream,
rmm::device_async_resource_ref mr)
cudf::memory_resources mr)
{
return sorted_order<sort_method::UNSTABLE>(input, column_order, null_precedence, stream, mr);
}
Expand Down Expand Up @@ -72,7 +72,7 @@ std::unique_ptr<column> sorted_order(table_view const& input,
std::vector<order> const& column_order,
std::vector<null_order> const& null_precedence,
cuda::stream_ref stream,
rmm::device_async_resource_ref mr)
cudf::memory_resources mr)
{
CUDF_FUNC_RANGE();
return detail::sorted_order(input, column_order, null_precedence, stream, mr);
Expand Down
4 changes: 2 additions & 2 deletions cpp/src/sort/sort.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -29,15 +29,15 @@ enum class sort_method : bool { STABLE, UNSTABLE };
* @param column_order Ascending or descending sort order
* @param null_precedence How null rows are to be ordered
* @param stream CUDA stream used for device memory operations and kernel launches
* @param mr Device memory resource used to allocate the returned column's device memory
* @param mr Memory resources used for temporary allocations and the returned column
* @return Sorted indices for the input column.
*/
template <sort_method method>
std::unique_ptr<column> sorted_order(column_view const& input,
order column_order,
null_order null_precedence,
cuda::stream_ref stream,
rmm::device_async_resource_ref mr);
cudf::memory_resources mr);

} // namespace detail
} // namespace cudf
Loading
Loading