diff --git a/.github/workflows/develop-bench.yml b/.github/workflows/develop-bench.yml index a8b1757f44a..5b3ec4f3e64 100644 --- a/.github/workflows/develop-bench.yml +++ b/.github/workflows/develop-bench.yml @@ -83,7 +83,7 @@ jobs: timeout-minutes: 120 runs-on: >- ${{ github.repository == 'vortex-data/vortex' - && format('runs-on={0}/runner=bench-dedicated/family=c8gd.metal-24xl/image=ubuntu24-full-arm64-pre-v2/extras=s3-cache/tag={1}', github.run_id, matrix.benchmark.id) + && format('runs-on={0}/runner=bench-dedicated/family=c8gd.metal-24xl/image=ubuntu24-full-arm64-pre-v2/extras=s3-cache/tag={1}{2}', github.run_id, matrix.benchmark.id, matrix.benchmark.variant_id) || 'ubuntu-latest' }} strategy: fail-fast: false @@ -93,6 +93,14 @@ jobs: name: Random Access build_args: "--features lance" v4_ingest: true + # Same benchmark, reading the data from S3 instead of local NVMe. + - id: random-access-bench + variant_id: "-s3" + name: Random Access (S3) + build_args: "--features lance" + v4_ingest: true + remote_data_dir: >- + s3://vortex-ci-benchmark-datasets/develop/random-access/ - id: compress-bench name: Compression build_args: "--features lance" @@ -135,6 +143,25 @@ jobs: extra_args: "--debuginfo-strip=false" parca_agent_version: "0.49.0" + - name: Setup AWS CLI + uses: aws-actions/configure-aws-credentials@e1253824e5c10ff9df46874f81ed3ec929e19cfd # v6 + with: + role-to-assume: arn:aws:iam::245040174862:role/GitHubBenchmarkRole + aws-region: us-east-1 + + - name: Upload benchmark data to S3 + if: matrix.benchmark.remote_data_dir != null + shell: bash + env: + AWS_REGION: "us-east-1" + run: | + set -Eeu -o pipefail -x + + target/release_debug/${{ matrix.benchmark.id }} --prepare-data \ + --formats parquet,vortex,lance + aws s3 rm --recursive "${{ matrix.benchmark.remote_data_dir }}" + aws s3 cp --recursive vortex-bench/data "${{ matrix.benchmark.remote_data_dir }}" + - name: Setup benchmark environment run: sudo bash scripts/setup-benchmark.sh @@ -144,8 +171,10 @@ jobs: env: RUST_BACKTRACE: full FLAT_LAYOUT_INLINE_ARRAY_NODE: "1" + AWS_REGION: "us-east-1" run: | - python3 scripts/random-access-split.py --emit-ingest-records + python3 scripts/random-access-split.py --emit-ingest-records \ + ${{ matrix.benchmark.remote_data_dir && format('--remote-data-dir {0}', matrix.benchmark.remote_data_dir) || '' }} - name: Run ${{ matrix.benchmark.name }} benchmark (per-dataset) if: matrix.benchmark.id == 'compress-bench' @@ -165,12 +194,6 @@ jobs: run: | python3 scripts/string-split.py - - name: Setup AWS CLI - uses: aws-actions/configure-aws-credentials@e1253824e5c10ff9df46874f81ed3ec929e19cfd # v6 - with: - role-to-assume: arn:aws:iam::245040174862:role/GitHubBenchmarkRole - aws-region: us-east-1 - - name: Upload Benchmark Results shell: bash run: | diff --git a/.github/workflows/pr-bench-dispatch.yml b/.github/workflows/pr-bench-dispatch.yml index 3cb60c619bf..fa9a701d758 100644 --- a/.github/workflows/pr-bench-dispatch.yml +++ b/.github/workflows/pr-bench-dispatch.yml @@ -35,6 +35,11 @@ jobs: uses: ./.github/workflows/pr-bench-random-access.yml secrets: inherit + all-random-access-s3-bench: + needs: remove-all-label + uses: ./.github/workflows/pr-bench-random-access-s3.yml + secrets: inherit + all-compression-bench: needs: remove-all-label uses: ./.github/workflows/pr-bench-compress.yml @@ -68,6 +73,22 @@ jobs: uses: ./.github/workflows/pr-bench-random-access.yml secrets: inherit + remove-random-access-s3-label: + runs-on: ubuntu-latest + timeout-minutes: 10 + if: github.event.label.name == 'action/bench-random-access-s3' + steps: + - uses: actions-ecosystem/action-remove-labels@2ce5d41b4b6aa8503e285553f75ed56e0a40bae0 # v1 + if: github.event.pull_request.head.repo.full_name == 'vortex-data/vortex' + with: + labels: action/bench-random-access-s3 + fail_on_error: true + + random-access-s3-bench: + needs: remove-random-access-s3-label + uses: ./.github/workflows/pr-bench-random-access-s3.yml + secrets: inherit + remove-compress-label: runs-on: ubuntu-latest timeout-minutes: 10 diff --git a/.github/workflows/pr-bench-random-access-s3.yml b/.github/workflows/pr-bench-random-access-s3.yml new file mode 100644 index 00000000000..d5d77cab2d3 --- /dev/null +++ b/.github/workflows/pr-bench-random-access-s3.yml @@ -0,0 +1,24 @@ +# Runs the random-access benchmark for a pull request, reading the data from S3. + +name: PR Random Access S3 Benchmark + +on: + workflow_call: { } + workflow_dispatch: { } + +permissions: + contents: read + pull-requests: write # for commenting on PRs + id-token: write # enables AWS-GitHub OIDC + +jobs: + bench: + uses: ./.github/workflows/pr-bench-runner.yml + secrets: inherit + with: + benchmark_id: random-access-bench + benchmark_name: Random Access (S3) + with_lance: true + variant_id: "-s3" + remote_data_dir: >- + s3://vortex-ci-benchmark-datasets/${{ github.ref_name }}/${{ github.run_id }}/random-access/ diff --git a/.github/workflows/pr-bench-runner.yml b/.github/workflows/pr-bench-runner.yml index ef4904c4830..a21678f5f46 100644 --- a/.github/workflows/pr-bench-runner.yml +++ b/.github/workflows/pr-bench-runner.yml @@ -4,7 +4,7 @@ name: PR Benchmark Runner concurrency: # The group causes runs to queue instead of running in parallel. - group: ${{ github.workflow }}-${{ github.head_ref || github.run_id }}-${{ inputs.benchmark_id }} + group: ${{ github.workflow }}-${{ github.head_ref || github.run_id }}-${{ inputs.benchmark_id }}${{ inputs.variant_id }} # Don't cancel benchmarks that are already running, instead just queue them up. cancel-in-progress: false @@ -21,6 +21,20 @@ on: required: false type: boolean default: false + variant_id: + description: >- + Suffix distinguishing runs of the same benchmark, e.g. "-s3". Keeps the PR comment + tag and the concurrency group of a variant separate from the default run. + required: false + type: string + default: "" + remote_data_dir: + description: >- + When set, the benchmark data is uploaded to this S3 prefix and read back from there + instead of local disk. Only supported by random-access-bench. + required: false + type: string + default: "" permissions: contents: read @@ -122,6 +136,26 @@ jobs: extra_args: "--debuginfo-strip=false" parca_agent_version: "0.49.0" + - name: Setup AWS CLI + if: github.event.pull_request.head.repo.fork == false + uses: aws-actions/configure-aws-credentials@e1253824e5c10ff9df46874f81ed3ec929e19cfd # v6 + with: + role-to-assume: arn:aws:iam::245040174862:role/GitHubBenchmarkRole + aws-region: us-east-1 + + - name: Upload benchmark data to S3 + if: inputs.remote_data_dir != '' && github.event.pull_request.head.repo.fork == false + shell: bash + env: + AWS_REGION: "us-east-1" + run: | + set -Eeu -o pipefail -x + + target/release_debug/${{ inputs.benchmark_id }} --prepare-data \ + --formats ${{ inputs.with_lance && 'parquet,vortex,lance' || 'parquet,vortex' }} + aws s3 rm --recursive "${{ inputs.remote_data_dir }}" + aws s3 cp --recursive vortex-bench/data "${{ inputs.remote_data_dir }}" + - name: Setup benchmark environment run: sudo bash scripts/setup-benchmark.sh @@ -131,8 +165,10 @@ jobs: env: RUST_BACKTRACE: full FLAT_LAYOUT_INLINE_ARRAY_NODE: "1" + AWS_REGION: "us-east-1" run: | - python3 scripts/random-access-split.py + python3 scripts/random-access-split.py \ + ${{ inputs.remote_data_dir != '' && format('--remote-data-dir {0}', inputs.remote_data_dir) || '' }} - name: Run ${{ inputs.benchmark_name }} benchmark (per-dataset) if: inputs.benchmark_id == 'compress-bench' @@ -152,13 +188,6 @@ jobs: run: | python3 scripts/string-split.py - - name: Setup AWS CLI - if: github.event.pull_request.head.repo.fork == false - uses: aws-actions/configure-aws-credentials@e1253824e5c10ff9df46874f81ed3ec929e19cfd # v6 - with: - role-to-assume: arn:aws:iam::245040174862:role/GitHubBenchmarkRole - aws-region: us-east-1 - - name: Install uv uses: spiraldb/actions/.github/actions/setup-uv@a746510eafaa926484c354541cfc49b2ec06cc63 # 0.18.6 with: @@ -181,7 +210,7 @@ jobs: uses: thollander/actions-comment-pull-request@24bffb9b452ba05a4f3f77933840a6a841d1b32b # v3 with: file-path: comment.md - comment-tag: bench-pr-comment-${{ inputs.benchmark_id }} + comment-tag: bench-pr-comment-${{ inputs.benchmark_id }}${{ inputs.variant_id }} - name: Comment PR on failure if: failure() && github.event_name == 'pull_request' && github.event.pull_request.head.repo.fork == false @@ -191,4 +220,11 @@ jobs: # BENCHMARK FAILED Benchmark `${{ inputs.benchmark_name }}` failed! Check the [workflow run](${{ github.server_url }}/${{ github.repository }}/actions/runs/${{ github.run_id }}) for details. - comment-tag: bench-pr-comment-${{ inputs.benchmark_id }} + comment-tag: bench-pr-comment-${{ inputs.benchmark_id }}${{ inputs.variant_id }} + + - name: Delete benchmark data from S3 + if: always() && inputs.remote_data_dir != '' && github.event.pull_request.head.repo.fork == false + shell: bash + env: + AWS_REGION: "us-east-1" + run: aws s3 rm --recursive "${{ inputs.remote_data_dir }}" diff --git a/Cargo.lock b/Cargo.lock index 090f170e388..2d2d38d70ce 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -16,7 +16,7 @@ checksum = "35f0f96ce78e38c3dc6d8948aa8163d06385be74000f3c7a95bf1eef35d3ea32" dependencies = [ "cipher", "cpubits", - "cpufeatures", + "cpufeatures 0.3.1", ] [[package]] @@ -1061,6 +1061,49 @@ version = "1.5.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f2032f911046de80f0a198e0901378627c33f59ea0ac00e363d481118bd70a53" +[[package]] +name = "aws-config" +version = "1.12.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b8d7b388a9fc3a6db15a5ec778c38b354eff1364882c94d08e0252f7a47dcaa4" +dependencies = [ + "aws-credential-types", + "aws-runtime", + "aws-sdk-sso", + "aws-sdk-ssooidc", + "aws-sdk-sts", + "aws-smithy-async", + "aws-smithy-http", + "aws-smithy-json", + "aws-smithy-runtime", + "aws-smithy-runtime-api", + "aws-smithy-schema", + "aws-smithy-types", + "aws-types", + "bytes", + "fastrand", + "hex", + "http", + "sha1 0.10.7", + "time", + "tokio", + "tracing", + "url", + "zeroize", +] + +[[package]] +name = "aws-credential-types" +version = "1.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e93964ffdaf57857f544be3666a5f57570bb699e934700f11b49708f61bb556e" +dependencies = [ + "aws-smithy-async", + "aws-smithy-runtime-api", + "aws-smithy-types", + "zeroize", +] + [[package]] name = "aws-lc-rs" version = "1.18.1" @@ -1084,6 +1127,339 @@ dependencies = [ "pkg-config", ] +[[package]] +name = "aws-runtime" +version = "1.10.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2b8a9911551b4ea6ca13805ef52ed96f7d2bbb43cc3b4a14cb0776a71f33cfaa" +dependencies = [ + "aws-credential-types", + "aws-sigv4", + "aws-smithy-async", + "aws-smithy-http", + "aws-smithy-runtime", + "aws-smithy-runtime-api", + "aws-smithy-types", + "aws-types", + "bytes", + "bytes-utils", + "fastrand", + "http", + "http-body", + "percent-encoding", + "pin-project-lite", + "tracing", + "uuid", +] + +[[package]] +name = "aws-sdk-sso" +version = "1.113.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "87184f36ade012f46cdc773c7be10031b4545a1aea70a9a22a044e8f36f1e5c5" +dependencies = [ + "arc-swap", + "aws-credential-types", + "aws-runtime", + "aws-smithy-async", + "aws-smithy-http", + "aws-smithy-json", + "aws-smithy-observability", + "aws-smithy-runtime", + "aws-smithy-runtime-api", + "aws-smithy-schema", + "aws-smithy-types", + "aws-types", + "bytes", + "fastrand", + "http", + "regex-lite", + "tracing", +] + +[[package]] +name = "aws-sdk-ssooidc" +version = "1.115.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "84edc7495d56e006dff0031cbc9782784902456f7024651c2da902a3d2336262" +dependencies = [ + "arc-swap", + "aws-credential-types", + "aws-runtime", + "aws-smithy-async", + "aws-smithy-http", + "aws-smithy-json", + "aws-smithy-observability", + "aws-smithy-runtime", + "aws-smithy-runtime-api", + "aws-smithy-schema", + "aws-smithy-types", + "aws-types", + "bytes", + "fastrand", + "http", + "regex-lite", + "tracing", +] + +[[package]] +name = "aws-sdk-sts" +version = "1.118.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "df7572569d018ca72923073494eac3c4671a0daea90b9e3923c2d27a3a02a2c1" +dependencies = [ + "arc-swap", + "aws-credential-types", + "aws-runtime", + "aws-smithy-async", + "aws-smithy-http", + "aws-smithy-json", + "aws-smithy-observability", + "aws-smithy-query", + "aws-smithy-runtime", + "aws-smithy-runtime-api", + "aws-smithy-schema", + "aws-smithy-types", + "aws-smithy-xml", + "aws-types", + "fastrand", + "http", + "regex-lite", + "tracing", +] + +[[package]] +name = "aws-sigv4" +version = "1.6.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2312577f088c9fbf4206dfdb884cf1de9407b43e1a923cbed5237775116fc24b" +dependencies = [ + "aws-credential-types", + "aws-smithy-http", + "aws-smithy-runtime-api", + "aws-smithy-types", + "bytes", + "form_urlencoded", + "hex", + "hmac", + "http", + "percent-encoding", + "sha2", + "time", + "tracing", +] + +[[package]] +name = "aws-smithy-async" +version = "1.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f02e407fb3b54891734224b9ffac8a71fdd35f542500fa1af95754a6b2beb316" +dependencies = [ + "futures-util", + "pin-project-lite", + "tokio", +] + +[[package]] +name = "aws-smithy-http" +version = "0.64.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "639b4d8f8555f24a9be649811c3eb0b4d4616f4d61daf0c32e28873bc1ea9af1" +dependencies = [ + "aws-smithy-runtime-api", + "aws-smithy-types", + "bytes", + "bytes-utils", + "futures-core", + "futures-util", + "http", + "http-body", + "http-body-util", + "percent-encoding", + "pin-project-lite", + "pin-utils", + "tracing", +] + +[[package]] +name = "aws-smithy-http-client" +version = "1.4.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7bd25384a4e437aa8d8f339afad4b69e786b936a7cb10db668a7aaf66717b1a8" +dependencies = [ + "aws-smithy-async", + "aws-smithy-runtime-api", + "aws-smithy-types", + "h2", + "http", + "hyper", + "hyper-rustls", + "hyper-util", + "pin-project-lite", + "rustls", + "rustls-native-certs", + "rustls-pki-types", + "tokio", + "tokio-rustls", + "tower", + "tracing", +] + +[[package]] +name = "aws-smithy-json" +version = "0.63.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3385d469edbe8b60cc72002784652b5efca39178192aa9cc4b44c9875c6bdc18" +dependencies = [ + "aws-smithy-runtime-api", + "aws-smithy-schema", + "aws-smithy-types", +] + +[[package]] +name = "aws-smithy-observability" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8e86338c869539a581bf161247762a6e87f92c5c075060057b5ed6d06632ed0c" +dependencies = [ + "aws-smithy-runtime-api", +] + +[[package]] +name = "aws-smithy-query" +version = "0.62.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f1d1d71f6562be974caa85442ecd90194c40fdb5df045f182a6c2e872ce95056" +dependencies = [ + "aws-smithy-runtime-api", + "aws-smithy-schema", + "aws-smithy-types", + "aws-smithy-xml", + "urlencoding", +] + +[[package]] +name = "aws-smithy-runtime" +version = "1.15.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "591e0024cdf2a8de8711860f7e1d5cfbca26b800ce13793a9aa2f4f5918e6d95" +dependencies = [ + "aws-smithy-async", + "aws-smithy-http", + "aws-smithy-http-client", + "aws-smithy-observability", + "aws-smithy-runtime-api", + "aws-smithy-schema", + "aws-smithy-types", + "bytes", + "fastrand", + "http", + "http-body", + "http-body-util", + "pin-project-lite", + "pin-utils", + "tokio", + "tracing", +] + +[[package]] +name = "aws-smithy-runtime-api" +version = "1.18.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9bb0deb40a69de59f7809251cfe1ab16de135f42c49d222d1202b4d0ce19720f" +dependencies = [ + "aws-smithy-async", + "aws-smithy-runtime-api-macros", + "aws-smithy-types", + "bytes", + "http", + "pin-project-lite", + "tokio", + "tracing", + "zeroize", +] + +[[package]] +name = "aws-smithy-runtime-api-macros" +version = "1.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "221eaa237ddf1ca79b60d1372aad77e47f9c0ea5b3ce5099da8c61d027dc77b3" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.119", +] + +[[package]] +name = "aws-smithy-schema" +version = "0.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e8f395d93304280b64b7632fea798d177e74897fe7f063416ce627cd6fa24829" +dependencies = [ + "aws-smithy-runtime-api", + "aws-smithy-types", + "http", +] + +[[package]] +name = "aws-smithy-types" +version = "1.8.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "69bb407740a197147da48238ecc94498493c9e85445732360cec180296ca45f1" +dependencies = [ + "base64-simd", + "bytes", + "bytes-utils", + "http", + "http-body", + "http-body-util", + "itoa", + "num-integer", + "pin-project-lite", + "pin-utils", + "ryu", + "serde", + "time", +] + +[[package]] +name = "aws-smithy-xml" +version = "0.62.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0b932c8d6dc127fc980eecd78f8694ae9b9551b69a93a7def2a199c1c0033daf" +dependencies = [ + "aws-smithy-runtime-api", + "aws-smithy-schema", + "aws-smithy-types", + "xmlparser", +] + +[[package]] +name = "aws-types" +version = "1.6.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "209f3a6d82a6e9e5f94abbed94c7a26e1c052341002bf57a5fb5481f625896fc" +dependencies = [ + "aws-credential-types", + "aws-smithy-async", + "aws-smithy-runtime-api", + "aws-smithy-schema", + "aws-smithy-types", + "rustc_version", + "tracing", +] + +[[package]] +name = "backon" +version = "1.6.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cffb0e931875b666fc4fcb20fee52e9bbd1ef836fd9e9e04ec21555f9f85f7ef" +dependencies = [ + "fastrand", + "gloo-timers", + "tokio", +] + [[package]] name = "base16ct" version = "1.0.0" @@ -1102,6 +1478,16 @@ version = "0.23.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ac07cdecf99051d9a5238b80f35af32cdeba5b336e55d957b318b50137e18da5" +[[package]] +name = "base64-simd" +version = "0.8.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "339abbe78e73178762e23bea9dfd08e697eb3f3301cd4be981c0f78ba5859195" +dependencies = [ + "outref", + "vsimd", +] + [[package]] name = "beamterm-core" version = "1.0.0" @@ -1256,7 +1642,7 @@ dependencies = [ "cc", "cfg-if", "constant_time_eq", - "cpufeatures", + "cpufeatures 0.3.1", ] [[package]] @@ -1424,6 +1810,16 @@ version = "1.12.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "fc652a48c352aef3ea3aed32080501cf3ef6ed5da78602a020c991775b0aff04" +[[package]] +name = "bytes-utils" +version = "0.1.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7dafe3a8757b027e2be6e4e5601ed563c55989fcf1546e933c66c8eb3a058d35" +dependencies = [ + "bytes", + "either", +] + [[package]] name = "bzip2" version = "0.6.1" @@ -1507,7 +1903,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "65c35e4b699c7e15ccbe7ee35c005e4fc0a278d22238a2857e6ce2dadeda1b06" dependencies = [ "cfg-if", - "cpufeatures", + "cpufeatures 0.3.1", "rand_core 0.10.1", ] @@ -1980,6 +2376,15 @@ version = "0.1.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "15b85f9c39137c3a891689859392b1bd49812121d0d61c9caf00d46ed5ce06ae" +[[package]] +name = "cpufeatures" +version = "0.2.17" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "59ed5838eebb26a2bb2e58f6d5b5316989ae9d08bab10e0e6d103e656d1b0280" +dependencies = [ + "libc", +] + [[package]] name = "cpufeatures" version = "0.3.1" @@ -2004,6 +2409,16 @@ version = "2.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "217698eaf96b4a3f0bc4f3662aaa55bdf913cd54d7204591faa790070c6d0853" +[[package]] +name = "crc-fast" +version = "1.10.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e75b2483e97a5a7da73ac68a05b629f9c53cff58d8ed1c77866079e18b00dba5" +dependencies = [ + "digest 0.10.7", + "spin", +] + [[package]] name = "crc32fast" version = "1.5.2" @@ -2156,6 +2571,16 @@ dependencies = [ "memchr", ] +[[package]] +name = "ctor" +version = "1.0.13" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "914a755b7c2d4af2bdcff7ce1739e2db9a1b81a9b07123d8015786ae03c0980d" +dependencies = [ + "link-section", + "linktime-proc-macro", +] + [[package]] name = "ctutils" version = "0.4.2" @@ -3564,7 +3989,7 @@ dependencies = [ "percent-encoding", "rand 0.9.5", "serde_json", - "sha1", + "sha1 0.11.0", "sha2", "twox-hash", "url", @@ -4032,7 +4457,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0fcb34a73abd2872401096bf65d13e8d26e0c9ca6166d5e2b11a314ba017cdef" dependencies = [ "const_for", - "cpufeatures", + "cpufeatures 0.3.1", "num-traits", "pastey", "seq-macro", @@ -4517,6 +4942,18 @@ version = "0.3.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e4eba85ea1d0a966a983acd07deee566e67395d2d96b6fb39e62b5a833f1eb0b" +[[package]] +name = "gloo-timers" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bbb143cf96099802033e0d4f4963b19fd2e0b728bcf076cd9cf7f6634f092994" +dependencies = [ + "futures-channel", + "futures-core", + "js-sys", + "wasm-bindgen", +] + [[package]] name = "glow" version = "0.17.0" @@ -5964,6 +6401,8 @@ dependencies = [ "arrow-array 58.4.0", "arrow-schema 58.4.0", "async-trait", + "aws-config", + "aws-credential-types", "byteorder", "bytes", "chrono", @@ -5975,6 +6414,8 @@ dependencies = [ "log", "moka", "object_store 0.14.2", + "object_store_opendal 0.60.3", + "opendal 0.59.3", "path_abs", "pin-project", "prost 0.14.4", @@ -6340,6 +6781,18 @@ dependencies = [ "bitflags 2.13.2", ] +[[package]] +name = "link-section" +version = "0.19.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "39c29a617ce3df32c08497bdc1ab6e2376e0b17948ac166a2fbe5977c5954cd9" + +[[package]] +name = "linktime-proc-macro" +version = "0.2.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7e57c38c1e860fd37c604281cdfb1dd2216977fd76a50f85ba2f388ef3219616" + [[package]] name = "linux-raw-sys" version = "0.12.1" @@ -7104,17 +7557,31 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f1796bc93603f78c5760a69f2d58badc9618d22adade0a95385bb2adbae4eb94" dependencies = [ "async-trait", + "aws-lc-rs", + "base64 0.23.1", "bytes", "chrono", + "crc-fast", + "form_urlencoded", "futures-channel", "futures-core", "futures-util", "http", + "http-body-util", "humantime", + "hyper", "itertools 0.15.0", + "md-5 0.11.0", "nix", "parking_lot", "percent-encoding", + "quick-xml 0.41.0", + "rand 0.10.3", + "reqwest 0.13.5", + "rustls-pki-types", + "serde", + "serde_json", + "serde_urlencoded", "thiserror 2.0.21", "tokio", "tracing", @@ -7137,7 +7604,24 @@ dependencies = [ "futures", "mea", "object_store 0.13.2", - "opendal", + "opendal 0.58.2", + "pin-project", + "tokio", +] + +[[package]] +name = "object_store_opendal" +version = "0.60.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "009a43e88f0f7b3b3c2c9de5384a0ac2cef8f33996b16f9e6738e24661618775" +dependencies = [ + "async-trait", + "asyncband 0.7.2", + "bytes", + "chrono", + "futures", + "object_store 0.14.2", + "opendal 0.59.3", "pin-project", "tokio", ] @@ -7179,12 +7663,28 @@ version = "0.58.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "33dbff14cc9bb085224256d6a81289d2f3202e85b06f408d42534b42162a4231" dependencies = [ - "opendal-core", + "opendal-core 0.58.2", "opendal-service-cos", "opendal-service-goosefs", "opendal-service-oss", ] +[[package]] +name = "opendal" +version = "0.59.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ef16fe6b19b73b49ce39e3e0f4cec27e54789eee5d5f752cb2d26bd3e9449686" +dependencies = [ + "ctor", + "opendal-core 0.59.3", + "opendal-http-transport-reqwest", + "opendal-layer-concurrent-limit", + "opendal-layer-logging", + "opendal-layer-retry", + "opendal-layer-timeout", + "opendal-service-s3", +] + [[package]] name = "opendal-core" version = "0.58.2" @@ -7211,6 +7711,89 @@ dependencies = [ "web-time", ] +[[package]] +name = "opendal-core" +version = "0.59.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ddff5c648dcb3922c2bec71e0d26e881ddf9c367cb7c9a59ed7f2d39a4f38218" +dependencies = [ + "anyhow", + "asyncband 0.7.2", + "base64 0.23.1", + "bytes", + "futures", + "http", + "jiff", + "log", + "md-5 0.11.0", + "percent-encoding", + "quick-xml 0.41.0", + "reqsign-core", + "serde", + "serde_json", + "tokio", + "url", + "uuid", + "web-time", +] + +[[package]] +name = "opendal-http-transport-reqwest" +version = "0.59.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "56fcee7ac4b4b815d38346ceb4f3a5d7ef18869c369e000cc7f622b51bcfa938" +dependencies = [ + "bytes", + "futures", + "http", + "http-body", + "opendal-core 0.59.3", + "reqwest 0.13.5", +] + +[[package]] +name = "opendal-layer-concurrent-limit" +version = "0.59.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b454e68e5f6b331bd079a2ee02ebe05078fc0ddaae6822767ef54ab3c2ac0aa8" +dependencies = [ + "asyncband 0.7.2", + "futures", + "http", + "opendal-core 0.59.3", +] + +[[package]] +name = "opendal-layer-logging" +version = "0.59.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e6fd30cc4be601d5adf4936f3bcc186be5e09046c85b4e0f4e6c47ac10966335" +dependencies = [ + "log", + "opendal-core 0.59.3", +] + +[[package]] +name = "opendal-layer-retry" +version = "0.59.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1b2da6d89ef0392177f0410dead4b0f50ae7eed56033e00e6c5d78d618b49194" +dependencies = [ + "backon", + "log", + "opendal-core 0.59.3", +] + +[[package]] +name = "opendal-layer-timeout" +version = "0.59.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "337e46d134eb08beafea92fbe32fdb96d80750b711fea21b158ed3f5cb7ec40c" +dependencies = [ + "opendal-core 0.59.3", + "tokio", +] + [[package]] name = "opendal-service-cos" version = "0.58.2" @@ -7220,7 +7803,7 @@ dependencies = [ "bytes", "http", "log", - "opendal-core", + "opendal-core 0.58.2", "quick-xml 0.41.0", "reqsign-core", "reqsign-file-read-tokio", @@ -7237,7 +7820,7 @@ dependencies = [ "bytes", "goosefs-sdk", "log", - "opendal-core", + "opendal-core 0.58.2", "serde", "tokio", ] @@ -7251,7 +7834,7 @@ dependencies = [ "bytes", "http", "log", - "opendal-core", + "opendal-core 0.58.2", "quick-xml 0.41.0", "reqsign-aliyun-oss", "reqsign-core", @@ -7259,6 +7842,27 @@ dependencies = [ "serde", ] +[[package]] +name = "opendal-service-s3" +version = "0.59.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9de51cba3b83c8130e211045fd2b9d454c4b73f557ad88b71244c2ab943b2b05" +dependencies = [ + "base64 0.23.1", + "bytes", + "crc-fast", + "http", + "log", + "md-5 0.11.0", + "opendal-core 0.59.3", + "quick-xml 0.41.0", + "reqsign-aws-v4", + "reqsign-core", + "reqsign-file-read-tokio", + "serde", + "url", +] + [[package]] name = "openssl-probe" version = "0.2.1" @@ -7369,6 +7973,12 @@ dependencies = [ "hashbrown 0.14.5", ] +[[package]] +name = "outref" +version = "0.5.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1a80800c0488c3a21695ea981a54918fbb37abf04f4d0720c453632255e2ff0e" + [[package]] name = "owo-colors" version = "4.4.0" @@ -7739,6 +8349,12 @@ version = "0.2.17" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a89322df9ebe1c1578d689c92318e070967d1042b512afbe49518723f4e6d5cd" +[[package]] +name = "pin-utils" +version = "0.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8b870d8c151b6f2fb93e84a13146138f05d02ed11c7e7c54f8826aaaf7c9f184" + [[package]] name = "ping" version = "0.7.1" @@ -8161,6 +8777,16 @@ dependencies = [ "serde", ] +[[package]] +name = "quick-xml" +version = "0.42.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "41b1177fdf999d2321d3fb46ff47159d9c1fb9ad66a4879f8c50a0b504615e9b" +dependencies = [ + "memchr", + "serde", +] + [[package]] name = "quick_cache" version = "0.6.24" @@ -8383,6 +9009,7 @@ dependencies = [ "rand_distr 0.6.0", "tabled", "tokio", + "url", "vortex-bench", ] @@ -8618,6 +9245,42 @@ dependencies = [ "serde_json", ] +[[package]] +name = "reqsign-aws-core" +version = "3.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d38c902d8ec745bd9eadd1abbcbe227f4c1fc2926b2ec8929f89234af3114500" +dependencies = [ + "bytes", + "form_urlencoded", + "hex", + "http", + "log", + "percent-encoding", + "quick-xml 0.42.0", + "reqsign-core", + "rust-ini", + "serde", + "serde_json", + "serde_urlencoded", + "sha1 0.11.0", +] + +[[package]] +name = "reqsign-aws-v4" +version = "3.3.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "906bfd028f96bfe73fdb892ce26f3b3f99b4f30feb9bd024be58bd6657f68b2b" +dependencies = [ + "bytes", + "http", + "log", + "quick-xml 0.42.0", + "reqsign-aws-core", + "reqsign-core", + "serde", +] + [[package]] name = "reqsign-core" version = "3.3.2" @@ -8634,7 +9297,7 @@ dependencies = [ "jiff", "log", "percent-encoding", - "sha1", + "sha1 0.11.0", "sha2", "windows-sys 0.61.2", ] @@ -9227,6 +9890,17 @@ dependencies = [ "syn 3.0.6", ] +[[package]] +name = "sha1" +version = "0.10.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a978451301f4db1d02937a4ab3ccce137717b81826e79b7d49ffe3244a13c3b8" +dependencies = [ + "cfg-if", + "cpufeatures 0.2.17", + "digest 0.10.7", +] + [[package]] name = "sha1" version = "0.11.0" @@ -9234,7 +9908,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "aacc4cc499359472b4abe1bf11d0b12e688af9a805fa5e3016f9a386dc2d0214" dependencies = [ "cfg-if", - "cpufeatures", + "cpufeatures 0.3.1", "digest 0.11.3", ] @@ -9245,7 +9919,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "446ba717509524cb3f22f17ecc096f10f4822d76ab5c0b9822c5f9c284e825f4" dependencies = [ "cfg-if", - "cpufeatures", + "cpufeatures 0.3.1", "digest 0.11.3", ] @@ -9475,6 +10149,12 @@ dependencies = [ "spatialbench", ] +[[package]] +name = "spin" +version = "0.10.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "023a211cb3138dbc438680b32560ad89f699977624c9f8dbb95a47d5b4c07dd3" + [[package]] name = "sqllogictest" version = "0.29.1" @@ -10495,6 +11175,12 @@ dependencies = [ "serde", ] +[[package]] +name = "urlencoding" +version = "2.1.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "daf8dba3b7eb870caf1ddeed7bc9d2a049f3cfdfae7cb521b087cc33ae4c49da" + [[package]] name = "utf8-ranges" version = "1.0.5" @@ -10754,6 +11440,7 @@ dependencies = [ "mimalloc", "noodles-bgzf", "noodles-vcf", + "object_store 0.13.2", "parking_lot", "parquet 56.2.1", "parquet 59.3.0", @@ -10867,8 +11554,8 @@ version = "0.1.0" dependencies = [ "http", "object_store 0.13.2", - "object_store_opendal", - "opendal", + "object_store_opendal 0.58.0", + "opendal 0.58.2", "parking_lot", "percent-encoding", "rstest", @@ -11865,6 +12552,12 @@ dependencies = [ "zstd 0.13.3", ] +[[package]] +name = "vsimd" +version = "0.8.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5c3082ca00d5a5ef149bb8b555a72ae84c9c59f7250f013ac822ac2e49b19c64" + [[package]] name = "walkdir" version = "2.5.0" @@ -12354,6 +13047,12 @@ dependencies = [ "rustix", ] +[[package]] +name = "xmlparser" +version = "0.13.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "66fee0b777b0f5ac1c69bb06d361268faafa61cd4682ae064a171c16c433e9e4" + [[package]] name = "xtask" version = "0.1.0" @@ -12516,7 +13215,7 @@ dependencies = [ "memchr", "pbkdf2", "ppmd-rust", - "sha1", + "sha1 0.11.0", "time", "typed-path", "zeroize", diff --git a/benchmarks/lance-bench/Cargo.toml b/benchmarks/lance-bench/Cargo.toml index 816cebaffa5..b7ce37c2f14 100644 --- a/benchmarks/lance-bench/Cargo.toml +++ b/benchmarks/lance-bench/Cargo.toml @@ -15,7 +15,8 @@ version.workspace = true publish = false [dependencies] -lance = { version = "12", default-features = false } +# `aws` registers Lance's s3:// object store provider, needed to benchmark against remote data. +lance = { version = "12", default-features = false, features = ["aws"] } lance-file = { version = "12", default-features = false } anyhow = { workspace = true } diff --git a/benchmarks/lance-bench/src/random_access.rs b/benchmarks/lance-bench/src/random_access.rs index f8d56b23f62..638ff517a83 100644 --- a/benchmarks/lance-bench/src/random_access.rs +++ b/benchmarks/lance-bench/src/random_access.rs @@ -89,11 +89,17 @@ pub struct LanceRandomAccessor { impl LanceRandomAccessor { /// Open a Lance dataset and return a ready-to-use accessor. pub async fn open(path: PathBuf, name: impl Into) -> anyhow::Result { - let dataset = Dataset::open( + Self::open_uri( path.to_str() .ok_or_else(|| anyhow!("Invalid dataset path"))?, + name, ) - .await?; + .await + } + + /// Open a Lance dataset from any URI (local path, `s3://...`, ...). + pub async fn open_uri(uri: &str, name: impl Into) -> anyhow::Result { + let dataset = Dataset::open(uri).await?; let projection = ProjectionRequest::from_schema(dataset.schema().clone()); Ok(Self { name: name.into(), diff --git a/benchmarks/random-access-bench/Cargo.toml b/benchmarks/random-access-bench/Cargo.toml index 1aeefd19116..9a356fbc1bf 100644 --- a/benchmarks/random-access-bench/Cargo.toml +++ b/benchmarks/random-access-bench/Cargo.toml @@ -23,6 +23,7 @@ rand = { workspace = true } rand_distr = { workspace = true } tabled = { workspace = true } tokio = { workspace = true, features = ["full"] } +url = { workspace = true } vortex-bench = { workspace = true } [features] diff --git a/benchmarks/random-access-bench/README.md b/benchmarks/random-access-bench/README.md index aba142b20c2..73f52f5b219 100644 --- a/benchmarks/random-access-bench/README.md +++ b/benchmarks/random-access-bench/README.md @@ -27,3 +27,29 @@ work in each timed iteration. CI drives the full matrix via ```bash cargo run -p random-access-bench --profile release_debug --features lance ``` + +## Running against S3 + +The same benchmark can read its data from an object store instead of local disk. The remote +directory must mirror the layout of the local data directory (`vortex-bench/data/`), so the +files are materialized locally first and then uploaded verbatim: + +```bash +cargo run -p random-access-bench --profile release_debug --features lance -- \ + --prepare-data --formats parquet,vortex,lance +aws s3 cp --recursive vortex-bench/data s3://my-bucket/my-prefix/ + +cargo run -p random-access-bench --profile release_debug --features lance -- \ + --remote-data-dir s3://my-bucket/my-prefix/ +``` + +Credentials and region come from the environment (`AWS_REGION`, `AWS_PROFILE`, ...). + +Remote measurements are named `...-tokio-s3` instead of `...-tokio-local-disk` and are +reported with `s3` storage, so they form a series separate from the local-disk numbers. Their +ingest records use an `-s3` dataset suffix, such as `taxi-s3/uniform`, for the same reason. Arrow +IPC has no object store reader and is skipped for remote runs. In CI +the variant runs from +[`pr-bench-random-access-s3.yml`](../../.github/workflows/pr-bench-random-access-s3.yml) +(label `action/bench-random-access-s3`) and from the `Random Access (S3)` matrix entry in +`develop-bench.yml`. diff --git a/benchmarks/random-access-bench/src/lib.rs b/benchmarks/random-access-bench/src/lib.rs index f06777c01da..357e688ed53 100644 --- a/benchmarks/random-access-bench/src/lib.rs +++ b/benchmarks/random-access-bench/src/lib.rs @@ -24,8 +24,10 @@ use vortex_bench::random_access::ArrowIpcRandomAccessor; use vortex_bench::random_access::BenchDataset; use vortex_bench::random_access::ParquetRandomAccessor; use vortex_bench::random_access::RandomAccessor; +use vortex_bench::random_access::RemoteDataDir; use vortex_bench::random_access::VortexRandomAccessor; use vortex_bench::utils::constants::STORAGE_NVME; +use vortex_bench::utils::constants::STORAGE_S3; use vortex_bench::v3; use crate::render::RandomAccessRun; @@ -150,10 +152,11 @@ async fn benchmark_random_access( time_limit_secs: u64, storage: &str, reopen: bool, + remote: Option<&RemoteDataDir>, ) -> Result { let time_limit = Duration::from_secs(time_limit_secs); let mut runs = Vec::new(); - let prepared_accessor = open_accessor(dataset, format).await?; + let prepared_accessor = open_accessor(dataset, format, remote).await?; let cached_accessor = (!reopen).then_some(prepared_accessor); if let Some(accessor) = &cached_accessor { @@ -173,7 +176,10 @@ async fn benchmark_random_access( let arr = if let Some(accessor) = &cached_accessor { accessor.take(indices).await? } else { - open_accessor(dataset, format).await?.take(indices).await? + open_accessor(dataset, format, remote) + .await? + .take(indices) + .await? }; runs.push(start.elapsed()); drop(arr); @@ -214,16 +220,29 @@ fn display_name(dataset: &str, pattern: Option) -> String { /// historical continuity with existing benchmark data. /// For other datasets, includes dataset and pattern: /// `random-access/{dataset}/{pattern}/{format}-tokio-local-disk`. -fn measurement_name(dataset: &str, pattern: Option, format: Format) -> String { +/// +/// Remote runs use a `-tokio-s3` suffix so their results form a separate series from the +/// local-disk ones. +fn measurement_name( + dataset: &str, + pattern: Option, + format: Format, + remote: bool, +) -> String { let fmt = format.ext(); + let suffix = source_suffix(remote); match pattern { - Some(p) => format!( - "random-access/{}/{}/{}-tokio-local-disk", - dataset, - p.name(), - fmt - ), - None => format!("random-access/{}-tokio-local-disk", fmt), + Some(p) => format!("random-access/{}/{}/{}-{}", dataset, p.name(), fmt, suffix), + None => format!("random-access/{}-{}", fmt, suffix), + } +} + +/// Name suffix identifying where the benchmarked files are read from. +fn source_suffix(remote: bool) -> &'static str { + if remote { + "tokio-s3" + } else { + "tokio-local-disk" } } @@ -235,11 +254,50 @@ fn v3_random_access_dataset_name(dataset: &str, pattern: Option) } fn push_v3_random_access_record(records: &mut Vec, run: &RandomAccessRun) { - let dataset = v3_random_access_dataset_name(&run.dataset, run.pattern); + // v3 records have no storage field, so S3 runs get their own dataset name. Otherwise they + // would share a measurement ID with, and overwrite, the local-disk rows. + let dataset = if run.timing.storage == STORAGE_S3 { + format!("{}-s3", run.dataset) + } else { + run.dataset.clone() + }; + let dataset = v3_random_access_dataset_name(&dataset, run.pattern); let open_mode = if run.reopen { "reopen" } else { "cached" }; records.push(v3::random_access_record(&run.timing, &dataset, open_mode)); } +/// Materialize the local data file for `dataset` in `format`, writing it if it is missing. +async fn dataset_path(dataset: &dyn BenchDataset, format: Format) -> Result { + match format { + #[cfg(feature = "lance")] + Format::Lance => { + use lance_bench::random_access; + match dataset.name() { + "taxi" => random_access::taxi_data_lance().await, + "feature-vectors" => random_access::feature_vectors_lance().await, + "nested-lists" => random_access::nested_lists_lance().await, + "nested-structs" => random_access::nested_structs_lance().await, + other => anyhow::bail!("Unknown dataset for Lance: {other}"), + } + } + format => dataset.path(format).await, + } +} + +/// Materialize every local data file needed for `datasets` and `formats`. +/// +/// Remote runs read data that must already exist in the remote directory, so this is run first +/// (locally) and the resulting files uploaded verbatim. +pub async fn prepare_data(datasets: &[Box], formats: &[Format]) -> Result<()> { + for dataset in datasets { + for format in formats { + let path = dataset_path(dataset.as_ref(), *format).await?; + println!("{}", path.display()); + } + } + Ok(()) +} + /// Open a random accessor for any supported format. /// /// For Vortex and Parquet, the path comes from [`BenchDataset::path`]. @@ -247,40 +305,50 @@ fn push_v3_random_access_record(records: &mut Vec, run: &RandomAcc async fn open_accessor( dataset: &dyn BenchDataset, format: Format, + remote: Option<&RemoteDataDir>, ) -> Result> { let name = format!( - "random-access/{}/{}-tokio-local-disk", + "random-access/{}/{}-{}", dataset.name(), - format.ext() + format.ext(), + source_suffix(remote.is_some()) ); match format { Format::ArrowIpc => { + if remote.is_some() { + anyhow::bail!("Arrow IPC random access is only supported from local disk"); + } let path = dataset.path(format).await?; Ok(Box::new(ArrowIpcRandomAccessor::open(path, name)?)) } Format::OnDiskVortex | Format::VortexCompact => { - let path = dataset.path(format).await?; - Ok(Box::new( - VortexRandomAccessor::open(path, name, format).await?, - )) + let path = dataset_path(dataset, format).await?; + Ok(match remote { + Some(remote) => Box::new( + VortexRandomAccessor::open_object_store(remote, &path, name, format).await?, + ), + None => Box::new(VortexRandomAccessor::open(path, name, format).await?), + }) } Format::Parquet => { - let path = dataset.path(format).await?; - Ok(Box::new(ParquetRandomAccessor::open(path, name).await?)) + let path = dataset_path(dataset, format).await?; + Ok(match remote { + Some(remote) => { + Box::new(ParquetRandomAccessor::open_object_store(remote, &path, name).await?) + } + None => Box::new(ParquetRandomAccessor::open(path, name).await?), + }) } #[cfg(feature = "lance")] Format::Lance => { use lance_bench::random_access; - let path = match dataset.name() { - "taxi" => random_access::taxi_data_lance().await?, - "feature-vectors" => random_access::feature_vectors_lance().await?, - "nested-lists" => random_access::nested_lists_lance().await?, - "nested-structs" => random_access::nested_structs_lance().await?, - other => anyhow::bail!("Unknown dataset for Lance: {other}"), - }; - Ok(Box::new( - random_access::LanceRandomAccessor::open(path, name).await?, - )) + let path = dataset_path(dataset, format).await?; + Ok(match remote { + Some(remote) => Box::new( + random_access::LanceRandomAccessor::open_uri(&remote.uri(&path)?, name).await?, + ), + None => Box::new(random_access::LanceRandomAccessor::open(path, name).await?), + }) } other => anyhow::bail!("open_accessor not implemented for {other}"), } @@ -313,6 +381,8 @@ pub struct RunConfig { pub output_path: Option, /// Optional path for benchmark ingest JSONL records. pub ingest_output: Option, + /// When set, read the data files from this remote directory instead of local disk. + pub remote_data_dir: Option, } /// Run random-access benchmarks with `config`. @@ -326,8 +396,16 @@ pub async fn run(config: RunConfig) -> Result<()> { display_format, output_path, ingest_output, + remote_data_dir, } = config; + let remote = remote_data_dir.as_ref(); + let storage = if remote.is_some() { + STORAGE_S3 + } else { + STORAGE_NVME + }; + let reopen_variants: &[bool] = match open_mode { OpenMode::Cached => &[false], OpenMode::Reopen => &[true], @@ -351,7 +429,7 @@ pub async fn run(config: RunConfig) -> Result<()> { for dataset in &datasets { for format in &formats { if dataset.name() == "taxi" { - let name = measurement_name(dataset.name(), None, *format); + let name = measurement_name(dataset.name(), None, *format, remote.is_some()); for &reopen in reopen_variants { let bench_name = if reopen { format!("{name}-footer") @@ -365,8 +443,9 @@ pub async fn run(config: RunConfig) -> Result<()> { None, &FIXED_TAXI_INDICES, time_limit, - STORAGE_NVME, + storage, reopen, + remote, ) .await?; @@ -378,7 +457,8 @@ pub async fn run(config: RunConfig) -> Result<()> { for pattern in &patterns { let indices = generate_indices(dataset.as_ref(), *pattern); - let name = measurement_name(dataset.name(), Some(*pattern), *format); + let name = + measurement_name(dataset.name(), Some(*pattern), *format, remote.is_some()); for &reopen in reopen_variants { let bench_name = if reopen { format!("{name}-footer") @@ -392,8 +472,9 @@ pub async fn run(config: RunConfig) -> Result<()> { Some(*pattern), &indices, time_limit, - STORAGE_NVME, + storage, reopen, + remote, ) .await?; @@ -496,6 +577,20 @@ mod tests { } } + #[test] + fn v3_random_access_records_keep_s3_runs_apart() { + let mut run = fake_run("taxi", Some(AccessPattern::Uniform), false); + run.timing.storage = STORAGE_S3.to_string(); + let mut records = Vec::new(); + + push_v3_random_access_record(&mut records, &run); + + match &records[0] { + v3::V3Record::RandomAccessTime(record) => assert_eq!(record.dataset, "taxi-s3/uniform"), + other => panic!("expected random-access record, got {other:?}"), + } + } + #[test] fn display_name_drops_format_extension() { assert_eq!(display_name("taxi", None), "random-access/taxi"); diff --git a/benchmarks/random-access-bench/src/main.rs b/benchmarks/random-access-bench/src/main.rs index b7f41533c79..37822576e99 100644 --- a/benchmarks/random-access-bench/src/main.rs +++ b/benchmarks/random-access-bench/src/main.rs @@ -9,6 +9,7 @@ use clap::ValueEnum; use random_access_bench::AccessPattern; use random_access_bench::OpenMode; use random_access_bench::RunConfig; +use url::Url; use vortex_bench::Format; use vortex_bench::datasets::feature_vectors::FeatureVectorsData; use vortex_bench::datasets::nested_lists::NestedListsData; @@ -16,6 +17,7 @@ use vortex_bench::datasets::nested_structs::NestedStructsData; use vortex_bench::datasets::taxi_data::TaxiData; use vortex_bench::display::DisplayFormat; use vortex_bench::random_access::BenchDataset; +use vortex_bench::random_access::RemoteDataDir; use vortex_bench::setup_logging_and_tracing; /// Which synthetic dataset to benchmark. @@ -85,6 +87,15 @@ struct Args { /// Whether to reopen the file on each iteration, use a cached handle, or run both. #[arg(long, value_enum, default_value_t = OpenMode::Both)] open_mode: OpenMode, + /// Read the data files from this remote directory (e.g. `s3://bucket/prefix/`) instead of + /// local disk. The directory must mirror the layout of the local benchmark data directory, + /// as produced by `--prepare-data`. + #[arg(long)] + remote_data_dir: Option, + /// Materialize the local data files for the selected datasets and formats, print their + /// paths, and exit without benchmarking. + #[arg(long)] + prepare_data: bool, } #[tokio::main] @@ -92,12 +103,18 @@ async fn main() -> Result<()> { let args = Args::parse(); setup_logging_and_tracing(args.verbose, args.tracing)?; + let datasets: Vec> = args + .datasets + .into_iter() + .map(DatasetArg::into_dataset) + .collect(); + + if args.prepare_data { + return random_access_bench::prepare_data(&datasets, &args.formats).await; + } + let run_config = RunConfig { - datasets: args - .datasets - .into_iter() - .map(DatasetArg::into_dataset) - .collect(), + datasets, formats: args.formats, patterns: args.patterns, time_limit: args.time_limit, @@ -105,6 +122,10 @@ async fn main() -> Result<()> { display_format: args.display_format, output_path: args.output_path, ingest_output: args.ingest_output, + remote_data_dir: args + .remote_data_dir + .map(RemoteDataDir::try_new) + .transpose()?, }; random_access_bench::run(run_config).await diff --git a/docs/developer-guide/benchmarking.md b/docs/developer-guide/benchmarking.md index ef6f8127ca1..90eb34c48c6 100644 --- a/docs/developer-guide/benchmarking.md +++ b/docs/developer-guide/benchmarking.md @@ -257,6 +257,8 @@ Benchmarks run automatically on all commits to `develop` and can be run on-deman - **Post-commit** -- compression, string encoding, random access, and SQL benchmarks run on every commit to `develop`, with results uploaded for historical tracking. - **Random access** -- `action/bench-random-access` runs only the random-access benchmark. +- **Random access (S3)** -- `action/bench-random-access-s3` runs the random-access benchmark + against data uploaded to S3 instead of local disk. - **Compression** -- `action/bench-compress` runs only the compression benchmark. - **String encoding** -- `action/bench-string` runs only the string encoding benchmark. - **GPU compression** -- `action/bench-gpu-compress` runs the allow-listed Vortex decompression diff --git a/scripts/random-access-split.py b/scripts/random-access-split.py index 3a5a41740b7..c33cacb0967 100755 --- a/scripts/random-access-split.py +++ b/scripts/random-access-split.py @@ -35,11 +35,13 @@ def drop_os_caches() -> None: pass -def run_combinations(emit_ingest_records: bool) -> None: +def run_combinations(emit_ingest_records: bool, remote_data_dir: str | None) -> None: PARTS_DIR.mkdir(parents=True, exist_ok=True) + # Arrow IPC has no object store reader, so a remote run would silently read local disk. + formats = [f for f in FORMATS if not (remote_data_dir and f == "arrow-ipc")] i = 0 for dataset in DATASETS: - for fmt in FORMATS: + for fmt in formats: for pattern in PATTERNS: for open_mode in OPEN_MODES: drop_os_caches() @@ -61,6 +63,8 @@ def run_combinations(emit_ingest_records: bool) -> None: "-o", str(PARTS_DIR / f"{i}.gh.json"), ] + if remote_data_dir: + args += ["--remote-data-dir", remote_data_dir] if emit_ingest_records: args += ["--ingest-jsonl", str(PARTS_DIR / f"{i}.ingest.jsonl")] print("+", " ".join(args), flush=True) @@ -109,9 +113,13 @@ def main() -> None: action="store_true", help="merge --ingest-jsonl records into results.ingest.jsonl", ) + parser.add_argument( + "--remote-data-dir", + help="read the benchmark data from this remote directory (e.g. s3://bucket/prefix/)", + ) args = parser.parse_args() - run_combinations(args.emit_ingest_records) + run_combinations(args.emit_ingest_records, args.remote_data_dir) merge(f"{PARTS_DIR}/*.gh.json", lambda record: record["name"], "results.json") if args.emit_ingest_records: merge( diff --git a/vortex-bench/Cargo.toml b/vortex-bench/Cargo.toml index a63f954a637..97681a7a138 100644 --- a/vortex-bench/Cargo.toml +++ b/vortex-bench/Cargo.toml @@ -53,8 +53,9 @@ itertools = { workspace = true } mimalloc = { workspace = true } noodles-bgzf = { workspace = true, features = ["async"] } noodles-vcf = { workspace = true, features = ["async"] } +object_store = { workspace = true, features = ["aws"] } parking_lot = { workspace = true } -parquet = { workspace = true, features = ["async"] } +parquet = { workspace = true, features = ["async", "object_store"] } rand = { workspace = true } regex = { workspace = true } reqwest = { workspace = true, features = ["stream"] } diff --git a/vortex-bench/src/random_access/mod.rs b/vortex-bench/src/random_access/mod.rs index c792096b795..ac7d48ca29c 100644 --- a/vortex-bench/src/random_access/mod.rs +++ b/vortex-bench/src/random_access/mod.rs @@ -2,20 +2,27 @@ // SPDX-FileCopyrightText: Copyright the Vortex contributors use std::fs::File; +use std::path::Path; use std::path::PathBuf; use std::sync::Arc; use anyhow::Result; +use anyhow::anyhow; use arrow_array::RecordBatch; use arrow_ipc::writer::FileWriter; use async_trait::async_trait; +use object_store::ObjectStore; +use object_store::aws::AmazonS3Builder; +use object_store::path::Path as ObjectStorePath; use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder; use parquet::basic::Compression; use parquet::basic::ZstdLevel; use parquet::file::properties::WriterProperties; +use url::Url; use vortex::array::ArrayRef; use crate::Format; +use crate::data_dir; use crate::idempotent; pub mod take; @@ -85,6 +92,69 @@ pub fn random_access_writer_properties() -> Result { .build()) } +/// A remote directory holding the same layout as the local benchmark data directory. +/// +/// Random access datasets are always materialized locally first, then uploaded verbatim, so a +/// remote object key is just the local path relative to [`data_dir`] appended to the URL path. +#[derive(Clone, Debug)] +pub struct RemoteDataDir { + url: Url, + store: Arc, +} + +impl RemoteDataDir { + /// Build an object store for `url` (e.g. `s3://bucket/prefix/`) from the ambient environment. + pub fn try_new(url: Url) -> Result { + let store: Arc = match url.scheme() { + "s3" => { + let bucket = url + .host_str() + .ok_or_else(|| anyhow!("remote data dir has no bucket: {url}"))?; + Arc::new( + AmazonS3Builder::from_env() + .with_bucket_name(bucket) + .build()?, + ) + } + other => return Err(anyhow!("unsupported remote data dir scheme: {other}")), + }; + Ok(Self { url, store }) + } + + /// The object store backing this directory. + pub fn store(&self) -> &Arc { + &self.store + } + + /// The object key of `local_path`, mirroring its location under the local data directory. + pub fn key(&self, local_path: &Path) -> Result { + let relative = local_path.strip_prefix(data_dir()).map_err(|_| { + anyhow!( + "{} is not inside the benchmark data directory", + local_path.display() + ) + })?; + let relative = relative + .to_str() + .ok_or_else(|| anyhow!("non-UTF-8 data path: {}", local_path.display()))?; + // `ObjectStorePath` drops the empty segments left by leading or trailing slashes. + Ok(ObjectStorePath::from(format!( + "{}/{relative}", + self.url.path() + ))) + } + + /// The fully qualified URL of `local_path` in this remote directory. + pub fn uri(&self, local_path: &Path) -> Result { + let scheme = self.url.scheme(); + let host = self + .url + .host_str() + .ok_or_else(|| anyhow!("remote data dir has no bucket: {}", self.url))?; + Ok(format!("{scheme}://{host}/{}", self.key(local_path)?)) + } +} + /// Trait for a benchmark dataset that knows how to prepare data files. #[async_trait] pub trait BenchDataset: Send + Sync { @@ -126,6 +196,46 @@ pub trait RandomAccessor: Send + Sync { mod tests { use super::*; + fn remote(url: &str) -> Result { + // `from_env` needs no credentials to construct the client. + RemoteDataDir::try_new(Url::parse(url)?) + } + + #[test] + fn key_mirrors_the_local_data_dir_layout() -> Result<()> { + let local = data_dir().join("random_access/taxi/taxi.vortex"); + + assert_eq!( + remote("s3://bucket/prefix/")?.key(&local)?.as_ref(), + "prefix/random_access/taxi/taxi.vortex" + ); + assert_eq!( + remote("s3://bucket/")?.key(&local)?.as_ref(), + "random_access/taxi/taxi.vortex" + ); + assert_eq!( + remote("s3://bucket/prefix/")?.uri(&local)?, + "s3://bucket/prefix/random_access/taxi/taxi.vortex" + ); + Ok(()) + } + + #[test] + fn key_rejects_paths_outside_the_data_dir() -> Result<()> { + assert!( + remote("s3://bucket/prefix/")? + .key(Path::new("/tmp/taxi.vortex")) + .is_err() + ); + Ok(()) + } + + #[test] + fn unsupported_scheme_is_rejected() -> Result<()> { + assert!(RemoteDataDir::try_new(Url::parse("gs://bucket/prefix/")?).is_err()); + Ok(()) + } + #[test] fn generated_parquet_is_zstd_level_3() -> Result<()> { let props = random_access_writer_properties()?; diff --git a/vortex-bench/src/random_access/take.rs b/vortex-bench/src/random_access/take.rs index 26cd93a9a78..93a66ccd49e 100644 --- a/vortex-bench/src/random_access/take.rs +++ b/vortex-bench/src/random_access/take.rs @@ -4,6 +4,7 @@ use std::collections::BTreeMap; use std::fs::File as StdFile; use std::iter::once; +use std::path::Path; use std::path::PathBuf; use std::sync::Arc; @@ -20,6 +21,8 @@ use parking_lot::Mutex; use parquet::arrow::ParquetRecordBatchStreamBuilder; use parquet::arrow::arrow_reader::ArrowReaderMetadata; use parquet::arrow::arrow_reader::ArrowReaderOptions; +use parquet::arrow::async_reader; +use parquet::arrow::async_reader::AsyncFileReader; use parquet::file::metadata::PageIndexPolicy; use stream::StreamExt; use tokio::fs::File; @@ -38,6 +41,7 @@ use crate::SESSION; use crate::random_access::ARROW_ROW_OFFSETS_METADATA_KEY; use crate::random_access::RandomAccessor; use crate::random_access::RandomAccessorRet; +use crate::random_access::RemoteDataDir; /// Random accessor for uncompressed Arrow IPC files. pub struct ArrowIpcRandomAccessor { @@ -134,7 +138,7 @@ pub struct VortexRandomAccessor { impl VortexRandomAccessor { /// Open a Vortex file and return a ready-to-use accessor. pub async fn open( - path: impl AsRef, + path: impl AsRef, name: impl Into, format: Format, ) -> anyhow::Result { @@ -149,6 +153,25 @@ impl VortexRandomAccessor { file, }) } + + /// Open a Vortex file stored in an object store and return a ready-to-use accessor. + pub async fn open_object_store( + remote: &RemoteDataDir, + path: &Path, + name: impl Into, + format: Format, + ) -> anyhow::Result { + let file = SESSION + .open_options() + .with_layout_reader_cache() + .open_object_store(remote.store(), remote.key(path)?) + .await?; + Ok(Self { + name: name.into(), + format, + file, + }) + } } #[async_trait] @@ -188,17 +211,64 @@ pub struct ParquetRandomAccessor { row_group_offsets: Vec, /// Cached Arrow reader metadata (footer) to avoid re-parsing on each take. arrow_metadata: ArrowReaderMetadata, - /// Path to the Parquet file (for re-opening on each take). - path: PathBuf, + /// Where to re-open the file from on each take. + source: ParquetSource, +} + +/// Reader for a Parquet file held in an object store. +#[expect( + deprecated, + reason = "arrow-rs deprecated this in favour of a hand-rolled AsyncFileReader; \ + keeping it holds the measured I/O path fixed (arrow-rs#10308)" +)] +type ObjectReader = async_reader::ParquetObjectReader; + +/// Backing store of a [`ParquetRandomAccessor`]. +enum ParquetSource { + /// Path to a local Parquet file. + Local(PathBuf), + /// Reader for a Parquet file held in an object store. + Object(ObjectReader), +} + +impl ParquetSource { + /// Open a fresh reader on the file, so each take pays the open cost like a cold reader. + async fn open(&self) -> anyhow::Result> { + Ok(match self { + ParquetSource::Local(path) => Box::new(File::open(path).await?), + ParquetSource::Object(reader) => Box::new(reader.clone()), + }) + } } impl ParquetRandomAccessor { /// Open a Parquet file, parse the footer, and return a ready-to-use accessor. pub async fn open(path: PathBuf, name: impl Into) -> anyhow::Result { let mut file = File::open(&path).await?; - let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required); - let arrow_metadata = ArrowReaderMetadata::load_async(&mut file, options).await?; + let arrow_metadata = load_metadata(&mut file).await?; + Ok(Self::new(name, arrow_metadata, ParquetSource::Local(path))) + } + + /// Open a Parquet file stored in an object store and return a ready-to-use accessor. + pub async fn open_object_store( + remote: &RemoteDataDir, + path: &Path, + name: impl Into, + ) -> anyhow::Result { + let mut reader = ObjectReader::new(Arc::clone(remote.store()), remote.key(path)?); + let arrow_metadata = load_metadata(&mut reader).await?; + Ok(Self::new( + name, + arrow_metadata, + ParquetSource::Object(reader), + )) + } + fn new( + name: impl Into, + arrow_metadata: ArrowReaderMetadata, + source: ParquetSource, + ) -> Self { let row_group_offsets = once(0) .chain( arrow_metadata @@ -213,15 +283,23 @@ impl ParquetRandomAccessor { }) .collect::>(); - Ok(Self { + Self { name: name.into(), row_group_offsets, arrow_metadata, - path, - }) + source, + } } } +/// Parse the Parquet footer, including the page index, from any async reader. +async fn load_metadata( + reader: &mut T, +) -> anyhow::Result { + let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required); + Ok(ArrowReaderMetadata::load_async(reader, options).await?) +} + #[async_trait] impl RandomAccessor for ParquetRandomAccessor { fn format(&self) -> Format { @@ -253,7 +331,7 @@ impl RandomAccessor for ParquetRandomAccessor { .collect_vec(); // Re-open the file but reuse cached metadata (avoids re-parsing the footer). - let file = File::open(&self.path).await?; + let file = self.source.open().await?; let builder = ParquetRecordBatchStreamBuilder::new_with_metadata(file, self.arrow_metadata.clone());