marsupialtail commented on code in PR #13669: URL: https://github.com/apache/arrow/pull/13669#discussion_r944139454
########## cpp/src/arrow/compute/exec/spilling_util.cc: ########## @@ -0,0 +1,305 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +#include "spilling_util.h" +#include <mutex> + +namespace arrow +{ +namespace compute +{ + struct ArrayInfo + { + int64_t num_children; + std::array<std::shared_ptr<Buffer>, 3> bufs; + std::array<size_t, 3> sizes; + std::shared_ptr<ArrayData> dictionary; + }; + + struct SpillFile::BatchInfo + { + int64_t start; + std::vector<ArrayInfo> arrays; + }; + +#ifdef _WIN32 +#include "windows_compatibility.h" + +const FileHandle kInvalidHandle = INVALID_HANDLE_VALUE; + +static Result<FileHandle> OpenTemporaryFile() +{ + constexpr DWORD kTempFileNameSize = MAX_PATH + 1; + wchar_t tmp_name_buf[kTempFileNameSize]; + wchar_t tmp_path_buf[kTempFileNameSize]; + + DWORD ret; + ret = GetTempPath2W(kTempFileNameSize, tmp_path_buf); + if(ret > kTempFileNameSize || ret == 0) + return Status::IOError(); + if(GetTempFileNameW(tmp_path_buf, L"ARROW_TMP", 0, tmp_name_buf) == 0) + return Status::IOError(); + + HANDLE file_handle = CreateFileA( + tmp_name_buf, + GENERIC_READ | GENERIC_WRITE | FILE_APPEND_DATA, + 0, + NULL, + CREATE_ALWAYS, + FILE_FLAG_NO_BUFFERING | FILE_FLAG_OVERLAPPED | FILE_FLAG_DELETE_ON_CLOSE, + NULL); + if(file_handle == INVALID_HANDLE_VALUE) + return Status::IOError("Failed to create temp file"); + return file_handle; +} + +static Status CloseTemporaryFile(FileHandle handle) +{ + if(!CloseHandle(handle)) + return Status::IOError("Failed to close temp file"); + return Status::OK(); +} + +static Status WriteBatch_PlatformSpecific(FileHandle handle, const SpillFile::BatchInfo &info) +{ + OVERLAPPED overlapped; + int64_t offset = info.start; + for(const ArrayInfo &arr : info.arrays) + { + for(size_t i = 0; i < arr.bufs.size(); i++) + { + if(info.bufs[i] != 0) + { + overlapped.Offset = static_cast<DWORD>(offset); + overlapped.OffsetHigh = static_cast<DWORD>(offset >> 32); + if(!WriteFile( + handle, + info.bufs[i]->data(), + info.bufs[i]->size(), + NULL, + &overlapped)) + return Status::IOError("Failed to spill!"); + offset += info.sizes[i]; + info.bufs[i].reset(); + } + } + } + return Status::OK(); +} +#else +#include <sys/stat.h> +#include <sys/types.h> +#include <sys/uio.h> +#include <unistd.h> +#include <fcntl.h> + +Result<FileHandle> OpenTemporaryFile() +{ + static std::once_flag generate_tmp_file_name_flag; + + constexpr int kFileNameSize = 1024; + static char name[kFileNameSize]; + + char *name_ptr = name; + std::call_once(generate_tmp_file_name_flag, [name_ptr]() noexcept + { + const char *selectors[] = { "TMPDIR", "TMP", "TEMP", "TEMPDIR" }; + constexpr size_t kNumSelectors = sizeof(selectors) / sizeof(selectors[0]); +#ifdef __ANDROID__ + const char *backup = "/data/local/tmp/"; +#else + const char *backup = "/tmp/"; Review Comment: What if I want to spill to an attached NVME SSD that is mounted on its own directory? ########## cpp/src/arrow/compute/exec/spilling_util.cc: ########## @@ -0,0 +1,305 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +#include "spilling_util.h" +#include <mutex> + +namespace arrow +{ +namespace compute +{ + struct ArrayInfo + { + int64_t num_children; + std::array<std::shared_ptr<Buffer>, 3> bufs; + std::array<size_t, 3> sizes; + std::shared_ptr<ArrayData> dictionary; + }; + + struct SpillFile::BatchInfo + { + int64_t start; + std::vector<ArrayInfo> arrays; + }; + +#ifdef _WIN32 +#include "windows_compatibility.h" + +const FileHandle kInvalidHandle = INVALID_HANDLE_VALUE; + +static Result<FileHandle> OpenTemporaryFile() +{ + constexpr DWORD kTempFileNameSize = MAX_PATH + 1; + wchar_t tmp_name_buf[kTempFileNameSize]; + wchar_t tmp_path_buf[kTempFileNameSize]; + + DWORD ret; + ret = GetTempPath2W(kTempFileNameSize, tmp_path_buf); + if(ret > kTempFileNameSize || ret == 0) + return Status::IOError(); + if(GetTempFileNameW(tmp_path_buf, L"ARROW_TMP", 0, tmp_name_buf) == 0) + return Status::IOError(); + + HANDLE file_handle = CreateFileA( + tmp_name_buf, + GENERIC_READ | GENERIC_WRITE | FILE_APPEND_DATA, + 0, + NULL, + CREATE_ALWAYS, + FILE_FLAG_NO_BUFFERING | FILE_FLAG_OVERLAPPED | FILE_FLAG_DELETE_ON_CLOSE, + NULL); + if(file_handle == INVALID_HANDLE_VALUE) + return Status::IOError("Failed to create temp file"); + return file_handle; +} + +static Status CloseTemporaryFile(FileHandle handle) +{ + if(!CloseHandle(handle)) + return Status::IOError("Failed to close temp file"); + return Status::OK(); +} + +static Status WriteBatch_PlatformSpecific(FileHandle handle, const SpillFile::BatchInfo &info) +{ + OVERLAPPED overlapped; + int64_t offset = info.start; + for(const ArrayInfo &arr : info.arrays) + { + for(size_t i = 0; i < arr.bufs.size(); i++) + { + if(info.bufs[i] != 0) + { + overlapped.Offset = static_cast<DWORD>(offset); + overlapped.OffsetHigh = static_cast<DWORD>(offset >> 32); + if(!WriteFile( + handle, + info.bufs[i]->data(), + info.bufs[i]->size(), + NULL, + &overlapped)) + return Status::IOError("Failed to spill!"); + offset += info.sizes[i]; + info.bufs[i].reset(); + } + } + } + return Status::OK(); +} +#else +#include <sys/stat.h> +#include <sys/types.h> +#include <sys/uio.h> +#include <unistd.h> +#include <fcntl.h> + +Result<FileHandle> OpenTemporaryFile() +{ + static std::once_flag generate_tmp_file_name_flag; + + constexpr int kFileNameSize = 1024; + static char name[kFileNameSize]; + + char *name_ptr = name; + std::call_once(generate_tmp_file_name_flag, [name_ptr]() noexcept + { + const char *selectors[] = { "TMPDIR", "TMP", "TEMP", "TEMPDIR" }; + constexpr size_t kNumSelectors = sizeof(selectors) / sizeof(selectors[0]); +#ifdef __ANDROID__ + const char *backup = "/data/local/tmp/"; +#else + const char *backup = "/tmp/"; +#endif + const char *tmp_dir = backup; + for(size_t i = 0; i < kNumSelectors; i++) + { + const char *env = getenv(selectors[i]); + if(env) + { + tmp_dir = env; + break; + } + } + size_t tmp_dir_length = std::strlen(tmp_dir); + + const char *tmp_name_template = "/ARROW_TMP_XXXXXX"; + size_t tmp_name_length = std::strlen(tmp_name_template); + + if((tmp_dir_length + tmp_name_length) >= kFileNameSize) + { + tmp_dir = backup; + tmp_dir_length = std::strlen(backup); + } + + std::strncpy(name_ptr, tmp_dir, kFileNameSize); + std::strncpy(name_ptr + tmp_dir_length, tmp_name_template, kFileNameSize - tmp_dir_length); + }); + +#ifdef __APPLE__ + int fd = mkstemp(name); + if(fd == kInvalidHandle) + return Status::IOError(strerror(errno)); + if(fcntl(fd, F_NOCACHE, 1) == -1) + return Status::IOError(strerror(errno)); +#else + int fd = mkostemp(name, O_DIRECT); + if(fd == kInvalidHandle) + return Status::IOError(strerror(errno)); +#endif + + if(unlink(name) != 0) + return Status::IOError(strerror(errno)); + return fd; +} + +static Status CloseTemporaryFile(FileHandle handle) +{ + if(close(handle) == -1) + return Status::IOError(strerror(errno)); + return Status::OK(); +} + +static Status WriteBatch_PlatformSpecific(FileHandle handle, SpillFile::BatchInfo &info) +{ + std::vector<struct iovec> ios; + for(const ArrayInfo &arr : info.arrays) + { + for(int i = 0; i < 3; i++) + { + if(arr.bufs[i]) + { + struct iovec io; + io.iov_base = static_cast<void *>(arr.bufs[i]->mutable_data()); + io.iov_len = static_cast<size_t>(arr.bufs[i]->size()); + ios.push_back(io); + } + } + } + + if(pwritev(handle, ios.data(), static_cast<int>(ios.size()), info.start) == -1) + return Status::IOError("Failed to spill!"); Review Comment: I seem to recall a discussion here where we talked about the performance of using pwritev versus things like IO uring where you were able to saturate NVME SSD bandwidth. Were you able to saturate SSD with pwritev? I understand that when you are spilling many batches there might be many pwritevs happening at the same time. Still I am curious how the perf compares to IO uring -- this is to satisfy my (and maybe other people's) curiosity not to point out a problem with your code. -- 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]
