marin-ma commented on code in PR #13067:
URL: https://github.com/apache/gluten/pull/13067#discussion_r4198810656
##########
cpp/velox/tests/ValueStreamDynamicFilterTest.cc:
##########
@@ -279,4 +287,86 @@ TEST_F(ValueStreamDynamicFilterTest, canAddDynamicFilter) {
ASSERT_EQ(end, nullptr);
}
+#ifdef GLUTEN_ENABLE_GPU
+TEST_F(ValueStreamDynamicFilterTest,
cudfValueStreamConvertsHostBatchToCudfVector) {
Review Comment:
Are these tests able to run without `cudf_velox::registerCudf();` and `
cudf_velox::unregisterCudf();`?
##########
cpp/velox/operators/plannodes/CudfVectorStream.h:
##########
@@ -115,11 +116,30 @@ class CudfVectorStream : public CudfVectorStreamBase {
VELOX_DCHECK(vp != nullptr);
auto cudfVector =
std::dynamic_pointer_cast<facebook::velox::cudf_velox::CudfVector>(vp);
if (cudfVector == nullptr) {
- // The vector may comes from BroadcastExchange, in this case, it's not a
CudfVector.
- vp->setType(outputType_);
- return vp;
+ // BroadcastExchange may return a host RowVector; upload it for GPU
operators.
+ auto stream =
facebook::velox::cudf_velox::cudfGlobalStreamPool().get_stream();
+ if (vp->childrenSize() == 0 || outputType_->size() == 0) {
+ // Preserve row count because zero-column cuDF tables cannot store it.
+ return std::make_shared<facebook::velox::cudf_velox::CudfVector>(
+ vp->pool(), outputType_, vp->size(),
std::make_unique<cudf::table>(), stream);
+ }
+ // Drop extra trailing columns added by broadcast exchange.
Review Comment:
Have you run into issue with trailing columns being added by broadcast
exchange? Seems the projection for the extra column is added by
`spark.gluten.velox.buildHashTableOncePerExecutor.enabled=true`, but the option
is not supported with cudf.
##########
cpp/velox/operators/plannodes/CudfVectorStream.h:
##########
@@ -115,11 +116,30 @@ class CudfVectorStream : public CudfVectorStreamBase {
VELOX_DCHECK(vp != nullptr);
auto cudfVector =
std::dynamic_pointer_cast<facebook::velox::cudf_velox::CudfVector>(vp);
if (cudfVector == nullptr) {
- // The vector may comes from BroadcastExchange, in this case, it's not a
CudfVector.
- vp->setType(outputType_);
- return vp;
+ // BroadcastExchange may return a host RowVector; upload it for GPU
operators.
+ auto stream =
facebook::velox::cudf_velox::cudfGlobalStreamPool().get_stream();
+ if (vp->childrenSize() == 0 || outputType_->size() == 0) {
+ // Preserve row count because zero-column cuDF tables cannot store it.
+ return std::make_shared<facebook::velox::cudf_velox::CudfVector>(
+ vp->pool(), outputType_, vp->size(),
std::make_unique<cudf::table>(), stream);
+ }
+ // Drop extra trailing columns added by broadcast exchange.
+ VELOX_CHECK_GE(
+ vp->childrenSize(),
+ outputType_->size(),
+ "Value stream batch has fewer columns than the declared output
type");
+ std::vector<facebook::velox::VectorPtr> children(
+ vp->children().begin(), vp->children().begin() +
outputType_->size());
+ for (auto& child : children) {
Review Comment:
Any reason to trigger the lazy vector loading here? The vectors should be
materialised by the arrow bridge immediately after in
`facebook::velox::cudf_velox::with_arrow::toCudfTable` anyway.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]