diff --git a/cpp/CMakeLists.txt b/cpp/CMakeLists.txt index c39bcbc28..d2fba4a33 100755 --- a/cpp/CMakeLists.txt +++ b/cpp/CMakeLists.txt @@ -389,6 +389,13 @@ if (BUILD_BENCHMARK) ${PROJECT_SRC_DIR}) target_link_libraries(read_backend_benchmark PRIVATE tsfile Threads::Threads) + + add_executable(alp_vs_gorilla_benchmark + bench_mark/alp_vs_gorilla_benchmark.cc) + target_include_directories(alp_vs_gorilla_benchmark PRIVATE + ${PROJECT_SRC_DIR}) + target_link_libraries(alp_vs_gorilla_benchmark PRIVATE + tsfile Threads::Threads) endif () if (BUILD_TOOLS) add_subdirectory(tools) diff --git a/cpp/bench_mark/alp_vs_gorilla_benchmark.cc b/cpp/bench_mark/alp_vs_gorilla_benchmark.cc new file mode 100644 index 000000000..c953a998c --- /dev/null +++ b/cpp/bench_mark/alp_vs_gorilla_benchmark.cc @@ -0,0 +1,373 @@ +/* + * 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 +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#include "common/allocator/byte_stream.h" +#include "common/db_common.h" +#include "encoding/decoder_factory.h" +#include "encoding/encoder_factory.h" +#include "encoding/ts2diff_encoder.h" +#include "utils/errno_define.h" + +namespace { + +int g_ts2diff_mpn = -1; + +using Clock = std::chrono::high_resolution_clock; + +template +struct BenchResult { + std::string encoding; + std::string pattern; + uint64_t value_count; + uint64_t raw_bytes; + uint64_t encoded_bytes; + double encode_ms; + double decode_ms; + double encode_mbps; + double decode_mbps; + bool correct; +}; + +template +bool BitEqual(T lhs, T rhs) { + return std::memcmp(&lhs, &rhs, sizeof(T)) == 0; +} + +int CopyEncodedBytes(common::ByteStream& stream, std::vector& bytes) { + bytes.clear(); + bytes.reserve(static_cast(stream.total_size())); + common::ByteStream::BufferIterator iter = stream.init_buffer_iterator(); + while (true) { + common::ByteStream::Buffer buffer = iter.get_next_buf(); + if (buffer.buf_ == NULL || buffer.len_ == 0) { + break; + } + const uint8_t* begin = reinterpret_cast(buffer.buf_); + bytes.insert(bytes.end(), begin, begin + buffer.len_); + } + return bytes.size() == stream.total_size() ? common::E_OK + : common::E_PARTIAL_READ; +} + +template +int EncodeOnce(common::TSEncoding encoding, common::TSDataType data_type, + const std::vector& values, std::vector& bytes, + double& elapsed_ms) { + storage::Encoder* encoder = + storage::EncoderFactory::alloc_value_encoder(encoding, data_type); + if (encoder == NULL) { + return common::E_NOT_SUPPORT; + } + if (encoding == common::TS_2DIFF && g_ts2diff_mpn >= 0) { + if (data_type == common::FLOAT) { + static_cast(encoder) + ->set_max_point_number(g_ts2diff_mpn); + } else if (data_type == common::DOUBLE) { + static_cast(encoder) + ->set_max_point_number(g_ts2diff_mpn); + } + } + common::ByteStream stream(1024, common::MOD_DEFAULT); + const Clock::time_point start = Clock::now(); + int ret = encoder->encode_batch( + &values[0], static_cast(values.size()), stream); + if (ret == common::E_OK) { + ret = encoder->flush(stream); + } + const Clock::time_point end = Clock::now(); + elapsed_ms = + std::chrono::duration_cast>( + end - start) + .count(); + if (ret == common::E_OK) { + ret = CopyEncodedBytes(stream, bytes); + } + storage::EncoderFactory::free(encoder); + return ret; +} + +template +int DecodeOnce(common::TSEncoding encoding, common::TSDataType data_type, + const std::vector& bytes, uint64_t count, + std::vector& decoded, double& elapsed_ms) { + common::ByteStream stream; + stream.wrap_from(reinterpret_cast(&bytes[0]), + static_cast(bytes.size())); + storage::Decoder* decoder = + storage::DecoderFactory::alloc_value_decoder(encoding, data_type); + if (decoder == NULL) { + return common::E_NOT_SUPPORT; + } + decoded.resize(static_cast(count)); + const Clock::time_point start = Clock::now(); + int ret = common::E_OK; + if (data_type == common::FLOAT) { + ret = decoder->read_exact_float( + reinterpret_cast(decoded.empty() ? NULL : &decoded[0]), + static_cast(count), stream); + } else { + ret = decoder->read_exact_double( + reinterpret_cast(decoded.empty() ? NULL : &decoded[0]), + static_cast(count), stream); + } + const Clock::time_point end = Clock::now(); + elapsed_ms = + std::chrono::duration_cast>( + end - start) + .count(); + storage::DecoderFactory::free(decoder); + return ret; +} + +template +BenchResult RunCase(const std::string& encoding_name, + common::TSEncoding encoding, const std::string& pattern, + const std::vector& values, uint32_t repetitions) { + const common::TSDataType data_type = + sizeof(T) == sizeof(float) ? common::FLOAT : common::DOUBLE; + BenchResult result; + result.encoding = encoding_name; + result.pattern = pattern; + result.value_count = values.size(); + result.raw_bytes = values.size() * sizeof(T); + result.encoded_bytes = 0; + result.encode_ms = std::numeric_limits::max(); + result.decode_ms = std::numeric_limits::max(); + result.encode_mbps = 0.0; + result.decode_mbps = 0.0; + result.correct = true; + + std::vector best_bytes; + std::vector decoded; + for (uint32_t rep = 0; rep < repetitions; ++rep) { + std::vector bytes; + double encode_ms = 0.0; + const int encode_ret = + EncodeOnce(encoding, data_type, values, bytes, encode_ms); + if (encode_ret != common::E_OK) { + result.correct = false; + return result; + } + double decode_ms = 0.0; + std::vector rep_decoded; + const int decode_ret = DecodeOnce( + encoding, data_type, bytes, values.size(), rep_decoded, decode_ms); + if (decode_ret != common::E_OK) { + result.correct = false; + return result; + } + if (rep == 0) { + best_bytes.swap(bytes); + decoded.swap(rep_decoded); + } + if (encode_ms < result.encode_ms) { + result.encode_ms = encode_ms; + } + if (decode_ms < result.decode_ms) { + result.decode_ms = decode_ms; + } + } + + for (size_t i = 0; i < values.size(); ++i) { + if (!BitEqual(values[i], decoded[i])) { + result.correct = false; + break; + } + } + + result.encoded_bytes = best_bytes.size(); + const double value_count = static_cast(values.size()); + if (result.encode_ms > 0.0) { + result.encode_mbps = + (static_cast(result.raw_bytes) / (1024.0 * 1024.0)) / + (result.encode_ms / 1000.0); + } + if (result.decode_ms > 0.0) { + result.decode_mbps = + (static_cast(result.raw_bytes) / (1024.0 * 1024.0)) / + (result.decode_ms / 1000.0); + } + return result; +} + +template +void PrintResult(const BenchResult& result) { + const double ratio = result.encoded_bytes == 0 + ? 0.0 + : static_cast(result.raw_bytes) / + static_cast(result.encoded_bytes); + std::cout << std::left << std::setw(8) << result.encoding << std::setw(12) + << result.pattern << std::right << std::setw(10) + << result.value_count << std::setw(14) << result.raw_bytes + << std::setw(14) << result.encoded_bytes << std::setw(10) + << std::fixed << std::setprecision(3) << ratio << std::setw(12) + << result.encode_ms << std::setw(12) << result.decode_ms + << std::setw(14) << result.encode_mbps << std::setw(14) + << result.decode_mbps << std::setw(10) + << (result.correct ? "yes" : "no") << std::endl; +} + +template +std::vector MakeDecimalData(uint64_t count, double scale, double offset) { + std::vector values(static_cast(count)); + for (uint64_t i = 0; i < count; ++i) { + const double raw = static_cast(i % 100000) * scale + offset; + values[static_cast(i)] = static_cast(raw); + } + return values; +} + +template +std::vector MakeSmoothData(uint64_t count) { + std::vector values(static_cast(count)); + for (uint64_t i = 0; i < count; ++i) { + const double x = static_cast(i) / 1000.0; + values[static_cast(i)] = + static_cast(std::sin(x) + 0.25 * std::cos(x * 0.5)); + } + return values; +} + +template +std::vector MakeRandomData(uint64_t count) { + std::vector values(static_cast(count)); + uint64_t state = 0x9e3779b97f4a7c15ULL; + for (uint64_t i = 0; i < count; ++i) { + state = state * 6364136223846793005ULL + 1442695040888963407ULL; + const uint64_t bits = state ^ (state >> 29); + if (sizeof(T) == sizeof(float)) { + uint32_t fbits = static_cast(bits); + float v = 0.0f; + std::memcpy(&v, &fbits, sizeof(v)); + if (!std::isfinite(v)) { + v = static_cast(i % 1024) * 0.001f; + } + values[static_cast(i)] = v; + } else { + double v = 0.0; + std::memcpy(&v, &bits, sizeof(v)); + if (!std::isfinite(v)) { + v = static_cast(i % 1024) * 0.001; + } + values[static_cast(i)] = v; + } + } + return values; +} + +template +std::vector MakeConstantData(uint64_t count) { + return std::vector(static_cast(count), + static_cast(3.141592653589793)); +} + +template +void RunAll(const std::string& type_name, uint32_t repetitions, + const std::string& codec_filter, + const std::string& pattern_filter) { + const uint64_t count = 1024 * 1024; + const common::TSDataType data_type = + sizeof(T) == sizeof(float) ? common::FLOAT : common::DOUBLE; + + struct Case { + std::string name; + std::vector values; + }; + std::vector cases; + cases.push_back(Case{"decimal", MakeDecimalData(count, 0.01, -1.0)}); + cases.push_back(Case{"smooth", MakeSmoothData(count)}); + cases.push_back(Case{"random", MakeRandomData(count)}); + cases.push_back(Case{"constant", MakeConstantData(count)}); + + std::cout << "==== " << type_name << " ====" << std::endl; + std::cout << std::left << std::setw(8) << "codec" << std::setw(12) + << "pattern" << std::right << std::setw(10) << "values" + << std::setw(14) << "raw(B)" << std::setw(14) << "encoded(B)" + << std::setw(10) << "ratio" << std::setw(12) << "enc(ms)" + << std::setw(12) << "dec(ms)" << std::setw(14) << "enc(MB/s)" + << std::setw(14) << "dec(MB/s)" << std::setw(10) << "ok" + << std::endl; + for (size_t c = 0; c < cases.size(); ++c) { + if (!pattern_filter.empty() && cases[c].name != pattern_filter) { + continue; + } + if (codec_filter.empty() || codec_filter == "ALP") { + PrintResult(RunCase("ALP", common::ALP, cases[c].name, + cases[c].values, repetitions)); + } + if (codec_filter.empty() || codec_filter == "GORILLA") { + PrintResult(RunCase("GORILLA", common::GORILLA, cases[c].name, + cases[c].values, repetitions)); + } + if (codec_filter.empty() || codec_filter == "TS_2DIFF") { + PrintResult(RunCase("TS_2DIFF", common::TS_2DIFF, cases[c].name, + cases[c].values, repetitions)); + } + if (codec_filter.empty() || codec_filter == "TS_2DIFF") { + PrintResult(RunCase("TS_2DIFF", common::TS_2DIFF, cases[c].name, + cases[c].values, repetitions)); + } + } + (void)data_type; +} + +} // namespace + +int main(int argc, char** argv) { + uint32_t repetitions = 3; + if (argc > 1) { + repetitions = static_cast(std::max(1, atoi(argv[1]))); + } + const std::string codec_filter = argc > 2 ? argv[2] : ""; + const std::string type_filter = argc > 3 ? argv[3] : ""; + const std::string pattern_filter = argc > 4 ? argv[4] : ""; + g_ts2diff_mpn = argc > 5 ? atoi(argv[5]) : -1; + std::cout << "ALP vs GORILLA vs TS_2DIFF codec benchmark, best of " + << repetitions << " runs"; + if (g_ts2diff_mpn >= 0) { + std::cout << ", TS_2DIFF mpn=" << g_ts2diff_mpn; + } else { + std::cout << ", TS_2DIFF mpn=2(float)/3(double)"; + } + std::cout << std::endl; + if (type_filter.empty() || type_filter == "FLOAT") { + if (g_ts2diff_mpn < 0) { + g_ts2diff_mpn = 2; + } + RunAll("FLOAT", repetitions, codec_filter, pattern_filter); + } + if (type_filter.empty() || type_filter == "DOUBLE") { + if (argc <= 5) { + g_ts2diff_mpn = 3; + } + RunAll("DOUBLE", repetitions, codec_filter, pattern_filter); + } + return 0; +} diff --git a/cpp/docs/alp-codec-design.md b/cpp/docs/alp-codec-design.md new file mode 100644 index 000000000..2e900db85 --- /dev/null +++ b/cpp/docs/alp-codec-design.md @@ -0,0 +1,158 @@ + + +# ALP Float/Double Codec Design + +Status: draft for the C++-first implementation on `colin/alp-codec`. + +## Scope + +- This work is C++-only. The Java implementation is intentionally untouched. +- ALP is an opt-in encoding for `FLOAT` and `DOUBLE` only. +- The existing Gorilla default is preserved. ALP is enabled through an + experimental build option until the format is accepted by the community. +- The format is language-neutral so a Java decoder can be added later without + changing encoded bytes. +- Python will reuse the C++ implementation through the C wrapper. + +## Relationship to TS_2DIFF + +ALP must not reuse the TS_2DIFF float-adaptation path. In particular it must not +use TS_2DIFF's page precision (`maxPointNumber`), page overflow bitmap, +`page_blocks_` buffering, or `PageMeta` parsing. + +ALP has its own block metadata because the algorithm needs it: exponent/factor +indices, frame-of-reference base, bit width, and exceptions. This metadata is +block-local and is part of the ALP payload. A TsFile page is only a container; +the ALP block boundary is a flush boundary. + +## Wire Format (version 1) + +All integers are little-endian. `value_size` is 4 for `FLOAT` and 8 for +`DOUBLE`. + +### Page header + +| field | size | description | +|---|---:|---| +| version | 1 | format version, currently 1 | +| data_type | 1 | 0 = float, 1 = double | +| reserved | 2 | zero | +| value_count | 4 | total values in the page | +| block_count | 4 | number of ALP blocks | +| reserved | 4 | zero | + +### Block header + +| field | size | description | +|---|---:|---| +| scheme | 1 | 0 = ALP, 1 = PLAIN fallback | +| factor | 1 | ALP factor index | +| exponent | 1 | ALP exponent index | +| flags | 1 | bit 0 = exceptions present | +| value_count | 4 | values in this block | +| exception_count | 4 | number of exceptions | +| body_bytes | 4 | padded body size in bytes | +| for_base | 8 | signed frame-of-reference base, int32 for float and int64 for double | +| bit_width | 1 | 0..32 for float, 0..64 for double | +| reserved | 3 | zero | + +The header is followed by: + +- for `scheme == ALP` with exceptions: an exception bitmap of + `ceil(value_count / 8)` bytes, followed by `exception_count` raw + `value_size` values in increasing position order; +- for `scheme == ALP`: the bit-packed body; +- for `scheme == PLAIN`: `value_count * value_size` raw values. + +The packed body is padded with zero bytes to a multiple of 32 bytes after +adding 16 bytes of guard space, so SIMD/word loads may safely read past the +logical end of the stream. `body_bytes` records the padded size. + +### Bit packing + +The body stores adjusted integers using an LSB-first word stream. For float +each word is 32 bits; for double each word is 64 bits. Value `i` occupies +`bit_width` bits starting at bit `i * bit_width`. A zero bit width means every +non-exception value equals `for_base`. + +This layout is deliberately different from TsFile's existing MSB-first pack8 +layout. It is defined by the scalar reference implementation, and SIMD kernels +must produce and consume the same bytes. + +## Encoding + +1. Buffer values per page. +2. Split the buffered values into blocks of at most 1024 values. +3. For each block, sample values and choose the best `(factor, exponent)` + combination using the estimated compressed size. +4. Encode each value with the chosen pair. A value is an exception when it is + non-finite, `-0.0`, outside the safe integer range, or when the reference + decode does not reproduce the original bit pattern. +5. Compute the frame-of-reference base and bit width over non-exception + integers. +6. Estimate whether ALP is smaller than PLAIN. If not, emit the block as + `PLAIN`. +7. Emit the block header, exception data, and packed body. + +The encoder's exception check must use the same decode routine as the decoder. +This is required for correctness: a value classified as a non-exception must +be bit-exactly reproducible by the decoder. + +## Decoding + +1. Parse and validate the page header. +2. For each block, parse and validate the block header. +3. For `PLAIN`, copy raw values. +4. For `ALP`, unpack adjusted integers, add `for_base`, apply the reference + decode, and overwrite exception positions from the exception array. +5. Validate that the decoded value count matches the page header. +6. Implement `read_batch_float` and `read_batch_double` so the SIMD path can + work on whole blocks rather than one value at a time. + +## SIMD Plan + +SIMDe is already a TsFile C++ dependency. The ALP kernels will use the existing +`ENABLE_SIMD` switch and `TsFile::SIMDe` target. + +- Write the kernels with `simde__m256i` and `simde_mm256_*`, matching the + existing encoding code style. +- Compile an AVX2 kernel translation unit with `-mavx2` or `/arch:AVX2` on + x86. Do not apply AVX2 flags to baseline TsFile translation units. +- On AArch64 the same SIMDe source maps to NEON without the AVX2 flag. +- Add a scalar fallback used when SIMD is disabled or unavailable. +- Dispatch once per block through a function table. Do not branch per value. +- Add an optional AVX-512 translation unit later. The wire format must not + depend on the availability of AVX-512. + +## Testing + +- Scalar round-trip tests for float and double. +- Boundary values: NaN, +Inf, -Inf, -0.0, subnormals, smallest/largest + magnitudes, exact zero, and values forcing exceptions. +- Constant blocks, single-value blocks, partial final blocks, and multi-block + pages. +- Scalar and SIMD must encode to identical bytes and decode to identical + values. +- Corrupt/truncated payload tests must return errors without out-of-bounds + access. +- Benchmarks compare ALP against Gorilla for decimal, normalized, smooth, + repeated, and high-entropy data. diff --git a/cpp/docs/alp-profiling-analysis.md b/cpp/docs/alp-profiling-analysis.md new file mode 100644 index 000000000..bc6b708c4 --- /dev/null +++ b/cpp/docs/alp-profiling-analysis.md @@ -0,0 +1,142 @@ + + +# ALP Profiling Analysis + +This note explains why the ALP codec still loses to Gorilla in some cases +even though the value kernels are SIMD-accelerated. The profiles were collected +with gperftools on the FLOAT decimal benchmark case. + +## Profile Setup + +```bash +CPUPROFILE=/tmp/alp_float.prof \ +CPUPROFILE_FREQUENCY=1000 \ +DYLD_INSERT_LIBRARIES=/opt/homebrew/lib/libprofiler.dylib \ +./cpp/build/alp/alp_vs_gorilla_benchmark 500 ALP FLOAT decimal + +/opt/homebrew/bin/pprof --svg \ + ./cpp/build/alp/alp_vs_gorilla_benchmark \ + /tmp/alp_float.prof > /tmp/alp_float.svg +``` + +The same command was repeated with `GORILLA` for comparison. + +## ALP FLOAT Decimal + +![ALP FLOAT decimal pprof](images/alp-float-decimal-encode-pprof.png) + +Flat profile: + +| Function | Flat | Cumulative | +| --- | ---: | ---: | +| `storage::AlpEncoderBase::flush` | 73.5% | 86.3% | +| `storage::alp::AlpAppendU32` | 7.8% | 7.8% | +| `storage::AlpDecoderBase::ensure_initialized` | 3.0% | 11.3% | +| `storage::alp::AlpDecodeValues` | 2.7% | 2.7% | +| `storage::alp::AlpUnpackBits` | 2.3% | 2.3% | +| `storage::AlpDecoderBase::read_batch_float` | 0.5% | 11.8% | + +The encode path dominates. Decode is already a small fraction of total time. + +## Gorilla FLOAT Decimal + +![Gorilla FLOAT decimal pprof](images/gorilla-float-decimal-pprof.png) + +Flat profile: + +| Function | Flat | Cumulative | +| --- | ---: | ---: | +| `storage::GorillaEncoder::compress_value` | 29.0% | 72.3% | +| `common::ByteStream::write_buf` | 27.2% | 43.3% | +| `storage::FloatGorillaDecoder::read_batch_float` | 22.0% | 22.0% | +| `storage::FloatGorillaEncoder::encode` | 2.2% | 74.5% | +| `storage::Encoder::encode_batch` | 1.9% | 76.4% | + +Gorilla spends most of its time in its scalar bit writer and in the ByteStream +write path. Its encode is not SIMD either; it is simply a very efficient +scalar bit stream for the FLOAT decimal pattern. + +## Per-Block Breakdown + +A separate microbenchmark measured one 1024-value block with the same decimal +pattern: + +| Stage | FLOAT (ms/block) | DOUBLE (ms/block) | +| --- | ---: | ---: | +| `AlpChooseFactorExponent` (two-stage) | 0.002 | 0.003 | +| `AlpEncodeValues` (SIMD) | 0.001 | 0.001 | +| `AlpPackBits` (word-based) | 0.002 | 0.002 | +| `AlpEncodePage` total | 0.005 | 0.007 | + +After replacing the per-bit packing loop with word-based packing and adding a +two-stage exponent/factor search, the block encode cost dropped to 0.005 ms +(FLOAT) and 0.007 ms (DOUBLE). The two-stage search is still the largest single +component (about 40% FLOAT and 43% DOUBLE), followed by word-based packing. +The SIMD value kernel is about 14-20% of the block encode time. + +## Why ALP Still Loses + +1. **The e/f search is still scalar.** + The two-stage search cut the cost substantially, but it remains the largest + single encode component. Stage 1 evaluates every candidate on an 8-value + sample; stage 2 runs the exact round-trip check on the shortlist. Further + gains require SIMD candidate evaluation or a cheaper estimator. + +2. **The SIMD value kernel is still a minority of encode time.** + `AlpEncodeValues` is about 14-20% of the block encode time. End-to-end + encode speed also depends on the search and packing decisions. + +3. **Non-decimal data falls back to PLAIN.** + Smooth/random/constant data often has too many exceptions for ALP to win. + The encoder then emits a PLAIN block. In that case SIMD cannot help because + the data is not ALP-encoded at all; Gorilla is usually the better codec for + those columns. + +4. **Gorilla has a dedicated constant path.** + Gorilla compresses constant blocks extremely well and decodes them very + quickly. ALP intentionally does not special-case that pattern, so Gorilla + wins the constant rows of the benchmark. + +5. **The benchmark is on Apple Silicon.** + SIMDe maps the AVX2-style kernels to NEON here. On x86, a native `-mavx2` + translation unit and runtime dispatch are still required to get the AVX2 + peak; without those flags SIMDe may emulate AVX2 with SSE2. + +## Next Optimizations + +1. **SIMD or cheaper exponent/factor selection.** + Evaluate the candidate pairs on the 32-value sample with SIMD, or reduce the + candidate set with a two-stage/early-exit estimate. This is now the largest + remaining encode cost for both FLOAT and DOUBLE. + +2. **Remove per-block allocations and copies.** + Reuse scratch vectors for `encoded`, `adjusted`, `bitmap`, `body`, and the + exception list instead of constructing them per block. Decode directly into + the caller's output buffer rather than an intermediate `values_` vector. + +3. **Adaptive codec selection.** + Choose Gorilla for constant/high-entropy columns and ALP for decimal-like + columns. This avoids the PLAIN fallback cases where Gorilla is better. + +4. **Native x86 dispatch.** + Add an AVX2 object library and CPU detection so x86 users get native AVX2 + without `-march=native`, plus an optional AVX-512 path for DOUBLE. diff --git a/cpp/docs/alp-vs-gorilla-benchmark.md b/cpp/docs/alp-vs-gorilla-benchmark.md new file mode 100644 index 000000000..b0c16a036 --- /dev/null +++ b/cpp/docs/alp-vs-gorilla-benchmark.md @@ -0,0 +1,162 @@ + + +# ALP vs Gorilla vs TS_2DIFF Codec Benchmark + +This document records the C++ codec-level comparison between ALP, Gorilla, +and TS_2DIFF on the `colin/alp-codec` branch. It is not an end-to-end TsFile +I/O benchmark; it measures the codec path through the same +`EncoderFactory` / `DecoderFactory` interfaces used by TsFile. + +## Environment + +| Item | Value | +| --- | --- | +| Machine | Apple Silicon (`arm64`) | +| Compiler | AppleClang 17.0.0 | +| Build | Release, `-O3` | +| SIMDe | `0.8.4-rc3` | +| Values per case | 1,048,576 | +| Repetitions | best of 5 | +| Compression | `UNCOMPRESSED` | + +On this machine SIMDe maps the AVX2-style kernels to NEON. The x86 native +AVX2 multi-version object library and runtime dispatch are a follow-up step; +the benchmark therefore reflects the portable SIMDe path, not the final +AVX2/AVX-512 peak. + +## Data Patterns + +| Pattern | Description | +| --- | --- | +| decimal | repeated fixed-scale decimals (`0.01` float steps, `0.001` double steps) | +| smooth | smooth `sin`/`cos` values | +| random | pseudo-random raw IEEE-754 bit patterns, finite values | +| constant | one repeated constant value | + +## Results + +### FLOAT + +| Codec | Pattern | Encoded bytes | Ratio | Encode MB/s | Decode MB/s | Correct | +| --- | --- | ---: | ---: | ---: | ---: | --- | +| ALP | decimal | 2,365,932 | 1.773 | 539.3 | 2646.9 | yes | +| Gorilla | decimal | 2,878,091 | 1.457 | 342.0 | 1389.3 | yes | +| ALP | smooth | 4,124,940 | 1.017 | 528.5 | 2399.5 | yes | +| Gorilla | smooth | 4,453,445 | 0.942 | 248.1 | 1384.7 | yes | +| ALP | random | 4,222,992 | 0.993 | 506.2 | 2329.7 | yes | +| Gorilla | random | 4,456,456 | 0.941 | 239.6 | 1378.9 | yes | +| ALP | constant | 28,688 | 146.204 | 1335.4 | 7393.7 | yes | +| Gorilla | constant | 131,082 | 31.998 | 1900.8 | 69214.9 | yes | + +### DOUBLE + +| Codec | Pattern | Encoded bytes | Ratio | Encode MB/s | Decode MB/s | Correct | +| --- | --- | ---: | ---: | ---: | ---: | --- | +| ALP | decimal | 1,426,248 | 5.882 | 1055.3 | 6470.7 | yes | +| Gorilla | decimal | 8,650,641 | 0.970 | 276.7 | 2546.8 | yes | +| ALP | smooth | 8,348,616 | 1.005 | 780.8 | 2228.6 | yes | +| Gorilla | smooth | 8,629,855 | 0.972 | 276.8 | 2600.1 | yes | +| ALP | random | 8,417,296 | 0.997 | 753.6 | 2118.6 | yes | +| Gorilla | random | 8,650,764 | 0.970 | 277.1 | 2644.3 | yes | +| ALP | constant | 8,417,296 | 0.997 | 861.9 | 2094.3 | yes | +| Gorilla | constant | 131,090 | 63.991 | 3824.2 | 76646.7 | yes | + +## Three-Way Decimal Comparison (ALP vs Gorilla vs TS_2DIFF) + +TS_2DIFF's FLOAT/DOUBLE path is the direct competitor for fixed-precision +decimal data. It scales the value by `10^maxPointNumber`, delta-encodes the +resulting integers, and stores page-level overflow/rounding metadata. The +benchmark now includes it as a reference row. + +### FLOAT decimal + +| Codec | Ratio | Encode MB/s | Decode MB/s | Lossless | +| --- | ---: | ---: | ---: | --- | +| ALP | 1.773 | 584.3 | 2749.7 | yes | +| Gorilla | 1.457 | 353.2 | 1333.2 | yes | +| TS_2DIFF (`mpn=2`) | **31.587** | **1419.7** | **1702.7** | **no** | + +### DOUBLE decimal + +| Codec | Ratio | Encode MB/s | Decode MB/s | Lossless | +| --- | ---: | ---: | ---: | --- | +| ALP | 5.882 | 1071.7 | **6398.5** | yes | +| Gorilla | 0.970 | 284.5 | 2633.7 | yes | +| TS_2DIFF (`mpn=3`) | **42.303** | **2839.6** | 3051.3 | **no** | + +The TS_2DIFF `Lossless = no` result is not a benchmark harness problem. A +standalone probe found: + +- FLOAT: a tiny residual value (`2.08e-17`) was reconstructed as `0.0`. +- DOUBLE: several fixed-3-decimal values were reconstructed with a 1-ULP + difference. + +TS_2DIFF is therefore excellent when the application accepts the scaled-value +semantics and does not require bit-exact reconstruction. ALP is the safer +choice when the original IEEE-754 bit pattern must be preserved exactly, +because its encoder performs a round-trip check and stores an exception when +the integer conversion cannot reproduce the input. + +The three-way trade-off is: + +- **TS_2DIFF**: best size on fixed-precision decimal data, fast encode/decode, + but the current C++ implementation is not bit-exact for all FLOAT/DOUBLE + values. +- **ALP**: bit-exact by construction, strong encode/decode throughput, better + ratio than Gorilla on decimal data, but lower ratio than TS_2DIFF because it + stores exceptions that TS_2DIFF approximates. +- **Gorilla**: best for constant and high-entropy shapes; not a decimal + specialist. + +## Interpretation + +- ALP is designed for decimal-like floating-point columns. On the decimal + cases it improves both compression and decode throughput: + - FLOAT decimal: 1.22x better ratio, 1.58x faster encode, and 1.91x faster + decode than Gorilla. + - DOUBLE decimal: 6.06x better ratio, 3.81x faster encode, and 2.54x faster + decode than Gorilla. +- ALP now wins both encode and decode on the decimal cases. Word-based packing + removed the largest FLOAT encode cost, and the two-stage exponent/factor + search cut the remaining scalar search cost substantially. +- ALP falls back to PLAIN for non-decimal smooth/random data. In those cases + Gorilla is the better choice, so ALP must remain opt-in. +- Gorilla is extremely fast on constant blocks. ALP does not try to beat that + case; the adaptive selection should leave constant columns on Gorilla. +- The benchmark includes the SIMDe value kernels, SIMD 4-lane unpack, and + word-based packing. Native x86 AVX2 multi-versioning and AVX-512 kernels are + expected to increase ALP throughput further on x86. + +## Reproduce + +```bash +cmake -S cpp -B cpp/build/alp \ + -DCMAKE_BUILD_TYPE=Release \ + -DBUILD_BENCHMARK=ON \ + -DBUILD_TEST=ON \ + -DENABLE_SIMD=ON \ + -DTSFILE_DEPENDENCY_SOURCE=BUNDLED + +cmake --build cpp/build/alp --target alp_vs_gorilla_benchmark -j + +./cpp/build/alp/alp_vs_gorilla_benchmark 5 +``` diff --git a/cpp/docs/images/alp-float-decimal-encode-pprof.png b/cpp/docs/images/alp-float-decimal-encode-pprof.png new file mode 100644 index 000000000..6a0aaeb15 Binary files /dev/null and b/cpp/docs/images/alp-float-decimal-encode-pprof.png differ diff --git a/cpp/docs/images/gorilla-float-decimal-pprof.png b/cpp/docs/images/gorilla-float-decimal-pprof.png new file mode 100644 index 000000000..19a062030 Binary files /dev/null and b/cpp/docs/images/gorilla-float-decimal-pprof.png differ diff --git a/cpp/src/common/db_common.h b/cpp/src/common/db_common.h index 0c2843b45..170b73bc3 100644 --- a/cpp/src/common/db_common.h +++ b/cpp/src/common/db_common.h @@ -77,6 +77,7 @@ enum TSEncoding : uint8_t { SPRINTZ = 12, RLBE = 13, CAMEL = 14, + ALP = 15, INVALID_ENCODING = 255 }; @@ -101,7 +102,7 @@ enum CompressionType : uint8_t { }; extern TSFILE_API const char* s_data_type_names[12]; -extern TSFILE_API const char* s_encoding_names[15]; +extern TSFILE_API const char* s_encoding_names[16]; extern TSFILE_API const char* s_compression_names[10]; } // namespace common @@ -159,7 +160,7 @@ FORCE_INLINE bool parse_data_type_name(const std::string& s, TSDataType& out) { } FORCE_INLINE const char* get_encoding_name(TSEncoding encoding) { - ASSERT(encoding >= PLAIN && encoding <= CAMEL); + ASSERT(encoding >= PLAIN && encoding <= ALP); return s_encoding_names[encoding]; } diff --git a/cpp/src/common/global.cc b/cpp/src/common/global.cc index bfc713037..a300ece89 100644 --- a/cpp/src/common/global.cc +++ b/cpp/src/common/global.cc @@ -147,10 +147,10 @@ const char* s_data_type_names[12] = {"BOOLEAN", "INT32", "INT64", "FLOAT", "DOUBLE", "TEXT", "VECTOR", "UNKNOWN", "TIMESTAMP", "DATE", "BLOB", "STRING"}; -const char* s_encoding_names[15] = { - "PLAIN", "DICTIONARY", "RLE", "DIFF", "TS_2DIFF", - "BITMAP", "GORILLA_V1", "REGULAR", "GORILLA", "ZIGZAG", - "FREQ", "CHIMP", "SPRINTZ", "RLBE", "CAMEL"}; +const char* s_encoding_names[16] = { + "PLAIN", "DICTIONARY", "RLE", "DIFF", "TS_2DIFF", "BITMAP", + "GORILLA_V1", "REGULAR", "GORILLA", "ZIGZAG", "FREQ", "CHIMP", + "SPRINTZ", "RLBE", "CAMEL", "ALP"}; const char* s_compression_names[10] = { "UNCOMPRESSED", "SNAPPY", "GZIP", "LZO", "SDT", diff --git a/cpp/src/common/global.h b/cpp/src/common/global.h index f1f6df19c..9a39ea7ce 100644 --- a/cpp/src/common/global.h +++ b/cpp/src/common/global.h @@ -48,7 +48,7 @@ FORCE_INLINE int set_global_time_data_type(uint8_t data_type) { } FORCE_INLINE int set_global_time_encoding(uint8_t encoding) { - ASSERT(encoding >= PLAIN && encoding <= CAMEL); + ASSERT(encoding >= PLAIN && encoding <= ALP); if (encoding != TS_2DIFF && encoding != PLAIN) { return E_NOT_SUPPORT; } @@ -101,7 +101,8 @@ FORCE_INLINE int set_datatype_encoding(uint8_t data_type, uint8_t encoding) { case FLOAT: if (encoding_type != PLAIN && encoding_type != TS_2DIFF && encoding_type != GORILLA && encoding_type != SPRINTZ && - encoding_type != CHIMP && encoding_type != RLBE) { + encoding_type != CHIMP && encoding_type != RLBE && + encoding_type != ALP) { return E_NOT_SUPPORT; } g_config_value_.float_encoding_type_ = encoding_type; @@ -111,7 +112,7 @@ FORCE_INLINE int set_datatype_encoding(uint8_t data_type, uint8_t encoding) { if (encoding_type != PLAIN && encoding_type != TS_2DIFF && encoding_type != GORILLA && encoding_type != SPRINTZ && encoding_type != CHIMP && encoding_type != RLBE && - encoding_type != CAMEL) { + encoding_type != CAMEL && encoding_type != ALP) { return E_NOT_SUPPORT; } g_config_value_.double_encoding_type_ = encoding_type; diff --git a/cpp/src/cwrapper/tsfile_cwrapper.h b/cpp/src/cwrapper/tsfile_cwrapper.h index 4656c5545..48568da52 100644 --- a/cpp/src/cwrapper/tsfile_cwrapper.h +++ b/cpp/src/cwrapper/tsfile_cwrapper.h @@ -60,6 +60,7 @@ typedef enum { TS_ENCODING_SPRINTZ = 12, TS_ENCODING_RLBE = 13, TS_ENCODING_CAMEL = 14, + TS_ENCODING_ALP = 15, TS_ENCODING_INVALID = 255 } TSEncoding; diff --git a/cpp/src/encoding/CMakeLists.txt b/cpp/src/encoding/CMakeLists.txt index 45c1f0bdd..60e099b6e 100644 --- a/cpp/src/encoding/CMakeLists.txt +++ b/cpp/src/encoding/CMakeLists.txt @@ -22,5 +22,5 @@ message("Running in src/encoding directory") add_custom_target( encoding_obj ALL ) -file(GLOB HEADERS "${CMAKE_CURRENT_SOURCE_DIR}/*.h") +file(GLOB_RECURSE HEADERS "${CMAKE_CURRENT_SOURCE_DIR}/*.h") copy_to_dir(${HEADERS} "encoding_obj") \ No newline at end of file diff --git a/cpp/src/encoding/alp/alp_constants.h b/cpp/src/encoding/alp/alp_constants.h new file mode 100644 index 000000000..ffb758ffd --- /dev/null +++ b/cpp/src/encoding/alp/alp_constants.h @@ -0,0 +1,132 @@ +/* + * 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. + */ + +#ifndef ENCODING_ALP_CONSTANTS_H +#define ENCODING_ALP_CONSTANTS_H + +#include + +namespace storage { +namespace alp { + +enum AlpDataType : uint8_t { ALP_FLOAT = 0, ALP_DOUBLE = 1 }; + +enum AlpScheme : uint8_t { + ALP_SCHEME_ALP = 0, + ALP_SCHEME_PLAIN = 1, + ALP_SCHEME_ALP_RD = 2 // reserved for a later format revision +}; + +static const uint8_t ALP_FORMAT_VERSION = 1; +static const uint32_t ALP_BLOCK_SIZE = 1024; +static const uint32_t ALP_SAMPLE_SIZE = 32; +static const uint32_t ALP_BODY_ALIGNMENT = 32; +static const uint8_t ALP_FLOAT_MAX_EXPONENT = 10; +static const uint8_t ALP_DOUBLE_MAX_EXPONENT = 18; + +static const float ALP_FLOAT_MAGIC = 12582912.0f; // 1.5 * 2^23 +static const double ALP_DOUBLE_MAGIC = 6755399441055744.0; // 1.5 * 2^52 + +// Values outside these bounds cannot be cast to the corresponding integer type +// without undefined behaviour. The bounds intentionally leave a small margin +// so the magic-number rounding itself cannot overflow. +static const float ALP_FLOAT_ENCODING_LOWER_LIMIT = -2147483520.0f; +static const float ALP_FLOAT_ENCODING_UPPER_LIMIT = 2147483520.0f; +static const double ALP_DOUBLE_ENCODING_LOWER_LIMIT = -9223372036854774784.0; +static const double ALP_DOUBLE_ENCODING_UPPER_LIMIT = 9223372036854774784.0; + +static const float ALP_FLOAT_FRAC[ALP_FLOAT_MAX_EXPONENT + 1] = { + 1.0f, 0.1f, 0.01f, 0.001f, 0.0001f, 0.00001f, + 0.000001f, 0.0000001f, 0.00000001f, 0.000000001f, 0.0000000001f}; + +static const float ALP_FLOAT_EXP[ALP_FLOAT_MAX_EXPONENT + 1] = { + 1.0f, 10.0f, 100.0f, 1000.0f, + 10000.0f, 100000.0f, 1000000.0f, 10000000.0f, + 100000000.0f, 1000000000.0f, 10000000000.0f}; + +static const double ALP_DOUBLE_FRAC[ALP_DOUBLE_MAX_EXPONENT + 1] = { + 1.0, + 0.1, + 0.01, + 0.001, + 0.0001, + 0.00001, + 0.000001, + 0.0000001, + 0.00000001, + 0.000000001, + 0.0000000001, + 0.00000000001, + 0.000000000001, + 0.0000000000001, + 0.00000000000001, + 0.000000000000001, + 0.0000000000000001, + 0.00000000000000001, + 0.000000000000000001}; + +static const double ALP_DOUBLE_EXP[ALP_DOUBLE_MAX_EXPONENT + 1] = { + 1.0, + 10.0, + 100.0, + 1000.0, + 10000.0, + 100000.0, + 1000000.0, + 10000000.0, + 100000000.0, + 1000000000.0, + 10000000000.0, + 100000000000.0, + 1000000000000.0, + 10000000000000.0, + 100000000000000.0, + 1000000000000000.0, + 10000000000000000.0, + 100000000000000000.0, + 1000000000000000000.0}; + +// Exact integer powers of ten used by the decoder. Up to 10^18 they are +// exactly representable as double; the float values are exactly representable +// up to 10^10. +static const int64_t ALP_FACT[ALP_DOUBLE_MAX_EXPONENT + 1] = { + 1LL, + 10LL, + 100LL, + 1000LL, + 10000LL, + 100000LL, + 1000000LL, + 10000000LL, + 100000000LL, + 1000000000LL, + 10000000000LL, + 100000000000LL, + 1000000000000LL, + 10000000000000LL, + 100000000000000LL, + 1000000000000000LL, + 10000000000000000LL, + 100000000000000000LL, + 1000000000000000000LL}; + +} // namespace alp +} // namespace storage + +#endif // ENCODING_ALP_CONSTANTS_H diff --git a/cpp/src/encoding/alp/alp_decoder.h b/cpp/src/encoding/alp/alp_decoder.h new file mode 100644 index 000000000..7fd34d7f4 --- /dev/null +++ b/cpp/src/encoding/alp/alp_decoder.h @@ -0,0 +1,246 @@ +/* + * 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. + */ + +#ifndef ENCODING_ALP_DECODER_H +#define ENCODING_ALP_DECODER_H + +#include + +#include +#include + +#include "alp_scalar.h" +#include "common/db_common.h" +#include "encoding/decoder.h" +#include "utils/errno_define.h" + +namespace storage { + +template +class AlpDecoderBase : public Decoder { + public: + AlpDecoderBase() : initialized_(false), read_index_(0) {} + ~AlpDecoderBase() override { destroy(); } + + bool owns_resources() const override { return true; } + + void destroy() override { + std::vector().swap(values_); + std::vector().swap(page_buffer_); + initialized_ = false; + read_index_ = 0; + } + + void reset() override { + values_.clear(); + page_buffer_.clear(); + initialized_ = false; + read_index_ = 0; + } + + bool has_remaining(const common::ByteStream& in) override { + if (initialized_) { + return read_index_ < values_.size(); + } + return in.has_remaining(); + } + + int read_boolean(bool&, common::ByteStream&) override { + return common::E_TYPE_NOT_MATCH; + } + int read_int32(int32_t&, common::ByteStream&) override { + return common::E_TYPE_NOT_MATCH; + } + int read_int64(int64_t&, common::ByteStream&) override { + return common::E_TYPE_NOT_MATCH; + } + int read_String(common::String&, common::PageArena&, + common::ByteStream&) override { + return common::E_TYPE_NOT_MATCH; + } + + int read_float(float& ret_value, common::ByteStream& in) override { + if (sizeof(T) != sizeof(float)) { + return common::E_TYPE_NOT_MATCH; + } + const int ret = ensure_initialized(in); + if (ret != common::E_OK) { + return ret; + } + if (read_index_ >= values_.size()) { + return common::E_NO_MORE_DATA; + } + ret_value = static_cast(values_[read_index_++]); + return common::E_OK; + } + + int read_double(double& ret_value, common::ByteStream& in) override { + if (sizeof(T) != sizeof(double)) { + return common::E_TYPE_NOT_MATCH; + } + const int ret = ensure_initialized(in); + if (ret != common::E_OK) { + return ret; + } + if (read_index_ >= values_.size()) { + return common::E_NO_MORE_DATA; + } + ret_value = static_cast(values_[read_index_++]); + return common::E_OK; + } + + int read_batch_float(float* out, int capacity, int& actual, + common::ByteStream& in) override { + if (sizeof(T) != sizeof(float)) { + return common::E_TYPE_NOT_MATCH; + } + return read_batch_impl(out, capacity, actual, in); + } + + int read_batch_double(double* out, int capacity, int& actual, + common::ByteStream& in) override { + if (sizeof(T) != sizeof(double)) { + return common::E_TYPE_NOT_MATCH; + } + return read_batch_impl(out, capacity, actual, in); + } + + int skip_float(int count, int& skipped, common::ByteStream& in) override { + if (sizeof(T) != sizeof(float)) { + return common::E_TYPE_NOT_MATCH; + } + return skip_impl(count, skipped, in); + } + + int skip_double(int count, int& skipped, common::ByteStream& in) override { + if (sizeof(T) != sizeof(double)) { + return common::E_TYPE_NOT_MATCH; + } + return skip_impl(count, skipped, in); + } + + private: + int ensure_initialized(common::ByteStream& in) { + if (initialized_) { + return common::E_OK; + } + initialized_ = true; + const uint64_t remaining = in.remaining_size(); + if (remaining == 0) { + return common::E_OK; + } + if (remaining > static_cast(UINT32_MAX)) { + return common::E_INVALID_ARG; + } + const uint8_t* data = NULL; + if (in.is_wrapped()) { + // The TsFile reader wraps the page value buffer directly. Decode + // from it without copying the whole page into page_buffer_. + data = reinterpret_cast(in.get_wrapped_buf() + + in.read_pos()); + in.wrapped_buf_advance_read_pos(remaining); + } else { + page_buffer_.resize(static_cast(remaining)); + uint32_t read_len = 0; + const int ret = in.read_buf( + &page_buffer_[0], static_cast(remaining), read_len); + if (ret != common::E_OK && ret != common::E_PARTIAL_READ) { + return ret; + } + if (read_len != remaining) { + return common::E_PARTIAL_READ; + } + data = &page_buffer_[0]; + } + const alp::AlpStatus status = + alp::AlpDecodePage(data, static_cast(remaining), values_); + if (status != alp::ALP_OK) { + values_.clear(); + return common::E_DECODE_ERR; + } + return common::E_OK; + } + + int read_batch_impl(float* out, int capacity, int& actual, + common::ByteStream& in) { + return read_batch_generic(out, capacity, actual, in); + } + + int read_batch_impl(double* out, int capacity, int& actual, + common::ByteStream& in) { + return read_batch_generic(out, capacity, actual, in); + } + + template + int read_batch_generic(Out* out, int capacity, int& actual, + common::ByteStream& in) { + actual = 0; + if (capacity <= 0 || out == NULL) { + return common::E_OK; + } + const int ret = ensure_initialized(in); + if (ret != common::E_OK) { + return ret; + } + const size_t remaining = values_.size() - read_index_; + const size_t count = std::min(static_cast(capacity), remaining); + for (size_t i = 0; i < count; ++i) { + out[i] = static_cast(values_[read_index_ + i]); + } + read_index_ += count; + actual = static_cast(count); + return common::E_OK; + } + + int skip_impl(int count, int& skipped, common::ByteStream& in) { + skipped = 0; + if (count <= 0) { + return common::E_OK; + } + const int ret = ensure_initialized(in); + if (ret != common::E_OK) { + return ret; + } + const size_t remaining = values_.size() - read_index_; + const size_t skipped_count = + std::min(static_cast(count), remaining); + read_index_ += skipped_count; + skipped = static_cast(skipped_count); + return common::E_OK; + } + + std::vector values_; + std::vector page_buffer_; + bool initialized_; + size_t read_index_; +}; + +class FloatAlpDecoder : public AlpDecoderBase { + public: + FloatAlpDecoder() {} +}; + +class DoubleAlpDecoder : public AlpDecoderBase { + public: + DoubleAlpDecoder() {} +}; + +} // namespace storage + +#endif // ENCODING_ALP_DECODER_H diff --git a/cpp/src/encoding/alp/alp_encoder.h b/cpp/src/encoding/alp/alp_encoder.h new file mode 100644 index 000000000..21a8120f1 --- /dev/null +++ b/cpp/src/encoding/alp/alp_encoder.h @@ -0,0 +1,153 @@ +/* + * 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. + */ + +#ifndef ENCODING_ALP_ENCODER_H +#define ENCODING_ALP_ENCODER_H + +#include + +#include + +#include "alp_scalar.h" +#include "common/db_common.h" +#include "encoding/encoder.h" +#include "utils/errno_define.h" + +namespace storage { + +template +class AlpEncoderBase : public Encoder { + public: + AlpEncoderBase() {} + ~AlpEncoderBase() override { destroy(); } + + void destroy() override { + std::vector().swap(values_); + std::vector().swap(page_buffer_); + } + + void reset() override { + values_.clear(); + page_buffer_.clear(); + } + + int encode(bool, common::ByteStream&) override { + return common::E_TYPE_NOT_MATCH; + } + int encode(int32_t, common::ByteStream&) override { + return common::E_TYPE_NOT_MATCH; + } + int encode(int64_t, common::ByteStream&) override { + return common::E_TYPE_NOT_MATCH; + } + int encode(common::String, common::ByteStream&) override { + return common::E_TYPE_NOT_MATCH; + } + + int encode(float value, common::ByteStream&) override { + if (sizeof(T) != sizeof(float)) { + return common::E_TYPE_NOT_MATCH; + } + return encode_value(static_cast(value)); + } + + int encode(double value, common::ByteStream&) override { + if (sizeof(T) != sizeof(double)) { + return common::E_TYPE_NOT_MATCH; + } + return encode_value(static_cast(value)); + } + + int encode_batch(const float* values, uint32_t count, + common::ByteStream&) override { + if (sizeof(T) != sizeof(float) || (values == NULL && count != 0)) { + return common::E_TYPE_NOT_MATCH; + } + return encode_batch_values(reinterpret_cast(values), count); + } + + int encode_batch(const double* values, uint32_t count, + common::ByteStream&) override { + if (sizeof(T) != sizeof(double) || (values == NULL && count != 0)) { + return common::E_TYPE_NOT_MATCH; + } + return encode_batch_values(reinterpret_cast(values), count); + } + + int flush(common::ByteStream& out_stream) override { + if (values_.empty()) { + return common::E_OK; + } + page_buffer_.reserve(values_.size() * sizeof(T) + + (values_.size() / alp::ALP_BLOCK_SIZE + 1) * + alp::ALP_BLOCK_HEADER_SIZE); + const alp::AlpStatus status = alp::AlpEncodePage( + values_.empty() ? NULL : &values_[0], + static_cast(values_.size()), page_buffer_); + values_.clear(); + if (status != alp::ALP_OK) { + return common::E_ENCODE_ERR; + } + if (page_buffer_.empty()) { + return common::E_OK; + } + const int ret = + out_stream.write_buf(page_buffer_.empty() ? NULL : &page_buffer_[0], + static_cast(page_buffer_.size())); + page_buffer_.clear(); + return ret; + } + + int get_max_byte_size() override { + return static_cast(values_.capacity() * sizeof(T) + + alp::ALP_BLOCK_SIZE * + (sizeof(T) + alp::ALP_BLOCK_HEADER_SIZE)); + } + + private: + int encode_value(T value) { + values_.push_back(value); + return common::E_OK; + } + + int encode_batch_values(const T* values, uint32_t count) { + if (count == 0) { + return common::E_OK; + } + values_.insert(values_.end(), values, values + count); + return common::E_OK; + } + + std::vector values_; + std::vector page_buffer_; +}; + +class FloatAlpEncoder : public AlpEncoderBase { + public: + FloatAlpEncoder() {} +}; + +class DoubleAlpEncoder : public AlpEncoderBase { + public: + DoubleAlpEncoder() {} +}; + +} // namespace storage + +#endif // ENCODING_ALP_ENCODER_H diff --git a/cpp/src/encoding/alp/alp_format.h b/cpp/src/encoding/alp/alp_format.h new file mode 100644 index 000000000..4e0bda540 --- /dev/null +++ b/cpp/src/encoding/alp/alp_format.h @@ -0,0 +1,210 @@ +/* + * 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. + */ + +#ifndef ENCODING_ALP_FORMAT_H +#define ENCODING_ALP_FORMAT_H + +#include +#include + +#include +#include + +#include "alp_constants.h" + +namespace storage { +namespace alp { + +enum AlpStatus : int { + ALP_OK = 0, + ALP_INVALID_ARGUMENT = -1, + ALP_MALFORMED = -2, + ALP_BUFFER_TOO_SMALL = -3, + ALP_TYPE_MISMATCH = -4 +}; + +struct AlpPageHeader { + uint8_t version; + uint8_t data_type; + uint16_t reserved; + uint32_t value_count; + uint32_t block_count; + uint32_t reserved2; +}; + +struct AlpBlockHeader { + uint8_t scheme; + uint8_t factor; + uint8_t exponent; + uint8_t flags; + uint32_t value_count; + uint32_t exception_count; + uint32_t body_bytes; + int64_t for_base; + uint8_t bit_width; + uint8_t reserved[3]; +}; + +static const uint32_t ALP_PAGE_HEADER_SIZE = 16; +static const uint32_t ALP_BLOCK_HEADER_SIZE = 32; + +inline uint32_t AlpAlignUp(uint32_t value, uint32_t alignment) { + if (alignment == 0) { + return value; + } + const uint32_t remainder = value % alignment; + return remainder == 0 ? value : value + (alignment - remainder); +} + +inline void AlpAppendU8(std::vector& out, uint8_t v) { + out.push_back(v); +} + +inline void AlpAppendU16(std::vector& out, uint16_t v) { + const size_t old_size = out.size(); + out.resize(old_size + 2); + out[old_size] = static_cast(v & 0xFFu); + out[old_size + 1] = static_cast((v >> 8) & 0xFFu); +} + +inline void AlpAppendU32(std::vector& out, uint32_t v) { + const size_t old_size = out.size(); + out.resize(old_size + 4); + out[old_size] = static_cast(v & 0xFFu); + out[old_size + 1] = static_cast((v >> 8) & 0xFFu); + out[old_size + 2] = static_cast((v >> 16) & 0xFFu); + out[old_size + 3] = static_cast((v >> 24) & 0xFFu); +} + +inline void AlpAppendU64(std::vector& out, uint64_t v) { + const size_t old_size = out.size(); + out.resize(old_size + 8); + for (int i = 0; i < 8; ++i) { + out[old_size + static_cast(i)] = + static_cast((v >> (8 * i)) & 0xFFu); + } +} + +inline bool AlpReadU8(const uint8_t* data, uint32_t size, uint32_t& pos, + uint8_t& out) { + if (pos + 1 > size) { + return false; + } + out = data[pos++]; + return true; +} + +inline bool AlpReadU16(const uint8_t* data, uint32_t size, uint32_t& pos, + uint16_t& out) { + if (pos + 2 > size) { + return false; + } + out = static_cast(data[pos]) | + (static_cast(data[pos + 1]) << 8); + pos += 2; + return true; +} + +inline bool AlpReadU32(const uint8_t* data, uint32_t size, uint32_t& pos, + uint32_t& out) { + if (pos + 4 > size) { + return false; + } + out = static_cast(data[pos]) | + (static_cast(data[pos + 1]) << 8) | + (static_cast(data[pos + 2]) << 16) | + (static_cast(data[pos + 3]) << 24); + pos += 4; + return true; +} + +inline bool AlpReadU64(const uint8_t* data, uint32_t size, uint32_t& pos, + uint64_t& out) { + if (pos + 8 > size) { + return false; + } + out = 0; + for (int i = 0; i < 8; ++i) { + out |= static_cast(data[pos + i]) << (8 * i); + } + pos += 8; + return true; +} + +inline void AlpWritePageHeader(std::vector& out, + const AlpPageHeader& header) { + AlpAppendU8(out, header.version); + AlpAppendU8(out, header.data_type); + AlpAppendU16(out, header.reserved); + AlpAppendU32(out, header.value_count); + AlpAppendU32(out, header.block_count); + AlpAppendU32(out, header.reserved2); +} + +inline bool AlpParsePageHeader(const uint8_t* data, uint32_t size, + uint32_t& pos, AlpPageHeader& header) { + return AlpReadU8(data, size, pos, header.version) && + AlpReadU8(data, size, pos, header.data_type) && + AlpReadU16(data, size, pos, header.reserved) && + AlpReadU32(data, size, pos, header.value_count) && + AlpReadU32(data, size, pos, header.block_count) && + AlpReadU32(data, size, pos, header.reserved2); +} + +inline void AlpWriteBlockHeader(std::vector& out, + const AlpBlockHeader& header) { + AlpAppendU8(out, header.scheme); + AlpAppendU8(out, header.factor); + AlpAppendU8(out, header.exponent); + AlpAppendU8(out, header.flags); + AlpAppendU32(out, header.value_count); + AlpAppendU32(out, header.exception_count); + AlpAppendU32(out, header.body_bytes); + AlpAppendU64(out, static_cast(header.for_base)); + AlpAppendU8(out, header.bit_width); + AlpAppendU8(out, header.reserved[0]); + AlpAppendU8(out, header.reserved[1]); + AlpAppendU8(out, header.reserved[2]); +} + +inline bool AlpParseBlockHeader(const uint8_t* data, uint32_t size, + uint32_t& pos, AlpBlockHeader& header) { + uint64_t base = 0; + if (!AlpReadU8(data, size, pos, header.scheme) || + !AlpReadU8(data, size, pos, header.factor) || + !AlpReadU8(data, size, pos, header.exponent) || + !AlpReadU8(data, size, pos, header.flags) || + !AlpReadU32(data, size, pos, header.value_count) || + !AlpReadU32(data, size, pos, header.exception_count) || + !AlpReadU32(data, size, pos, header.body_bytes) || + !AlpReadU64(data, size, pos, base) || + !AlpReadU8(data, size, pos, header.bit_width) || + !AlpReadU8(data, size, pos, header.reserved[0]) || + !AlpReadU8(data, size, pos, header.reserved[1]) || + !AlpReadU8(data, size, pos, header.reserved[2])) { + return false; + } + header.for_base = static_cast(base); + return true; +} + +} // namespace alp +} // namespace storage + +#endif // ENCODING_ALP_FORMAT_H diff --git a/cpp/src/encoding/alp/alp_scalar.h b/cpp/src/encoding/alp/alp_scalar.h new file mode 100644 index 000000000..9d78e8acc --- /dev/null +++ b/cpp/src/encoding/alp/alp_scalar.h @@ -0,0 +1,658 @@ +/* + * 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. + */ + +#ifndef ENCODING_ALP_SCALAR_H +#define ENCODING_ALP_SCALAR_H + +#include + +#include +#include +#include +#include +#include + +#include "alp_format.h" +#include "alp_traits.h" +#ifdef ENABLE_SIMD +#include "alp_simde.h" +#endif + +namespace storage { +namespace alp { + +template +inline uint8_t AlpComputeBitWidth(Unsigned value) { + uint8_t width = 0; + while (value > 0) { + ++width; + value >>= 1; + } + return width; +} + +template +inline void AlpPackBits(const Unsigned* values, uint32_t count, + uint8_t bit_width, std::vector& body) { + body.clear(); + if (bit_width == 0 || count == 0) { + return; + } + const uint32_t raw_bytes = static_cast( + (static_cast(count) * bit_width + 7) / 8); + body.assign(AlpAlignUp(raw_bytes + 16, ALP_BODY_ALIGNMENT), 0); + + if (bit_width == sizeof(Unsigned) * 8) { + for (uint32_t i = 0; i < count; ++i) { + const Unsigned value = values[i]; + for (uint32_t b = 0; b < sizeof(Unsigned); ++b) { + body[i * sizeof(Unsigned) + b] = + static_cast((value >> (8 * b)) & 0xFFu); + } + } + return; + } + + // Word-based packing: build a 128-bit window for each value and OR it + // into two adjacent 64-bit words. This removes the per-bit loop while + // keeping the byte layout identical to the scalar reference. + const uint64_t mask = bit_width >= 64 + ? ~static_cast(0) + : ((static_cast(1) << bit_width) - 1); + uint64_t bit_pos = 0; + for (uint32_t i = 0; i < count; ++i) { + const uint64_t value = static_cast(values[i]) & mask; + const uint32_t byte_pos = static_cast(bit_pos >> 3); + const uint32_t bit_offset = static_cast(bit_pos & 7u); + const uint64_t low = value << bit_offset; + const uint64_t high = + bit_offset == 0 ? 0 : (value >> (64 - bit_offset)); + uint64_t dst = 0; + std::memcpy(&dst, body.data() + byte_pos, sizeof(dst)); + dst |= low; + std::memcpy(body.data() + byte_pos, &dst, sizeof(dst)); + if (high != 0) { + std::memcpy(&dst, body.data() + byte_pos + sizeof(dst), + sizeof(dst)); + dst |= high; + std::memcpy(body.data() + byte_pos + sizeof(dst), &dst, + sizeof(dst)); + } + bit_pos += bit_width; + } +} + +template +inline bool AlpUnpackBits(const uint8_t* body, uint32_t body_bytes, + uint32_t count, uint8_t bit_width, + std::vector& out) { + out.assign(count, 0); + if (bit_width == 0 || count == 0) { + return true; + } + const uint32_t raw_bytes = static_cast( + (static_cast(count) * bit_width + 7) / 8); + if (body_bytes < raw_bytes) { + return false; + } + + if (bit_width == sizeof(Unsigned) * 8) { + for (uint32_t i = 0; i < count; ++i) { + Unsigned value = 0; + for (uint32_t b = 0; b < sizeof(Unsigned); ++b) { + value |= static_cast(body[i * sizeof(Unsigned) + b]) + << (8 * b); + } + out[i] = value; + } + return true; + } + + // The body is padded with at least 16 zero bytes past the logical end. + // The SIMD path handles groups of four; the scalar tail uses two + // overlapping 64-bit loads and remains portable C++11. + uint32_t i = 0; +#ifdef ENABLE_SIMD + const uint32_t simd_count = count & ~static_cast(3); + if (simd_count > 0 && + !AlpSimdUnpackBits(body, simd_count, bit_width, &out[0])) { + return false; + } + i = simd_count; +#endif + const uint64_t mask = bit_width >= 64 + ? ~static_cast(0) + : ((static_cast(1) << bit_width) - 1); + uint64_t bit_pos = static_cast(i) * bit_width; + for (; i < count; ++i) { + const uint32_t byte_pos = static_cast(bit_pos >> 3); + const uint32_t bit_offset = static_cast(bit_pos & 7u); + uint64_t low = 0; + uint64_t high = 0; + std::memcpy(&low, body + byte_pos, sizeof(low)); + std::memcpy(&high, body + byte_pos + sizeof(low), sizeof(high)); + uint64_t value = low >> bit_offset; + if (bit_offset != 0) { + value |= high << (64 - bit_offset); + } + out[i] = static_cast(value & mask); + bit_pos += bit_width; + } + return true; +} + +template +inline uint64_t AlpEstimatedBodyBytes(uint32_t count, uint8_t bit_width) { + if (bit_width == 0) { + return 0; + } + const uint64_t raw_bytes = + (static_cast(count) * bit_width + 7) / 8; + return AlpAlignUp(static_cast(raw_bytes + 16), + ALP_BODY_ALIGNMENT); +} + +template +inline AlpStatus AlpEncodeValuesScalar( + const T* values, uint32_t count, uint8_t factor, uint8_t exponent, + typename AlpTypeTraits::Encoded* encoded, uint8_t* bitmap, T* exceptions, + uint32_t* exception_count) { + typedef AlpTypeTraits Traits; + typedef typename Traits::Encoded Encoded; + uint32_t ex_count = 0; + for (uint32_t i = 0; i < count; ++i) { + Encoded value = 0; + if (!Traits::EncodeValue(values[i], factor, exponent, &value) || + !Traits::BitwiseEqual(Traits::DecodeValue(value, factor, exponent), + values[i])) { + bitmap[i >> 3] |= static_cast(1u << (i & 7u)); + exceptions[ex_count++] = values[i]; + encoded[i] = 0; + } else { + encoded[i] = value; + } + } + *exception_count = ex_count; + return ALP_OK; +} + +template +inline AlpStatus AlpDecodeValuesScalar( + const typename AlpTypeTraits::Encoded* encoded, uint32_t count, + uint8_t factor, uint8_t exponent, const uint8_t* bitmap, + const T* exceptions, uint32_t exception_count, T* out) { + typedef AlpTypeTraits Traits; + uint32_t exception_index = 0; + for (uint32_t i = 0; i < count; ++i) { + T value = Traits::DecodeValue(encoded[i], factor, exponent); + if (bitmap != NULL && + (bitmap[i >> 3] & static_cast(1u << (i & 7u))) != 0) { + if (exception_index >= exception_count) { + return ALP_MALFORMED; + } + value = exceptions[exception_index++]; + } + out[i] = value; + } + return exception_index == exception_count ? ALP_OK : ALP_MALFORMED; +} + +template +inline AlpStatus AlpEncodeValues(const T* values, uint32_t count, + uint8_t factor, uint8_t exponent, + typename AlpTypeTraits::Encoded* encoded, + uint8_t* bitmap, T* exceptions, + uint32_t* exception_count) { +#ifdef ENABLE_SIMD + return AlpSimdEncodeValues(values, count, factor, exponent, encoded, bitmap, + exceptions, exception_count); +#else + return AlpEncodeValuesScalar(values, count, factor, exponent, encoded, + bitmap, exceptions, exception_count); +#endif +} + +template +inline AlpStatus AlpDecodeValues( + const typename AlpTypeTraits::Encoded* encoded, uint32_t count, + uint8_t factor, uint8_t exponent, const uint8_t* bitmap, + const T* exceptions, uint32_t exception_count, T* out) { +#ifdef ENABLE_SIMD + return AlpSimdDecodeValues(encoded, count, factor, exponent, bitmap, + exceptions, exception_count, out); +#else + return AlpDecodeValuesScalar(encoded, count, factor, exponent, bitmap, + exceptions, exception_count, out); +#endif +} + +template +inline uint64_t AlpEstimateCandidate(const T* samples, uint32_t sample_count, + uint32_t block_count, uint8_t factor, + uint8_t exponent, bool verify) { + typedef AlpTypeTraits Traits; + typedef typename Traits::Encoded Encoded; + typedef typename Traits::Unsigned Unsigned; + + uint32_t sample_exceptions = 0; + bool has_value = false; + Encoded min_value = 0; + Encoded max_value = 0; + for (uint32_t i = 0; i < sample_count; ++i) { + Encoded encoded = 0; + if (!Traits::EncodeValue(samples[i], factor, exponent, &encoded)) { + ++sample_exceptions; + continue; + } + if (verify && + !Traits::BitwiseEqual( + Traits::DecodeValue(encoded, factor, exponent), samples[i])) { + ++sample_exceptions; + continue; + } + if (!has_value) { + min_value = encoded; + max_value = encoded; + has_value = true; + } else { + min_value = std::min(min_value, encoded); + max_value = std::max(max_value, encoded); + } + } + + const uint64_t exception_estimate = + sample_count == 0 ? block_count + : static_cast(sample_exceptions) * + block_count / sample_count; + uint8_t bit_width = 0; + if (has_value) { + const Unsigned adjusted_max = + static_cast(max_value) - static_cast(min_value); + bit_width = AlpComputeBitWidth(adjusted_max); + } + const uint64_t body_bytes = + AlpEstimatedBodyBytes(block_count, bit_width); + const uint64_t bitmap_bytes = (static_cast(block_count) + 7) / 8; + return body_bytes + bitmap_bytes + exception_estimate * Traits::ValueSize(); +} + +template +inline bool AlpChooseFactorExponent(const T* values, uint32_t count, + uint8_t& best_factor, + uint8_t& best_exponent) { + typedef AlpTypeTraits Traits; + + const uint32_t sample_count = std::min(count, ALP_SAMPLE_SIZE); + std::vector samples(sample_count); + for (uint32_t i = 0; i < sample_count; ++i) { + const uint32_t index = + sample_count == count ? i : (i * count) / sample_count; + samples[i] = values[index]; + } + + // Stage 1: cheap screen over a small prefix of the sample. The decode + // round-trip check is skipped here; stage 2 performs the exact check on + // the shortlisted candidates. + static const uint32_t kStage1Samples = 8; + static const uint32_t kShortlist = 8; + const uint32_t stage1_count = std::min(sample_count, kStage1Samples); + uint64_t top_estimate[kShortlist]; + uint8_t top_factor[kShortlist]; + uint8_t top_exponent[kShortlist]; + uint32_t top_count = 0; + + for (uint8_t e = 0; e <= Traits::MaxExponent(); ++e) { + for (uint8_t f = 0; f <= e; ++f) { + const uint64_t estimate = + AlpEstimateCandidate(samples.empty() ? NULL : &samples[0], + stage1_count, count, f, e, true); + int pos = static_cast(top_count); + for (uint32_t i = 0; i < top_count; ++i) { + if (estimate < top_estimate[i] || + (estimate == top_estimate[i] && + (e > top_exponent[i] || + (e == top_exponent[i] && f > top_factor[i])))) { + pos = static_cast(i); + break; + } + } + if (pos < static_cast(kShortlist)) { + const uint32_t last = std::min(top_count, kShortlist - 1); + for (uint32_t j = last; j > static_cast(pos); --j) { + top_estimate[j] = top_estimate[j - 1]; + top_factor[j] = top_factor[j - 1]; + top_exponent[j] = top_exponent[j - 1]; + } + top_estimate[pos] = estimate; + top_factor[pos] = f; + top_exponent[pos] = e; + if (top_count < kShortlist) { + ++top_count; + } + } + } + } + + if (top_count == 0) { + return false; + } + + // Stage 2: exact round-trip check on the shortlisted candidates using the + // full sample. + uint64_t best_size = std::numeric_limits::max(); + bool found = false; + for (uint32_t i = 0; i < top_count; ++i) { + const uint8_t f = top_factor[i]; + const uint8_t e = top_exponent[i]; + const uint64_t estimate = + AlpEstimateCandidate(samples.empty() ? NULL : &samples[0], + sample_count, count, f, e, true); + if (!found || estimate < best_size || + (estimate == best_size && + (e > best_exponent || (e == best_exponent && f > best_factor)))) { + found = true; + best_size = estimate; + best_factor = f; + best_exponent = e; + } + } + return found; +} + +template +inline AlpStatus AlpEncodeBlock(const T* values, uint32_t count, + std::vector& out) { + typedef AlpTypeTraits Traits; + typedef typename Traits::Encoded Encoded; + typedef typename Traits::Unsigned Unsigned; + + if (values == NULL || count == 0 || count > ALP_BLOCK_SIZE) { + return ALP_INVALID_ARGUMENT; + } + + uint8_t factor = 0; + uint8_t exponent = 0; + if (!AlpChooseFactorExponent(values, count, factor, exponent)) { + return ALP_MALFORMED; + } + + std::vector encoded(count, 0); + std::vector bitmap((count + 7) / 8, 0); + std::vector exceptions(count); + uint32_t exception_count = 0; + AlpStatus encode_status = + AlpEncodeValues(values, count, factor, exponent, &encoded[0], + &bitmap[0], &exceptions[0], &exception_count); + if (encode_status != ALP_OK) { + return encode_status; + } + exceptions.resize(exception_count); + + bool has_value = false; + Encoded min_value = 0; + Encoded max_value = 0; + for (uint32_t i = 0; i < count; ++i) { + if ((bitmap[i >> 3] & static_cast(1u << (i & 7u))) != 0) { + continue; + } + const Encoded value = encoded[i]; + if (!has_value) { + min_value = value; + max_value = value; + has_value = true; + } else { + min_value = std::min(min_value, value); + max_value = std::max(max_value, value); + } + } + + const Encoded for_base = has_value ? min_value : 0; + const Unsigned adjusted_max = has_value + ? static_cast(max_value) - + static_cast(min_value) + : 0; + const uint8_t bit_width = AlpComputeBitWidth(adjusted_max); + + std::vector adjusted(count, 0); + for (uint32_t i = 0; i < count; ++i) { + if ((bitmap[i >> 3] & static_cast(1u << (i & 7u))) == 0) { + adjusted[i] = static_cast(encoded[i]) - + static_cast(for_base); + } + } + + std::vector body; + AlpPackBits(adjusted.empty() ? NULL : &adjusted[0], count, bit_width, body); + + const uint64_t bitmap_bytes = exception_count == 0 ? 0 : bitmap.size(); + const uint64_t alp_payload_bytes = + bitmap_bytes + + static_cast(exceptions.size()) * Traits::ValueSize() + + body.size(); + const uint64_t plain_payload_bytes = + static_cast(count) * Traits::ValueSize(); + const bool use_plain = + !has_value || alp_payload_bytes >= plain_payload_bytes; + + AlpBlockHeader header; + memset(&header, 0, sizeof(header)); + header.value_count = count; + header.factor = factor; + header.exponent = exponent; + header.bit_width = bit_width; + header.for_base = static_cast(for_base); + + if (use_plain) { + header.scheme = ALP_SCHEME_PLAIN; + header.exception_count = 0; + header.body_bytes = static_cast(plain_payload_bytes); + header.bit_width = 0; + header.for_base = 0; + AlpWriteBlockHeader(out, header); + for (uint32_t i = 0; i < count; ++i) { + Traits::AppendValue(out, values[i]); + } + return ALP_OK; + } + + header.scheme = ALP_SCHEME_ALP; + header.flags = exception_count == 0 ? 0 : 1; + header.exception_count = exception_count; + header.body_bytes = static_cast(body.size()); + AlpWriteBlockHeader(out, header); + + if (exception_count > 0) { + out.insert(out.end(), bitmap.begin(), bitmap.end()); + for (uint32_t i = 0; i < exceptions.size(); ++i) { + Traits::AppendValue(out, exceptions[i]); + } + } + out.insert(out.end(), body.begin(), body.end()); + return ALP_OK; +} + +template +inline AlpStatus AlpDecodeBlock(const uint8_t* data, uint32_t size, + uint32_t& pos, std::vector& out) { + typedef AlpTypeTraits Traits; + typedef typename Traits::Encoded Encoded; + typedef typename Traits::Unsigned Unsigned; + + AlpBlockHeader header; + if (!AlpParseBlockHeader(data, size, pos, header)) { + return ALP_MALFORMED; + } + if (header.value_count == 0 || header.value_count > ALP_BLOCK_SIZE || + header.exception_count > header.value_count || + header.bit_width > Traits::MaxBitWidth()) { + return ALP_MALFORMED; + } + + if (header.scheme == ALP_SCHEME_PLAIN) { + if (header.body_bytes != + static_cast(header.value_count * Traits::ValueSize())) { + return ALP_MALFORMED; + } + for (uint32_t i = 0; i < header.value_count; ++i) { + T value; + if (!Traits::ReadValue(data, size, pos, &value)) { + return ALP_BUFFER_TOO_SMALL; + } + out.push_back(value); + } + return ALP_OK; + } + + if (header.scheme != ALP_SCHEME_ALP) { + return ALP_MALFORMED; + } + + std::vector bitmap; + if (header.exception_count > 0) { + const uint32_t bitmap_bytes = (header.value_count + 7) / 8; + if (pos + bitmap_bytes > size) { + return ALP_BUFFER_TOO_SMALL; + } + bitmap.assign(data + pos, data + pos + bitmap_bytes); + pos += bitmap_bytes; + } + + std::vector exceptions; + exceptions.reserve(header.exception_count); + for (uint32_t i = 0; i < header.exception_count; ++i) { + T value; + if (!Traits::ReadValue(data, size, pos, &value)) { + return ALP_BUFFER_TOO_SMALL; + } + exceptions.push_back(value); + } + + const uint32_t expected_raw_bytes = + header.bit_width == 0 + ? 0 + : static_cast((static_cast(header.value_count) * + header.bit_width + + 7) / + 8); + const uint32_t expected_body_bytes = + header.bit_width == 0 + ? 0 + : AlpAlignUp(expected_raw_bytes + 16, ALP_BODY_ALIGNMENT); + if (header.body_bytes != expected_body_bytes || + pos + header.body_bytes > size) { + return ALP_MALFORMED; + } + + std::vector adjusted; + if (!AlpUnpackBits(data + pos, header.body_bytes, header.value_count, + header.bit_width, adjusted)) { + return ALP_MALFORMED; + } + pos += header.body_bytes; + + std::vector encoded(header.value_count); + for (uint32_t i = 0; i < header.value_count; ++i) { + const Unsigned raw = + adjusted[i] + static_cast(header.for_base); + std::memcpy(&encoded[i], &raw, sizeof(encoded[i])); + } + + std::vector decoded(header.value_count); + const uint8_t* bitmap_ptr = header.exception_count > 0 + ? (bitmap.empty() ? NULL : &bitmap[0]) + : NULL; + const T* exception_ptr = header.exception_count > 0 + ? (exceptions.empty() ? NULL : &exceptions[0]) + : NULL; + const AlpStatus decode_status = AlpDecodeValues( + &encoded[0], header.value_count, header.factor, header.exponent, + bitmap_ptr, exception_ptr, header.exception_count, &decoded[0]); + if (decode_status != ALP_OK) { + return decode_status; + } + out.insert(out.end(), decoded.begin(), decoded.end()); + return ALP_OK; +} + +template +inline AlpStatus AlpEncodePage(const T* values, uint32_t count, + std::vector& out) { + out.clear(); + if (values == NULL && count != 0) { + return ALP_INVALID_ARGUMENT; + } + const uint32_t block_count = + count == 0 ? 0 : (count + ALP_BLOCK_SIZE - 1) / ALP_BLOCK_SIZE; + + AlpPageHeader page; + page.version = ALP_FORMAT_VERSION; + page.data_type = AlpTypeTraits::DataType(); + page.reserved = 0; + page.value_count = count; + page.block_count = block_count; + page.reserved2 = 0; + AlpWritePageHeader(out, page); + + uint32_t offset = 0; + for (uint32_t b = 0; b < block_count; ++b) { + const uint32_t n = std::min(ALP_BLOCK_SIZE, count - offset); + const AlpStatus status = AlpEncodeBlock(values + offset, n, out); + if (status != ALP_OK) { + return status; + } + offset += n; + } + return ALP_OK; +} + +template +inline AlpStatus AlpDecodePage(const uint8_t* data, uint32_t size, + std::vector& out) { + out.clear(); + if (data == NULL && size != 0) { + return ALP_INVALID_ARGUMENT; + } + uint32_t pos = 0; + AlpPageHeader page; + if (!AlpParsePageHeader(data, size, pos, page)) { + return ALP_MALFORMED; + } + if (page.version != ALP_FORMAT_VERSION || + page.data_type != AlpTypeTraits::DataType() || page.reserved != 0 || + page.reserved2 != 0) { + return ALP_MALFORMED; + } + + out.reserve(page.value_count); + for (uint32_t b = 0; b < page.block_count; ++b) { + const AlpStatus status = AlpDecodeBlock(data, size, pos, out); + if (status != ALP_OK) { + return status; + } + } + if (out.size() != page.value_count || pos != size) { + return ALP_MALFORMED; + } + return ALP_OK; +} + +} // namespace alp +} // namespace storage + +#endif // ENCODING_ALP_SCALAR_H diff --git a/cpp/src/encoding/alp/alp_simde.h b/cpp/src/encoding/alp/alp_simde.h new file mode 100644 index 000000000..1f7ae1b15 --- /dev/null +++ b/cpp/src/encoding/alp/alp_simde.h @@ -0,0 +1,393 @@ +/* + * 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. + */ + +#ifndef ENCODING_ALP_SIMDE_H +#define ENCODING_ALP_SIMDE_H + +#include + +#include "alp_traits.h" +#include "simde/x86/avx2.h" +#include "simde/x86/avx512/cvt.h" +#include "simde/x86/avx512/cvtt.h" +#include "simde/x86/sse4.1.h" + +namespace storage { +namespace alp { + +template +struct AlpSimdKernel; + +template <> +struct AlpSimdKernel { + static AlpStatus EncodeValues(const float* values, uint32_t count, + uint8_t factor, uint8_t exponent, + int32_t* encoded, uint8_t* bitmap, + float* exceptions, + uint32_t* exception_count) { + typedef AlpTypeTraits Traits; + const float scale = ALP_FLOAT_EXP[exponent] * ALP_FLOAT_FRAC[factor]; + const float factor_value = static_cast(ALP_FACT[factor]); + const float frac_value = ALP_FLOAT_FRAC[exponent]; + + const simde__m256 vscale = simde_mm256_set1_ps(scale); + const simde__m256 vmagic = simde_mm256_set1_ps(ALP_FLOAT_MAGIC); + const simde__m256 vfactor = simde_mm256_set1_ps(factor_value); + const simde__m256 vfrac = simde_mm256_set1_ps(frac_value); + const simde__m256 vlower = + simde_mm256_set1_ps(ALP_FLOAT_ENCODING_LOWER_LIMIT); + const simde__m256 vupper = + simde_mm256_set1_ps(ALP_FLOAT_ENCODING_UPPER_LIMIT); + const simde__m256i vnegzero = + simde_mm256_set1_epi32(static_cast(0x80000000u)); + const simde__m256i vall = simde_mm256_set1_epi32(-1); + + uint32_t ex_count = 0; + uint32_t i = 0; + for (; i + 7 < count; i += 8) { + const simde__m256 v = simde_mm256_loadu_ps(values + i); + const simde__m256 scaled = simde_mm256_mul_ps(v, vscale); + const simde__m256 rounded = + simde_mm256_sub_ps(simde_mm256_add_ps(scaled, vmagic), vmagic); + const simde__m256i vi = simde_mm256_cvttps_epi32(rounded); + const simde__m256 back = simde_mm256_mul_ps( + simde_mm256_mul_ps(simde_mm256_cvtepi32_ps(vi), vfactor), + vfrac); + + const simde__m256 in_range = simde_mm256_and_ps( + simde_mm256_cmp_ps(scaled, vlower, SIMDE_CMP_GE_OQ), + simde_mm256_cmp_ps(scaled, vupper, SIMDE_CMP_LE_OQ)); + const simde__m256 rounded_in_range = simde_mm256_and_ps( + simde_mm256_cmp_ps(rounded, vlower, SIMDE_CMP_GE_OQ), + simde_mm256_cmp_ps(rounded, vupper, SIMDE_CMP_LE_OQ)); + const simde__m256 ordered = + simde_mm256_cmp_ps(v, v, SIMDE_CMP_ORD_Q); + const simde__m256 equal = + simde_mm256_cmp_ps(back, v, SIMDE_CMP_EQ_OQ); + const simde__m256i is_negzero = + simde_mm256_cmpeq_epi32(simde_mm256_castps_si256(v), vnegzero); + const simde__m256 not_negzero = simde_mm256_castsi256_ps( + simde_mm256_xor_si256(is_negzero, vall)); + const simde__m256 ok = simde_mm256_and_ps( + simde_mm256_and_ps(in_range, rounded_in_range), + simde_mm256_and_ps(ordered, + simde_mm256_and_ps(equal, not_negzero))); + const int mask = simde_mm256_movemask_ps(ok); + + simde_mm256_storeu_si256( + reinterpret_cast(encoded + i), vi); + for (int lane = 0; lane < 8; ++lane) { + if ((mask & (1 << lane)) == 0) { + const uint32_t index = i + static_cast(lane); + bitmap[index >> 3] |= + static_cast(1u << (index & 7u)); + exceptions[ex_count++] = values[index]; + encoded[index] = 0; + } + } + } + + for (; i < count; ++i) { + int32_t value = 0; + if (!Traits::EncodeValue(values[i], factor, exponent, &value) || + !Traits::BitwiseEqual( + Traits::DecodeValue(value, factor, exponent), values[i])) { + bitmap[i >> 3] |= static_cast(1u << (i & 7u)); + exceptions[ex_count++] = values[i]; + } else { + encoded[i] = value; + } + } + *exception_count = ex_count; + return ALP_OK; + } + + static AlpStatus DecodeValues(const int32_t* encoded, uint32_t count, + uint8_t factor, uint8_t exponent, + const uint8_t* bitmap, + const float* exceptions, + uint32_t exception_count, float* out) { + typedef AlpTypeTraits Traits; + const simde__m256 vfactor = + simde_mm256_set1_ps(static_cast(ALP_FACT[factor])); + const simde__m256 vfrac = simde_mm256_set1_ps(ALP_FLOAT_FRAC[exponent]); + + uint32_t i = 0; + for (; i + 7 < count; i += 8) { + const simde__m256i vi = simde_mm256_loadu_si256( + reinterpret_cast(encoded + i)); + const simde__m256 value = simde_mm256_mul_ps( + simde_mm256_mul_ps(simde_mm256_cvtepi32_ps(vi), vfactor), + vfrac); + simde_mm256_storeu_ps(out + i, value); + } + for (; i < count; ++i) { + out[i] = Traits::DecodeValue(encoded[i], factor, exponent); + } + + uint32_t exception_index = 0; + for (uint32_t index = 0; index < count; ++index) { + if (bitmap != NULL && + (bitmap[index >> 3] & + static_cast(1u << (index & 7u))) != 0) { + if (exception_index >= exception_count) { + return ALP_MALFORMED; + } + out[index] = exceptions[exception_index++]; + } + } + return exception_index == exception_count ? ALP_OK : ALP_MALFORMED; + } +}; + +template <> +struct AlpSimdKernel { + static AlpStatus EncodeValues(const double* values, uint32_t count, + uint8_t factor, uint8_t exponent, + int64_t* encoded, uint8_t* bitmap, + double* exceptions, + uint32_t* exception_count) { + typedef AlpTypeTraits Traits; + const double scale = ALP_DOUBLE_EXP[exponent] * ALP_DOUBLE_FRAC[factor]; + const double factor_value = static_cast(ALP_FACT[factor]); + const double frac_value = ALP_DOUBLE_FRAC[exponent]; + + const simde__m128d vscale = simde_mm_set1_pd(scale); + const simde__m128d vmagic = simde_mm_set1_pd(ALP_DOUBLE_MAGIC); + const simde__m128d vfactor = simde_mm_set1_pd(factor_value); + const simde__m128d vfrac = simde_mm_set1_pd(frac_value); + const simde__m128d vlower = + simde_mm_set1_pd(ALP_DOUBLE_ENCODING_LOWER_LIMIT); + const simde__m128d vupper = + simde_mm_set1_pd(ALP_DOUBLE_ENCODING_UPPER_LIMIT); + const simde__m128i vnegzero = + simde_mm_set1_epi64x(static_cast(0x8000000000000000ULL)); + const simde__m128i vall = simde_mm_set1_epi32(-1); + + uint32_t ex_count = 0; + uint32_t i = 0; + for (; i + 1 < count; i += 2) { + const simde__m128d v = simde_mm_loadu_pd(values + i); + const simde__m128d scaled = simde_mm_mul_pd(v, vscale); + const simde__m128d rounded = + simde_mm_sub_pd(simde_mm_add_pd(scaled, vmagic), vmagic); + const simde__m128i vi = simde_mm_cvttpd_epi64(rounded); + const simde__m128d back = simde_mm_mul_pd( + simde_mm_mul_pd(simde_mm_cvtepi64_pd(vi), vfactor), vfrac); + + const simde__m128d in_range = simde_mm_and_pd( + simde_mm_cmp_pd(scaled, vlower, SIMDE_CMP_GE_OQ), + simde_mm_cmp_pd(scaled, vupper, SIMDE_CMP_LE_OQ)); + const simde__m128d rounded_in_range = simde_mm_and_pd( + simde_mm_cmp_pd(rounded, vlower, SIMDE_CMP_GE_OQ), + simde_mm_cmp_pd(rounded, vupper, SIMDE_CMP_LE_OQ)); + const simde__m128d ordered = simde_mm_cmp_pd(v, v, SIMDE_CMP_ORD_Q); + const simde__m128d equal = + simde_mm_cmp_pd(back, v, SIMDE_CMP_EQ_OQ); + const simde__m128i is_negzero = + simde_mm_cmpeq_epi64(simde_mm_castpd_si128(v), vnegzero); + const simde__m128d not_negzero = + simde_mm_castsi128_pd(simde_mm_xor_si128(is_negzero, vall)); + const simde__m128d ok = simde_mm_and_pd( + simde_mm_and_pd(in_range, rounded_in_range), + simde_mm_and_pd(ordered, simde_mm_and_pd(equal, not_negzero))); + const int mask = simde_mm_movemask_pd(ok); + + simde_mm_storeu_si128(reinterpret_cast(encoded + i), + vi); + for (int lane = 0; lane < 2; ++lane) { + if ((mask & (1 << lane)) == 0) { + const uint32_t index = i + static_cast(lane); + bitmap[index >> 3] |= + static_cast(1u << (index & 7u)); + exceptions[ex_count++] = values[index]; + encoded[index] = 0; + } + } + } + + for (; i < count; ++i) { + int64_t value = 0; + if (!Traits::EncodeValue(values[i], factor, exponent, &value) || + !Traits::BitwiseEqual( + Traits::DecodeValue(value, factor, exponent), values[i])) { + bitmap[i >> 3] |= static_cast(1u << (i & 7u)); + exceptions[ex_count++] = values[i]; + } else { + encoded[i] = value; + } + } + *exception_count = ex_count; + return ALP_OK; + } + + static AlpStatus DecodeValues(const int64_t* encoded, uint32_t count, + uint8_t factor, uint8_t exponent, + const uint8_t* bitmap, + const double* exceptions, + uint32_t exception_count, double* out) { + typedef AlpTypeTraits Traits; + const simde__m128d vfactor = + simde_mm_set1_pd(static_cast(ALP_FACT[factor])); + const simde__m128d vfrac = simde_mm_set1_pd(ALP_DOUBLE_FRAC[exponent]); + + uint32_t i = 0; + for (; i + 1 < count; i += 2) { + const simde__m128i vi = simde_mm_loadu_si128( + reinterpret_cast(encoded + i)); + const simde__m128d value = simde_mm_mul_pd( + simde_mm_mul_pd(simde_mm_cvtepi64_pd(vi), vfactor), vfrac); + simde_mm_storeu_pd(out + i, value); + } + for (; i < count; ++i) { + out[i] = Traits::DecodeValue(encoded[i], factor, exponent); + } + + uint32_t exception_index = 0; + for (uint32_t index = 0; index < count; ++index) { + if (bitmap != NULL && + (bitmap[index >> 3] & + static_cast(1u << (index & 7u))) != 0) { + if (exception_index >= exception_count) { + return ALP_MALFORMED; + } + out[index] = exceptions[exception_index++]; + } + } + return exception_index == exception_count ? ALP_OK : ALP_MALFORMED; + } +}; + +template +inline AlpStatus AlpSimdEncodeValues( + const T* values, uint32_t count, uint8_t factor, uint8_t exponent, + typename AlpTypeTraits::Encoded* encoded, uint8_t* bitmap, T* exceptions, + uint32_t* exception_count) { + return AlpSimdKernel::EncodeValues(values, count, factor, exponent, + encoded, bitmap, exceptions, + exception_count); +} + +template +inline AlpStatus AlpSimdDecodeValues( + const typename AlpTypeTraits::Encoded* encoded, uint32_t count, + uint8_t factor, uint8_t exponent, const uint8_t* bitmap, + const T* exceptions, uint32_t exception_count, T* out) { + return AlpSimdKernel::DecodeValues(encoded, count, factor, exponent, + bitmap, exceptions, exception_count, + out); +} + +template +struct AlpSimdUnpackKernel; + +template <> +struct AlpSimdUnpackKernel { + static bool Unpack(const uint8_t* body, uint32_t count, uint8_t bit_width, + uint32_t* out) { + if (bit_width == 0) { + for (uint32_t i = 0; i < count; ++i) { + out[i] = 0; + } + return true; + } + const uint64_t mask64 = + bit_width >= 32 ? static_cast(0xFFFFFFFFu) + : ((static_cast(1) << bit_width) - 1); + const simde__m256i mask = + simde_mm256_set1_epi64x(static_cast(mask64)); + const simde__m256i perm = + simde_mm256_setr_epi32(0, 2, 4, 6, 0, 0, 0, 0); + for (uint32_t i = 0; i + 3 < count; i += 4) { + int32_t byte_offsets[4]; + int64_t bit_offsets[4]; + for (int lane = 0; lane < 4; ++lane) { + const uint64_t bit_pos = + static_cast(i + lane) * bit_width; + byte_offsets[lane] = static_cast(bit_pos >> 3); + bit_offsets[lane] = static_cast(bit_pos & 7u); + } + const simde__m128i indices = simde_mm_loadu_si128( + reinterpret_cast(byte_offsets)); + const simde__m256i windows = simde_mm256_i32gather_epi64( + reinterpret_cast(body), indices, 1); + const simde__m256i shifts = simde_mm256_loadu_si256( + reinterpret_cast(bit_offsets)); + const simde__m256i shifted = + simde_mm256_srlv_epi64(windows, shifts); + const simde__m256i masked = simde_mm256_and_si256(shifted, mask); + const simde__m256i compact = + simde_mm256_permutevar8x32_epi32(masked, perm); + simde_mm_storeu_si128(reinterpret_cast(out + i), + simde_mm256_castsi256_si128(compact)); + } + return true; + } +}; + +template <> +struct AlpSimdUnpackKernel { + static bool Unpack(const uint8_t* body, uint32_t count, uint8_t bit_width, + uint64_t* out) { + if (bit_width == 0) { + for (uint32_t i = 0; i < count; ++i) { + out[i] = 0; + } + return true; + } + const uint64_t mask64 = + bit_width >= 64 ? ~static_cast(0) + : ((static_cast(1) << bit_width) - 1); + const simde__m256i mask = + simde_mm256_set1_epi64x(static_cast(mask64)); + for (uint32_t i = 0; i + 3 < count; i += 4) { + int32_t byte_offsets[4]; + int64_t bit_offsets[4]; + for (int lane = 0; lane < 4; ++lane) { + const uint64_t bit_pos = + static_cast(i + lane) * bit_width; + byte_offsets[lane] = static_cast(bit_pos >> 3); + bit_offsets[lane] = static_cast(bit_pos & 7u); + } + const simde__m128i indices = simde_mm_loadu_si128( + reinterpret_cast(byte_offsets)); + const simde__m256i windows = simde_mm256_i32gather_epi64( + reinterpret_cast(body), indices, 1); + const simde__m256i shifts = simde_mm256_loadu_si256( + reinterpret_cast(bit_offsets)); + const simde__m256i shifted = + simde_mm256_srlv_epi64(windows, shifts); + const simde__m256i masked = simde_mm256_and_si256(shifted, mask); + simde_mm256_storeu_si256(reinterpret_cast(out + i), + masked); + } + return true; + } +}; + +template +inline bool AlpSimdUnpackBits(const uint8_t* body, uint32_t count, + uint8_t bit_width, Unsigned* out) { + return AlpSimdUnpackKernel::Unpack(body, count, bit_width, out); +} + +} // namespace alp +} // namespace storage + +#endif // ENCODING_ALP_SIMDE_H diff --git a/cpp/src/encoding/alp/alp_traits.h b/cpp/src/encoding/alp/alp_traits.h new file mode 100644 index 000000000..e1046df69 --- /dev/null +++ b/cpp/src/encoding/alp/alp_traits.h @@ -0,0 +1,158 @@ +/* + * 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. + */ + +#ifndef ENCODING_ALP_TRAITS_H +#define ENCODING_ALP_TRAITS_H + +#include + +#include +#include + +#include "alp_format.h" + +namespace storage { +namespace alp { + +template +struct AlpTypeTraits; + +template <> +struct AlpTypeTraits { + typedef int32_t Encoded; + typedef uint32_t Unsigned; + + static uint8_t DataType() { return ALP_FLOAT; } + static uint8_t MaxExponent() { return ALP_FLOAT_MAX_EXPONENT; } + static uint8_t MaxBitWidth() { return 32; } + static uint32_t ValueSize() { return 4; } + + static bool EncodeValue(float value, uint8_t factor, uint8_t exponent, + int32_t* out) { + if (!std::isfinite(value) || (value == 0.0f && std::signbit(value))) { + return false; + } + const float scaled = + value * ALP_FLOAT_EXP[exponent] * ALP_FLOAT_FRAC[factor]; + if (!std::isfinite(scaled) || scaled < ALP_FLOAT_ENCODING_LOWER_LIMIT || + scaled > ALP_FLOAT_ENCODING_UPPER_LIMIT) { + return false; + } + const float rounded = scaled + ALP_FLOAT_MAGIC - ALP_FLOAT_MAGIC; + if (!std::isfinite(rounded) || + rounded < ALP_FLOAT_ENCODING_LOWER_LIMIT || + rounded > ALP_FLOAT_ENCODING_UPPER_LIMIT) { + return false; + } + *out = static_cast(rounded); + return true; + } + + static float DecodeValue(int32_t encoded, uint8_t factor, + uint8_t exponent) { + const float factor_value = static_cast(ALP_FACT[factor]); + return static_cast(encoded) * factor_value * + ALP_FLOAT_FRAC[exponent]; + } + + static void AppendValue(std::vector& out, float value) { + uint32_t bits = 0; + std::memcpy(&bits, &value, sizeof(bits)); + AlpAppendU32(out, bits); + } + + static bool ReadValue(const uint8_t* data, uint32_t size, uint32_t& pos, + float* value) { + uint32_t bits = 0; + if (!AlpReadU32(data, size, pos, bits)) { + return false; + } + std::memcpy(value, &bits, sizeof(bits)); + return true; + } + + static bool BitwiseEqual(float lhs, float rhs) { + return std::memcmp(&lhs, &rhs, sizeof(float)) == 0; + } +}; + +template <> +struct AlpTypeTraits { + typedef int64_t Encoded; + typedef uint64_t Unsigned; + + static uint8_t DataType() { return ALP_DOUBLE; } + static uint8_t MaxExponent() { return ALP_DOUBLE_MAX_EXPONENT; } + static uint8_t MaxBitWidth() { return 64; } + static uint32_t ValueSize() { return 8; } + + static bool EncodeValue(double value, uint8_t factor, uint8_t exponent, + int64_t* out) { + if (!std::isfinite(value) || (value == 0.0 && std::signbit(value))) { + return false; + } + const double scaled = + value * ALP_DOUBLE_EXP[exponent] * ALP_DOUBLE_FRAC[factor]; + if (!std::isfinite(scaled) || + scaled < ALP_DOUBLE_ENCODING_LOWER_LIMIT || + scaled > ALP_DOUBLE_ENCODING_UPPER_LIMIT) { + return false; + } + const double rounded = scaled + ALP_DOUBLE_MAGIC - ALP_DOUBLE_MAGIC; + if (!std::isfinite(rounded) || + rounded < ALP_DOUBLE_ENCODING_LOWER_LIMIT || + rounded > ALP_DOUBLE_ENCODING_UPPER_LIMIT) { + return false; + } + *out = static_cast(rounded); + return true; + } + + static double DecodeValue(int64_t encoded, uint8_t factor, + uint8_t exponent) { + const double factor_value = static_cast(ALP_FACT[factor]); + return static_cast(encoded) * factor_value * + ALP_DOUBLE_FRAC[exponent]; + } + + static void AppendValue(std::vector& out, double value) { + uint64_t bits = 0; + std::memcpy(&bits, &value, sizeof(bits)); + AlpAppendU64(out, bits); + } + + static bool ReadValue(const uint8_t* data, uint32_t size, uint32_t& pos, + double* value) { + uint64_t bits = 0; + if (!AlpReadU64(data, size, pos, bits)) { + return false; + } + std::memcpy(value, &bits, sizeof(bits)); + return true; + } + + static bool BitwiseEqual(double lhs, double rhs) { + return std::memcmp(&lhs, &rhs, sizeof(double)) == 0; + } +}; + +} // namespace alp +} // namespace storage + +#endif // ENCODING_ALP_TRAITS_H diff --git a/cpp/src/encoding/decoder.h b/cpp/src/encoding/decoder.h index 035023f8c..b41f8fa2c 100644 --- a/cpp/src/encoding/decoder.h +++ b/cpp/src/encoding/decoder.h @@ -30,6 +30,8 @@ class Decoder { Decoder() {} virtual ~Decoder() {} virtual void reset() = 0; + virtual void destroy() {} + virtual bool owns_resources() const { return false; } virtual bool has_remaining(const common::ByteStream& buffer) = 0; virtual int read_boolean(bool& ret_value, common::ByteStream& in) = 0; virtual int read_int32(int32_t& ret_value, common::ByteStream& in) = 0; diff --git a/cpp/src/encoding/decoder_factory.h b/cpp/src/encoding/decoder_factory.h index 16111a690..953e9883d 100644 --- a/cpp/src/encoding/decoder_factory.h +++ b/cpp/src/encoding/decoder_factory.h @@ -20,6 +20,7 @@ #ifndef ENCODING_DECODER_FACTORY_H #define ENCODING_DECODER_FACTORY_H +#include "alp/alp_decoder.h" #include "camel_decoder.h" #include "chimp_decoder.h" #include "common/global.h" @@ -188,6 +189,16 @@ class DecoderFactory { return nullptr; } + case ALP: + switch (data_type) { + case FLOAT: + ALLOC_AND_RETURN_DECODER(FloatAlpDecoder); + case DOUBLE: + ALLOC_AND_RETURN_DECODER(DoubleAlpDecoder); + default: + return nullptr; + } + default: // Not supported encoding return nullptr; @@ -195,7 +206,12 @@ class DecoderFactory { return nullptr; } - static void free(Decoder* decoder) { common::mem_free(decoder); } + static void free(Decoder* decoder) { + if (decoder != nullptr && decoder->owns_resources()) { + decoder->destroy(); + } + common::mem_free(decoder); + } }; } // end namespace storage #endif // ENCODING_DECODER_FACTORY_H diff --git a/cpp/src/encoding/encoder_factory.h b/cpp/src/encoding/encoder_factory.h index 052cbf8ee..655187a76 100644 --- a/cpp/src/encoding/encoder_factory.h +++ b/cpp/src/encoding/encoder_factory.h @@ -20,6 +20,7 @@ #ifndef ENCODING_ENCODER_FACTORY_H #define ENCODING_ENCODER_FACTORY_H +#include "alp/alp_encoder.h" #include "camel_encoder.h" #include "chimp_encoder.h" #include "common/global.h" @@ -201,6 +202,16 @@ class EncoderFactory { return nullptr; } + case ALP: + switch (data_type) { + case FLOAT: + ALLOC_AND_RETURN_ENCODER(FloatAlpEncoder); + case DOUBLE: + ALLOC_AND_RETURN_ENCODER(DoubleAlpEncoder); + default: + return nullptr; + } + case DIFF: case BITMAP: case GORILLA_V1: diff --git a/cpp/src/reader/chunk_reader.cc b/cpp/src/reader/chunk_reader.cc index 271b5c206..41d7524de 100644 --- a/cpp/src/reader/chunk_reader.cc +++ b/cpp/src/reader/chunk_reader.cc @@ -153,7 +153,7 @@ int ChunkReader::alloc_compressor_and_value_decoder( // fall through DecoderFactory and be reported as an allocation failure. // Distinguish malformed file metadata from a genuine OOM before creating // any decoder object. - if (encoding < common::PLAIN || encoding > common::CAMEL) { + if (encoding < common::PLAIN || encoding > common::ALP) { return E_TSFILE_CORRUPTED; } if (value_decoder_ != nullptr) { diff --git a/cpp/test/encoding/alp_codec_test.cc b/cpp/test/encoding/alp_codec_test.cc new file mode 100644 index 000000000..afbc2d88e --- /dev/null +++ b/cpp/test/encoding/alp_codec_test.cc @@ -0,0 +1,218 @@ +/* + * 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 +#include +#include +#include + +#include "encoding/alp/alp_scalar.h" +#include "gtest/gtest.h" + +namespace storage { +namespace alp { + +template +static bool BitEqual(T lhs, T rhs) { + return std::memcmp(&lhs, &rhs, sizeof(T)) == 0; +} + +template +static void ExpectRoundTrip(const std::vector& values) { + std::vector encoded; + ASSERT_EQ(ALP_OK, + AlpEncodePage(values.empty() ? NULL : &values[0], + static_cast(values.size()), encoded)); + std::vector decoded; + ASSERT_EQ(ALP_OK, + AlpDecodePage(encoded.empty() ? NULL : &encoded[0], + static_cast(encoded.size()), decoded)); + ASSERT_EQ(values.size(), decoded.size()); + for (size_t i = 0; i < values.size(); ++i) { + EXPECT_TRUE(BitEqual(values[i], decoded[i])) << "index " << i; + } + + // Encoding must be deterministic. + std::vector encoded_again; + ASSERT_EQ(ALP_OK, AlpEncodePage(values.empty() ? NULL : &values[0], + static_cast(values.size()), + encoded_again)); + EXPECT_EQ(encoded, encoded_again); +} + +TEST(AlpCodecTest, FloatDecimalRoundTrip) { + std::vector values; + for (int i = 0; i < 2048; ++i) { + values.push_back(static_cast(i % 100) / 100.0f - 0.5f); + } + ExpectRoundTrip(values); +} + +TEST(AlpCodecTest, DoubleDecimalRoundTrip) { + std::vector values; + for (int i = 0; i < 2048; ++i) { + values.push_back(static_cast(i % 1000) / 1000.0 - 0.5); + } + ExpectRoundTrip(values); +} + +TEST(AlpCodecTest, SpecialValuesRoundTrip) { + const float nan_f = std::numeric_limits::quiet_NaN(); + const float inf_f = std::numeric_limits::infinity(); + const double nan_d = std::numeric_limits::quiet_NaN(); + const double inf_d = std::numeric_limits::infinity(); + + std::vector floats = {0.0f, + -0.0f, + 1.0f, + -1.0f, + nan_f, + inf_f, + -inf_f, + std::numeric_limits::denorm_min(), + std::numeric_limits::max(), + std::numeric_limits::lowest()}; + ExpectRoundTrip(floats); + + std::vector doubles = {0.0, + -0.0, + 1.0, + -1.0, + nan_d, + inf_d, + -inf_d, + std::numeric_limits::denorm_min(), + std::numeric_limits::max(), + std::numeric_limits::lowest()}; + ExpectRoundTrip(doubles); +} + +TEST(AlpCodecTest, PartialAndEmptyBlocks) { + ExpectRoundTrip(std::vector()); + ExpectRoundTrip(std::vector(1, 1.25f)); + ExpectRoundTrip(std::vector(ALP_BLOCK_SIZE - 1, 2.5f)); + ExpectRoundTrip(std::vector(ALP_BLOCK_SIZE, 3.75f)); + ExpectRoundTrip(std::vector(ALP_BLOCK_SIZE + 1, 4.125f)); +} + +TEST(AlpCodecTest, HighEntropyFallsBackButRoundTrips) { + std::vector values; + uint64_t state = 0x123456789abcdef0ULL; + for (int i = 0; i < 4096; ++i) { + state = state * 6364136223846793005ULL + 1442695040888963407ULL; + uint64_t bits = state; + double value = 0.0; + std::memcpy(&value, &bits, sizeof(value)); + if (std::isfinite(value)) { + values.push_back(value); + } + } + ExpectRoundTrip(values); +} + +TEST(AlpCodecTest, TruncatedPayloadReturnsError) { + std::vector values; + for (int i = 0; i < 128; ++i) { + values.push_back(static_cast(i) * 0.25); + } + std::vector encoded; + ASSERT_EQ(ALP_OK, + AlpEncodePage(&values[0], static_cast(values.size()), + encoded)); + std::vector decoded; + EXPECT_NE(ALP_OK, AlpDecodePage(&encoded[0], + static_cast(encoded.size() / 2), + decoded)); +} + +template +static AlpBlockHeader FirstBlockHeader(const std::vector& encoded) { + AlpBlockHeader header; + std::memset(&header, 0, sizeof(header)); + uint32_t pos = ALP_PAGE_HEADER_SIZE; + EXPECT_TRUE(AlpParseBlockHeader(encoded.empty() ? NULL : &encoded[0], + static_cast(encoded.size()), pos, + header)); + return header; +} + +TEST(AlpCodecTest, DecimalDataUsesAlpScheme) { + std::vector values(ALP_BLOCK_SIZE); + for (size_t i = 0; i < values.size(); ++i) { + values[i] = static_cast(i % 128) * 0.01f; + } + std::vector encoded; + ASSERT_EQ(ALP_OK, + AlpEncodePage(&values[0], static_cast(values.size()), + encoded)); + const AlpBlockHeader header = FirstBlockHeader(encoded); + EXPECT_EQ(ALP_SCHEME_ALP, header.scheme); + EXPECT_LT(header.body_bytes, values.size() * sizeof(float)); +} + +TEST(AlpCodecTest, HighEntropyUsesPlainScheme) { + std::vector values(ALP_BLOCK_SIZE); + uint64_t state = 0x9e3779b97f4a7c15ULL; + for (size_t i = 0; i < values.size(); ++i) { + state = state * 6364136223846793005ULL + 1442695040888963407ULL; + uint64_t bits = state ^ (state >> 29); + std::memcpy(&values[i], &bits, sizeof(double)); + if (!std::isfinite(values[i])) { + values[i] = static_cast(i); + } + } + std::vector encoded; + ASSERT_EQ(ALP_OK, + AlpEncodePage(&values[0], static_cast(values.size()), + encoded)); + const AlpBlockHeader header = FirstBlockHeader(encoded); + EXPECT_EQ(ALP_SCHEME_PLAIN, header.scheme); +} + +TEST(AlpCodecTest, SmallIntegerBitWidths) { + std::vector values(ALP_BLOCK_SIZE); + for (size_t i = 0; i < values.size(); ++i) { + values[i] = static_cast(i % 256); + } + std::vector encoded; + ASSERT_EQ(ALP_OK, + AlpEncodePage(&values[0], static_cast(values.size()), + encoded)); + const AlpBlockHeader header = FirstBlockHeader(encoded); + EXPECT_EQ(ALP_SCHEME_ALP, header.scheme); + EXPECT_EQ(8, header.bit_width); + EXPECT_EQ(0, header.for_base); + + std::vector signed_values(ALP_BLOCK_SIZE); + for (size_t i = 0; i < signed_values.size(); ++i) { + signed_values[i] = static_cast(static_cast(i % 256) - 128); + } + std::vector signed_encoded; + ASSERT_EQ(ALP_OK, AlpEncodePage(&signed_values[0], + static_cast(signed_values.size()), + signed_encoded)); + const AlpBlockHeader signed_header = + FirstBlockHeader(signed_encoded); + EXPECT_EQ(ALP_SCHEME_ALP, signed_header.scheme); + EXPECT_EQ(8, signed_header.bit_width); + EXPECT_EQ(-128, signed_header.for_base); +} + +} // namespace alp +} // namespace storage diff --git a/cpp/test/encoding/alp_integration_test.cc b/cpp/test/encoding/alp_integration_test.cc new file mode 100644 index 000000000..4ba4cbaf1 --- /dev/null +++ b/cpp/test/encoding/alp_integration_test.cc @@ -0,0 +1,150 @@ +/* + * 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 +#include +#include + +#include "common/allocator/byte_stream.h" +#include "common/db_common.h" +#include "encoding/decoder_factory.h" +#include "encoding/encoder_factory.h" +#include "gtest/gtest.h" + +namespace storage { + +template +static bool BitEqual(T lhs, T rhs) { + return std::memcmp(&lhs, &rhs, sizeof(T)) == 0; +} + +static int ReadAll(common::ByteStream& stream, std::vector& bytes) { + bytes.clear(); + bytes.reserve(static_cast(stream.total_size())); + common::ByteStream::BufferIterator iter = stream.init_buffer_iterator(); + while (true) { + common::ByteStream::Buffer buffer = iter.get_next_buf(); + if (buffer.buf_ == NULL || buffer.len_ == 0) { + break; + } + const uint8_t* begin = reinterpret_cast(buffer.buf_); + bytes.insert(bytes.end(), begin, begin + buffer.len_); + } + return bytes.size() == stream.total_size() ? common::E_OK + : common::E_PARTIAL_READ; +} + +TEST(AlpIntegrationTest, EncodingNameIsRegistered) { + EXPECT_STREQ("ALP", common::get_encoding_name(common::ALP)); +} + +TEST(AlpIntegrationTest, FactoryRejectsNonFloatTypes) { + EXPECT_EQ(nullptr, + EncoderFactory::alloc_value_encoder(common::ALP, common::INT32)); + EXPECT_EQ(nullptr, + EncoderFactory::alloc_value_encoder(common::ALP, common::INT64)); + EXPECT_EQ(nullptr, + DecoderFactory::alloc_value_decoder(common::ALP, common::INT32)); + EXPECT_EQ(nullptr, + DecoderFactory::alloc_value_decoder(common::ALP, common::INT64)); +} + +TEST(AlpIntegrationTest, FloatEncoderDecoderRoundTrip) { + std::vector values; + for (int i = 0; i < 4096; ++i) { + values.push_back(static_cast(i % 257) * 0.01f - 1.25f); + } + values.push_back(std::numeric_limits::quiet_NaN()); + values.push_back(-0.0f); + + Encoder* encoder = + EncoderFactory::alloc_value_encoder(common::ALP, common::FLOAT); + ASSERT_NE(nullptr, encoder); + common::ByteStream encoded(1024, common::MOD_DEFAULT); + ASSERT_EQ(common::E_OK, + encoder->encode_batch( + &values[0], static_cast(values.size()), encoded)); + ASSERT_EQ(common::E_OK, encoder->flush(encoded)); + EncoderFactory::free(encoder); + + std::vector bytes; + ASSERT_EQ(common::E_OK, ReadAll(encoded, bytes)); + ASSERT_FALSE(bytes.empty()); + + common::ByteStream wrapped; + wrapped.wrap_from(reinterpret_cast(&bytes[0]), + static_cast(bytes.size())); + Decoder* decoder = + DecoderFactory::alloc_value_decoder(common::ALP, common::FLOAT); + ASSERT_NE(nullptr, decoder); + std::vector decoded(values.size()); + int actual = 0; + ASSERT_EQ(common::E_OK, + decoder->read_exact_float( + &decoded[0], static_cast(values.size()), wrapped)); + (void)actual; + DecoderFactory::free(decoder); + + ASSERT_EQ(values.size(), decoded.size()); + for (size_t i = 0; i < values.size(); ++i) { + EXPECT_TRUE(BitEqual(values[i], decoded[i])) << "index " << i; + } +} + +TEST(AlpIntegrationTest, DoubleEncoderDecoderRoundTrip) { + std::vector values; + for (int i = 0; i < 4096; ++i) { + values.push_back(static_cast(i % 509) * 0.001 - 0.25); + } + values.push_back(std::numeric_limits::quiet_NaN()); + values.push_back(-0.0); + + Encoder* encoder = + EncoderFactory::alloc_value_encoder(common::ALP, common::DOUBLE); + ASSERT_NE(nullptr, encoder); + common::ByteStream encoded(1024, common::MOD_DEFAULT); + ASSERT_EQ(common::E_OK, + encoder->encode_batch( + &values[0], static_cast(values.size()), encoded)); + ASSERT_EQ(common::E_OK, encoder->flush(encoded)); + EncoderFactory::free(encoder); + + std::vector bytes; + ASSERT_EQ(common::E_OK, ReadAll(encoded, bytes)); + ASSERT_FALSE(bytes.empty()); + + common::ByteStream wrapped; + wrapped.wrap_from(reinterpret_cast(&bytes[0]), + static_cast(bytes.size())); + Decoder* decoder = + DecoderFactory::alloc_value_decoder(common::ALP, common::DOUBLE); + ASSERT_NE(nullptr, decoder); + std::vector decoded(values.size()); + ASSERT_EQ(common::E_OK, + decoder->read_exact_double( + &decoded[0], static_cast(values.size()), wrapped)); + DecoderFactory::free(decoder); + + ASSERT_EQ(values.size(), decoded.size()); + for (size_t i = 0; i < values.size(); ++i) { + EXPECT_TRUE(BitEqual(values[i], decoded[i])) << "index " << i; + } +} + +} // namespace storage diff --git a/cpp/test/encoding/alp_simde_test.cc b/cpp/test/encoding/alp_simde_test.cc new file mode 100644 index 000000000..92ef240d6 --- /dev/null +++ b/cpp/test/encoding/alp_simde_test.cc @@ -0,0 +1,152 @@ +/* + * 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 +#include +#include + +#include "encoding/alp/alp_scalar.h" +#include "gtest/gtest.h" + +#ifdef ENABLE_SIMD +#include "encoding/alp/alp_simde.h" +#endif + +namespace storage { +namespace alp { + +#ifdef ENABLE_SIMD + +template +static void ScalarEncodeValues(const T* values, uint32_t count, uint8_t factor, + uint8_t exponent, + typename AlpTypeTraits::Encoded* encoded, + uint8_t* bitmap, T* exceptions, + uint32_t* exception_count) { + typedef AlpTypeTraits Traits; + uint32_t ex_count = 0; + for (uint32_t i = 0; i < count; ++i) { + typename Traits::Encoded value = 0; + if (!Traits::EncodeValue(values[i], factor, exponent, &value) || + !Traits::BitwiseEqual(Traits::DecodeValue(value, factor, exponent), + values[i])) { + bitmap[i >> 3] |= static_cast(1u << (i & 7u)); + exceptions[ex_count++] = values[i]; + } else { + encoded[i] = value; + } + } + *exception_count = ex_count; +} + +template +static void ExpectSimdMatchesScalar(const std::vector& values) { + typedef AlpTypeTraits Traits; + ASSERT_FALSE(values.empty()); + uint8_t factor = 0; + uint8_t exponent = 0; + ASSERT_TRUE(AlpChooseFactorExponent( + &values[0], static_cast(values.size()), factor, exponent)); + + const uint32_t count = static_cast(values.size()); + std::vector scalar_encoded(count, 0); + std::vector simd_encoded(count, 0); + std::vector scalar_bitmap((count + 7) / 8, 0); + std::vector simd_bitmap((count + 7) / 8, 0); + std::vector scalar_exceptions(count); + std::vector simd_exceptions(count); + uint32_t scalar_exception_count = 0; + uint32_t simd_exception_count = 0; + + ScalarEncodeValues(&values[0], count, factor, exponent, &scalar_encoded[0], + &scalar_bitmap[0], &scalar_exceptions[0], + &scalar_exception_count); + ASSERT_EQ(ALP_OK, + AlpSimdEncodeValues(&values[0], count, factor, exponent, + &simd_encoded[0], &simd_bitmap[0], + &simd_exceptions[0], &simd_exception_count)); + + EXPECT_EQ(scalar_exception_count, simd_exception_count); + EXPECT_EQ(scalar_bitmap, simd_bitmap); + EXPECT_EQ(scalar_encoded, simd_encoded); + ASSERT_EQ(scalar_exception_count, simd_exception_count); + for (uint32_t i = 0; i < scalar_exception_count; ++i) { + EXPECT_TRUE( + Traits::BitwiseEqual(scalar_exceptions[i], simd_exceptions[i])); + } + + std::vector simd_decoded(count); + ASSERT_EQ(ALP_OK, + AlpSimdDecodeValues( + &scalar_encoded[0], count, factor, exponent, + scalar_exception_count == 0 ? NULL : &scalar_bitmap[0], + scalar_exception_count == 0 ? NULL : &scalar_exceptions[0], + scalar_exception_count, &simd_decoded[0])); + for (uint32_t i = 0; i < count; ++i) { + EXPECT_TRUE(Traits::BitwiseEqual(values[i], simd_decoded[i])) + << "index " << i; + } +} + +TEST(AlpSimdeTest, FloatValuesMatchScalar) { + std::vector values; + for (int i = 0; i < 1024; ++i) { + values.push_back(static_cast(i % 251) * 0.01f - 1.25f); + } + ExpectSimdMatchesScalar(values); +} + +TEST(AlpSimdeTest, DoubleValuesMatchScalar) { + std::vector values; + for (int i = 0; i < 1024; ++i) { + values.push_back(static_cast(i % 509) * 0.001 - 0.25); + } + ExpectSimdMatchesScalar(values); +} + +TEST(AlpSimdeTest, SpecialValuesMatchScalar) { + const float nan_f = std::numeric_limits::quiet_NaN(); + const double nan_d = std::numeric_limits::quiet_NaN(); + std::vector floats = {0.0f, + -0.0f, + 1.0f, + -1.0f, + nan_f, + std::numeric_limits::infinity(), + -std::numeric_limits::infinity(), + std::numeric_limits::denorm_min(), + std::numeric_limits::max()}; + ExpectSimdMatchesScalar(floats); + + std::vector doubles = {0.0, + -0.0, + 1.0, + -1.0, + nan_d, + std::numeric_limits::infinity(), + -std::numeric_limits::infinity(), + std::numeric_limits::denorm_min(), + std::numeric_limits::max()}; + ExpectSimdMatchesScalar(doubles); +} + +#endif // ENABLE_SIMD + +} // namespace alp +} // namespace storage diff --git a/cpp/test/encoding/alp_tsfile_test.cc b/cpp/test/encoding/alp_tsfile_test.cc new file mode 100644 index 000000000..1c0d86cf4 --- /dev/null +++ b/cpp/test/encoding/alp_tsfile_test.cc @@ -0,0 +1,127 @@ +/* + * 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 +#include +#include +#include +#include + +#ifdef _WIN32 +#include +#else +#include +#endif + +#include "common/path.h" +#include "common/record.h" +#include "common/schema.h" +#include "gtest/gtest.h" +#include "reader/qds_without_timegenerator.h" +#include "reader/tsfile_reader.h" +#include "writer/tsfile_writer.h" + +namespace storage { + +class AlpTsFileTest : public ::testing::Test { + protected: + void SetUp() override { + libtsfile_init(); + writer_ = new TsFileWriter(); +#ifdef _WIN32 + const int pid = _getpid(); +#else + const int pid = static_cast(getpid()); +#endif + file_name_ = + std::string("alp_tsfile_test_") + std::to_string(pid) + ".tsfile"; + std::remove(file_name_.c_str()); + int flags = O_RDWR | O_CREAT | O_TRUNC; +#ifdef _WIN32 + flags |= O_BINARY; +#endif + ASSERT_EQ(common::E_OK, writer_->open(file_name_, flags, 0666)); + } + + void TearDown() override { + delete writer_; + writer_ = nullptr; + std::remove(file_name_.c_str()); + libtsfile_destroy(); + } + + TsFileWriter* writer_ = nullptr; + std::string file_name_; +}; + +TEST_F(AlpTsFileTest, FloatAndDoubleRoundTrip) { + const std::string device = "alp_device"; + const std::string float_name = "f"; + const std::string double_name = "d"; + ASSERT_EQ( + common::E_OK, + writer_->register_timeseries( + device, MeasurementSchema(float_name, common::FLOAT, common::ALP, + common::UNCOMPRESSED))); + ASSERT_EQ( + common::E_OK, + writer_->register_timeseries( + device, MeasurementSchema(double_name, common::DOUBLE, common::ALP, + common::UNCOMPRESSED))); + + const int row_count = 20000; + const int64_t base_time = 1700000000000LL; + for (int i = 0; i < row_count; ++i) { + TsRecord record(base_time + i, device); + const float fv = static_cast(i % 1000) * 0.01f - 5.0f; + const double dv = static_cast(i % 10000) * 0.0001 - 0.5; + record.add_point(float_name, fv); + record.add_point(double_name, dv); + ASSERT_EQ(common::E_OK, writer_->write_record(record)); + } + ASSERT_EQ(common::E_OK, writer_->flush()); + ASSERT_EQ(common::E_OK, writer_->close()); + + TsFileReader reader; + ASSERT_EQ(common::E_OK, reader.open(file_name_)); + std::vector select_list = {device + "." + float_name, + device + "." + double_name}; + ResultSet* result = nullptr; + ASSERT_EQ(common::E_OK, reader.query(select_list, base_time, + base_time + row_count, result)); + ASSERT_NE(nullptr, result); + auto* qds = static_cast(result); + + int rows = 0; + bool has_next = false; + while (IS_SUCC(qds->next(has_next)) && has_next) { + const int i = rows; + const float expected_f = static_cast(i % 1000) * 0.01f - 5.0f; + const double expected_d = static_cast(i % 10000) * 0.0001 - 0.5; + EXPECT_EQ(expected_f, qds->get_value(device + "." + float_name)); + EXPECT_EQ(expected_d, + qds->get_value(device + "." + double_name)); + ++rows; + } + EXPECT_EQ(row_count, rows); + reader.destroy_query_data_set(qds); + ASSERT_EQ(common::E_OK, reader.close()); +} + +} // namespace storage diff --git a/python/tests/test_write_and_read.py b/python/tests/test_write_and_read.py index da5288478..05e8d6edd 100644 --- a/python/tests/test_write_and_read.py +++ b/python/tests/test_write_and_read.py @@ -137,6 +137,64 @@ def test_row_record_write_and_read(): os.remove("record_write_and_read.tsfile") +def test_alp_encoding_round_trip(): + assert TSEncoding.ALP == 15 + + file_path = "alp_encoding_round_trip.tsfile" + try: + if os.path.exists(file_path): + os.remove(file_path) + + writer = TsFileWriter(file_path) + writer.register_timeseries( + "root.alp", + TimeseriesSchema( + "f", TSDataType.FLOAT, TSEncoding.ALP, Compressor.UNCOMPRESSED + ), + ) + writer.register_timeseries( + "root.alp", + TimeseriesSchema( + "d", TSDataType.DOUBLE, TSEncoding.ALP, Compressor.UNCOMPRESSED + ), + ) + + row_count = 32 + for i in range(row_count): + writer.write_row_record( + RowRecord( + "root.alp", + i, + [ + Field("f", i * 0.1, TSDataType.FLOAT), + Field("d", i * 0.1, TSDataType.DOUBLE), + ], + ) + ) + writer.close() + + reader = TsFileReader(file_path) + result = reader.query_timeseries("root.alp", ["f", "d"], 0, 100) + values = [] + while result.next(): + values.append( + ( + result.get_value_by_index(2), + result.get_value_by_index(3), + ) + ) + result.close() + reader.close() + + assert len(values) == row_count + for i, (f, d) in enumerate(values): + assert f == pytest.approx(i * 0.1, abs=1e-6) + assert d == pytest.approx(i * 0.1, abs=1e-12) + finally: + if os.path.exists(file_path): + os.remove(file_path) + + def test_tree_query_to_dataframe_variants(): file_path = "tree_query_to_dataframe.tsfile" device_ids = [ diff --git a/python/tsfile/constants.py b/python/tsfile/constants.py index 237a14f2c..0b21510fa 100644 --- a/python/tsfile/constants.py +++ b/python/tsfile/constants.py @@ -200,6 +200,7 @@ class TSEncoding(IntEnum): SPRINTZ = 12 RLBE = 13 CAMEL = 14 + ALP = 15 @unique diff --git a/python/tsfile/tsfile_cpp.pxd b/python/tsfile/tsfile_cpp.pxd index b0b2d6b78..59a2b1806 100644 --- a/python/tsfile/tsfile_cpp.pxd +++ b/python/tsfile/tsfile_cpp.pxd @@ -72,6 +72,7 @@ cdef extern from "cwrapper/tsfile_cwrapper.h": TS_ENCODING_SPRINTZ = 12, TS_ENCODING_RLBE = 13, TS_ENCODING_CAMEL = 14, + TS_ENCODING_ALP = 15, TS_ENCODING_INVALID = 255 ctypedef enum CompressionType: diff --git a/python/tsfile/tsfile_py_cpp.pyx b/python/tsfile/tsfile_py_cpp.pyx index c8f155ae8..21b089254 100644 --- a/python/tsfile/tsfile_py_cpp.pyx +++ b/python/tsfile/tsfile_py_cpp.pyx @@ -128,6 +128,7 @@ cdef dict TS_ENCODING_MAP = { TSEncodingPy.SPRINTZ: TSEncoding.TS_ENCODING_SPRINTZ, TSEncodingPy.RLBE: TSEncoding.TS_ENCODING_RLBE, TSEncodingPy.CAMEL: TSEncoding.TS_ENCODING_CAMEL, + TSEncodingPy.ALP: TSEncoding.TS_ENCODING_ALP, } cdef dict COMPRESSION_TYPE_MAP = {