From 230fbc14003d1493ea8b08ae3a9f2bc45bd8a6ac Mon Sep 17 00:00:00 2001 From: zhixingheyi-tian Date: Tue, 19 Jul 2022 13:12:16 +0800 Subject: [PATCH 1/2] Add close() --- .../scala/com/intel/oap/execution/ArrowColumnarToRowExec.scala | 3 +++ 1 file changed, 3 insertions(+) diff --git a/native-sql-engine/core/src/main/scala/com/intel/oap/execution/ArrowColumnarToRowExec.scala b/native-sql-engine/core/src/main/scala/com/intel/oap/execution/ArrowColumnarToRowExec.scala index 6a34b882c..54b5cc5de 100644 --- a/native-sql-engine/core/src/main/scala/com/intel/oap/execution/ArrowColumnarToRowExec.scala +++ b/native-sql-engine/core/src/main/scala/com/intel/oap/execution/ArrowColumnarToRowExec.scala @@ -140,6 +140,9 @@ case class ArrowColumnarToRowExec(child: SparkPlan) extends ColumnarToRowTransit arrowSchema, batch.numRows, bufAddrs.toArray, bufSizes.toArray, SparkMemoryUtils.contextMemoryPool().getNativeInstanceId) + // Because nativeConvertColumnarToRow didn't need Buffers of batch anymore, so close immediately + batch.close() + convertTime += NANOSECONDS.toMillis(System.nanoTime() - beforeConvert) new Iterator[InternalRow] { From 89f3f3d4824eb3a0a10730d5d53b42bae1efc8c9 Mon Sep 17 00:00:00 2001 From: zhixingheyi-tian Date: Wed, 20 Jul 2022 17:08:37 +0800 Subject: [PATCH 2/2] Add reset in native --- .../scala/com/intel/oap/execution/ArrowColumnarToRowExec.scala | 2 +- .../cpp/src/operators/columnar_to_row_converter.cc | 3 +++ 2 files changed, 4 insertions(+), 1 deletion(-) diff --git a/native-sql-engine/core/src/main/scala/com/intel/oap/execution/ArrowColumnarToRowExec.scala b/native-sql-engine/core/src/main/scala/com/intel/oap/execution/ArrowColumnarToRowExec.scala index 54b5cc5de..c792cd3df 100644 --- a/native-sql-engine/core/src/main/scala/com/intel/oap/execution/ArrowColumnarToRowExec.scala +++ b/native-sql-engine/core/src/main/scala/com/intel/oap/execution/ArrowColumnarToRowExec.scala @@ -141,7 +141,7 @@ case class ArrowColumnarToRowExec(child: SparkPlan) extends ColumnarToRowTransit SparkMemoryUtils.contextMemoryPool().getNativeInstanceId) // Because nativeConvertColumnarToRow didn't need Buffers of batch anymore, so close immediately - batch.close() +// batch.close() convertTime += NANOSECONDS.toMillis(System.nanoTime() - beforeConvert) diff --git a/native-sql-engine/cpp/src/operators/columnar_to_row_converter.cc b/native-sql-engine/cpp/src/operators/columnar_to_row_converter.cc index 51626f66e..60a3101cd 100644 --- a/native-sql-engine/cpp/src/operators/columnar_to_row_converter.cc +++ b/native-sql-engine/cpp/src/operators/columnar_to_row_converter.cc @@ -608,6 +608,9 @@ arrow::Status ColumnarToRowConverter::Write() { support_avx512_); } + // Because didn't need rb_ anymore here, so reset immediately. + rb_.reset(); + return arrow::Status::OK(); }