From 295753733aa1de34498313fe4c70f03e912fca8e Mon Sep 17 00:00:00 2001 From: Thomas Peiselt Date: Fri, 20 Feb 2026 11:24:17 +0100 Subject: [PATCH 1/7] Working build --- .github/actions/setup-builder/action.yaml | 2 +- .github/workflows/arrow.yml | 19 --- .github/workflows/arrow_flight.yml | 14 -- .github/workflows/dev.yml | 61 --------- .github/workflows/docs.yml | 99 -------------- .github/workflows/integration.yml | 159 ---------------------- .github/workflows/miri.sh | 20 --- .github/workflows/miri.yaml | 62 --------- .github/workflows/parquet.yml | 55 -------- .github/workflows/rust.yml | 111 ++++----------- Cargo.toml | 3 +- arrow-array/Cargo.toml | 2 +- arrow-avro/Cargo.toml | 2 +- arrow-cast/Cargo.toml | 6 +- arrow-data/Cargo.toml | 2 +- arrow-flight/Cargo.toml | 6 +- arrow-ipc/Cargo.toml | 2 +- arrow-json/Cargo.toml | 4 +- arrow-ord/Cargo.toml | 2 +- arrow-row/Cargo.toml | 2 +- arrow/Cargo.toml | 4 +- parquet/Cargo.toml | 6 +- 22 files changed, 50 insertions(+), 593 deletions(-) delete mode 100644 .github/workflows/dev.yml delete mode 100644 .github/workflows/docs.yml delete mode 100644 .github/workflows/integration.yml delete mode 100755 .github/workflows/miri.sh delete mode 100644 .github/workflows/miri.yaml diff --git a/.github/actions/setup-builder/action.yaml b/.github/actions/setup-builder/action.yaml index 20da777ec0e5..a6274aae5683 100644 --- a/.github/actions/setup-builder/action.yaml +++ b/.github/actions/setup-builder/action.yaml @@ -21,7 +21,7 @@ inputs: rust-version: description: 'version of rust to install (e.g. stable)' required: false - default: 'stable' + default: '1.86.0' target: description: 'target architecture(s)' required: false diff --git a/.github/workflows/arrow.yml b/.github/workflows/arrow.yml index 0b90a78577e5..d8a27988d525 100644 --- a/.github/workflows/arrow.yml +++ b/.github/workflows/arrow.yml @@ -133,25 +133,6 @@ jobs: run: cargo check -p arrow --no-default-features --all-targets --features chrono-tz - # test the arrow crate builds against wasm32 in nightly rust - wasm32-build: - name: Build wasm32 - runs-on: ubuntu-latest - container: - image: amd64/rust - steps: - - uses: actions/checkout@v4 - with: - submodules: true - - name: Setup Rust toolchain - uses: ./.github/actions/setup-builder - with: - target: wasm32-unknown-unknown,wasm32-wasip1 - - name: Build wasm32-unknown-unknown - run: cargo build -p arrow --no-default-features --features=json,csv,ipc,ffi --target wasm32-unknown-unknown - - name: Build wasm32-wasip1 - run: cargo build -p arrow --no-default-features --features=json,csv,ipc,ffi --target wasm32-wasip1 - clippy: name: Clippy runs-on: ubuntu-latest diff --git a/.github/workflows/arrow_flight.yml b/.github/workflows/arrow_flight.yml index 79627448ca40..b50366b4f0b9 100644 --- a/.github/workflows/arrow_flight.yml +++ b/.github/workflows/arrow_flight.yml @@ -62,20 +62,6 @@ jobs: run: | cargo test -p arrow-flight --features=flight-sql-experimental,tls --examples - vendor: - name: Verify Vendored Code - runs-on: ubuntu-latest - container: - image: amd64/rust - steps: - - uses: actions/checkout@v4 - - name: Setup Rust toolchain - uses: ./.github/actions/setup-builder - - name: Run gen - run: ./arrow-flight/regen.sh - - name: Verify workspace clean (if this fails, run ./arrow-flight/regen.sh and check in results) - run: git diff --exit-code - clippy: name: Clippy runs-on: ubuntu-latest diff --git a/.github/workflows/dev.yml b/.github/workflows/dev.yml deleted file mode 100644 index b28e8c20cfe7..000000000000 --- a/.github/workflows/dev.yml +++ /dev/null @@ -1,61 +0,0 @@ -# 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. - -name: dev - -concurrency: - group: ${{ github.repository }}-${{ github.head_ref || github.sha }}-${{ github.workflow }} - cancel-in-progress: true - -# trigger for all PRs and changes to main -on: - push: - branches: - - main - pull_request: - -env: - ARCHERY_DOCKER_USER: ${{ secrets.DOCKERHUB_USER }} - ARCHERY_DOCKER_PASSWORD: ${{ secrets.DOCKERHUB_TOKEN }} - -jobs: - - rat: - name: Release Audit Tool (RAT) - runs-on: ubuntu-latest - steps: - - uses: actions/checkout@v4 - - name: Setup Python - uses: actions/setup-python@v5 - with: - python-version: 3.8 - - name: Audit licenses - run: ./dev/release/run-rat.sh . - - prettier: - name: Markdown format - runs-on: ubuntu-latest - steps: - - uses: actions/checkout@v4 - - uses: actions/setup-node@v4 - with: - node-version: "14" - - name: Prettier check - run: | - # if you encounter error, run the command below and commit the changes - npx prettier@2.3.2 --write {arrow,arrow-flight,dev,arrow-integration-testing,parquet}/**/*.md README.md CODE_OF_CONDUCT.md CONTRIBUTING.md - git diff --exit-code diff --git a/.github/workflows/docs.yml b/.github/workflows/docs.yml deleted file mode 100644 index d6ec0622f6ed..000000000000 --- a/.github/workflows/docs.yml +++ /dev/null @@ -1,99 +0,0 @@ -# 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. - -name: docs - -concurrency: - group: ${{ github.repository }}-${{ github.head_ref || github.sha }}-${{ github.workflow }} - cancel-in-progress: true - -# trigger for all PRs and changes to main -on: - push: - branches: - - main - pull_request: - -jobs: - - # test doc links still work - docs: - name: Rustdocs are clean - runs-on: ubuntu-latest - strategy: - matrix: - arch: [ amd64 ] - rust: [ nightly ] - container: - image: ${{ matrix.arch }}/rust - env: - RUSTDOCFLAGS: "-Dwarnings --enable-index-page -Zunstable-options" - steps: - - uses: actions/checkout@v4 - with: - submodules: true - - name: Install python dev - run: | - apt update - apt install -y libpython3.11-dev - - name: Setup Rust toolchain - uses: ./.github/actions/setup-builder - with: - rust-version: ${{ matrix.rust }} - - name: Run cargo doc - run: cargo doc --document-private-items --no-deps --workspace --all-features - - name: Fix file permissions - shell: sh - run: | - chmod -c -R +rX "target/doc" | - while read line; do - echo "::warning title=Invalid file permissions automatically fixed::$line" - done - - name: Upload artifacts - uses: actions/upload-pages-artifact@v3 - with: - name: crate-docs - path: target/doc - - deploy: - # Only deploy if a push to main - if: github.ref_name == 'main' && github.event_name == 'push' - needs: docs - permissions: - contents: write - runs-on: ubuntu-latest - steps: - - uses: actions/checkout@v4 - - name: Download crate docs - uses: actions/download-artifact@v4 - with: - name: crate-docs - path: website/build - - name: Prepare website - run: | - tar -xf website/build/artifact.tar -C website/build - rm website/build/artifact.tar - cp .asf.yaml ./website/build/.asf.yaml - - name: Deploy to gh-pages - uses: peaceiris/actions-gh-pages@v4.0.0 - if: github.event_name == 'push' && github.ref_name == 'main' - with: - github_token: ${{ secrets.GITHUB_TOKEN }} - publish_dir: website/build - publish_branch: asf-site - # Avoid accumulating history of in progress API jobs: https://github.com/apache/arrow-rs/issues/5908 - force_orphan: true diff --git a/.github/workflows/integration.yml b/.github/workflows/integration.yml deleted file mode 100644 index df73d635ee90..000000000000 --- a/.github/workflows/integration.yml +++ /dev/null @@ -1,159 +0,0 @@ -# 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. - -name: integration - -concurrency: - group: ${{ github.repository }}-${{ github.head_ref || github.sha }}-${{ github.workflow }} - cancel-in-progress: true - -# trigger for all PRs that touch certain files and changes to main -on: - push: - branches: - - main - pull_request: - paths: - - .github/** - - arrow-array/** - - arrow-buffer/** - - arrow-cast/** - - arrow-csv/** - - arrow-data/** - - arrow-integration-test/** - - arrow-integration-testing/** - - arrow-ipc/** - - arrow-json/** - - arrow-avro/** - - arrow-ord/** - - arrow-pyarrow-integration-testing/** - - arrow-schema/** - - arrow-select/** - - arrow-sort/** - - arrow-string/** - - arrow/** - -jobs: - integration: - name: Archery test With other arrows - runs-on: ubuntu-latest - container: - image: apache/arrow-dev:amd64-conda-integration - env: - ARROW_USE_CCACHE: OFF - ARROW_CPP_EXE_PATH: /build/cpp/debug - ARROW_NANOARROW_PATH: /build/nanoarrow - ARROW_RUST_EXE_PATH: /build/rust/debug - BUILD_DOCS_CPP: OFF - ARROW_INTEGRATION_CPP: ON - ARROW_INTEGRATION_CSHARP: ON - ARROW_INTEGRATION_GO: ON - ARROW_INTEGRATION_JAVA: ON - ARROW_INTEGRATION_JS: ON - ARCHERY_INTEGRATION_TARGET_IMPLEMENTATIONS: "rust" - ARCHERY_INTEGRATION_WITH_NANOARROW: "1" - # https://github.com/apache/arrow/pull/38403/files#r1371281630 - ARCHERY_INTEGRATION_WITH_RUST: "1" - # These are necessary because the github runner overrides $HOME - # https://github.com/actions/runner/issues/863 - RUSTUP_HOME: /root/.rustup - CARGO_HOME: /root/.cargo - defaults: - run: - shell: bash - steps: - # This is necessary so that actions/checkout can find git - - name: Export conda path - run: echo "/opt/conda/envs/arrow/bin" >> $GITHUB_PATH - # This is necessary so that Rust can find cargo - - name: Export cargo path - run: echo "/root/.cargo/bin" >> $GITHUB_PATH - - name: Check rustup - run: which rustup - - name: Check cmake - run: which cmake - - name: Checkout Arrow - uses: actions/checkout@v4 - with: - repository: apache/arrow - submodules: true - fetch-depth: 0 - - name: Checkout Arrow Rust - uses: actions/checkout@v4 - with: - path: rust - fetch-depth: 0 - - name: Checkout Arrow nanoarrow - uses: actions/checkout@v4 - with: - repository: apache/arrow-nanoarrow - path: nanoarrow - fetch-depth: 0 - - name: Build - run: conda run --no-capture-output ci/scripts/integration_arrow_build.sh $PWD /build - - name: Run - run: conda run --no-capture-output ci/scripts/integration_arrow.sh $PWD /build - - # test FFI against the C-Data interface exposed by pyarrow - pyarrow-integration-test: - name: Pyarrow C Data Interface - runs-on: ubuntu-latest - strategy: - matrix: - rust: [stable] - # PyArrow 15 was the first version to introduce StringView/BinaryView support - pyarrow: ["15", "16", "17"] - steps: - - uses: actions/checkout@v4 - with: - submodules: true - - name: Setup Rust toolchain - run: | - rustup toolchain install ${{ matrix.rust }} - rustup default ${{ matrix.rust }} - rustup component add rustfmt clippy - - name: Cache Cargo - uses: actions/cache@v4 - with: - path: /home/runner/.cargo - key: cargo-maturin-cache- - - name: Cache Rust dependencies - uses: actions/cache@v4 - with: - path: /home/runner/target - # this key is not equal because maturin uses different compilation flags. - key: ${{ runner.os }}-${{ matrix.arch }}-target-maturin-cache-${{ matrix.rust }}- - - uses: actions/setup-python@v5 - with: - python-version: '3.8' - - name: Upgrade pip and setuptools - run: pip install --upgrade pip setuptools wheel virtualenv - - name: Create virtualenv and install dependencies - run: | - virtualenv venv - source venv/bin/activate - pip install maturin toml pytest pytz pyarrow==${{ matrix.pyarrow }} - - name: Run Rust tests - run: | - source venv/bin/activate - cargo test -p arrow --test pyarrow --features pyarrow - - name: Run tests - run: | - source venv/bin/activate - cd arrow-pyarrow-integration-testing - maturin develop - pytest -v . diff --git a/.github/workflows/miri.sh b/.github/workflows/miri.sh deleted file mode 100755 index 86be2100ee67..000000000000 --- a/.github/workflows/miri.sh +++ /dev/null @@ -1,20 +0,0 @@ -#!/bin/bash -# -# Script -# -# Must be run with nightly rust for example -# rustup default nightly - -set -e - -export MIRIFLAGS="-Zmiri-disable-isolation" -cargo miri setup -cargo clean - -echo "Starting Arrow MIRI run..." -cargo miri test -p arrow-buffer -cargo miri test -p arrow-data --features ffi -cargo miri test -p arrow-schema --features ffi -cargo miri test -p arrow-ord -cargo miri test -p arrow-array -cargo miri test -p arrow-arith \ No newline at end of file diff --git a/.github/workflows/miri.yaml b/.github/workflows/miri.yaml deleted file mode 100644 index ce67546a104b..000000000000 --- a/.github/workflows/miri.yaml +++ /dev/null @@ -1,62 +0,0 @@ -# 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. - -name: miri - -concurrency: - group: ${{ github.repository }}-${{ github.head_ref || github.sha }}-${{ github.workflow }} - cancel-in-progress: true - -# trigger for all PRs that touch certain files and changes to main -on: - push: - branches: - - main - pull_request: - paths: - - .github/** - - arrow-array/** - - arrow-buffer/** - - arrow-cast/** - - arrow-csv/** - - arrow-data/** - - arrow-ipc/** - - arrow-json/** - - arrow-avro/** - - arrow-schema/** - - arrow-select/** - - arrow-string/** - - arrow/** - -jobs: - miri-checks: - name: MIRI - runs-on: ubuntu-latest - steps: - - uses: actions/checkout@v4 - with: - submodules: true - - name: Setup Rust toolchain - run: | - rustup toolchain install nightly --component miri - rustup override set nightly - cargo miri setup - - name: Run Miri Checks - env: - RUST_BACKTRACE: full - RUST_LOG: "trace" - run: bash .github/workflows/miri.sh diff --git a/.github/workflows/parquet.yml b/.github/workflows/parquet.yml index 96c7ab8f4e3a..31a3c684281f 100644 --- a/.github/workflows/parquet.yml +++ b/.github/workflows/parquet.yml @@ -114,61 +114,6 @@ jobs: - name: Check compilation --no-default-features --features encryption --features async run: cargo check -p parquet --no-default-features --features encryption --features async - # test the parquet crate builds against wasm32 in stable rust - wasm32-build: - name: Build wasm32 - runs-on: ubuntu-latest - container: - image: amd64/rust - steps: - - uses: actions/checkout@v4 - with: - submodules: true - - name: Setup Rust toolchain - uses: ./.github/actions/setup-builder - with: - target: wasm32-unknown-unknown,wasm32-wasip1 - - name: Install clang # Needed for zlib compilation - run: apt-get update && apt-get install -y clang gcc-multilib - - name: Build wasm32-unknown-unknown - run: cargo build -p parquet --target wasm32-unknown-unknown - - name: Build wasm32-wasip1 - run: cargo build -p parquet --target wasm32-wasip1 - - pyspark-integration-test: - name: PySpark Integration Test - runs-on: ubuntu-latest - strategy: - matrix: - rust: [ stable ] - steps: - - uses: actions/checkout@v4 - - name: Setup Python - uses: actions/setup-python@v5 - with: - python-version: "3.10" - cache: "pip" - - name: Install Python dependencies - run: | - cd parquet/pytest - pip install -r requirements.txt - - name: Black check the test files - run: | - cd parquet/pytest - black --check *.py --verbose - - name: Setup Rust toolchain - run: | - rustup toolchain install ${{ matrix.rust }} - rustup default ${{ matrix.rust }} - - name: Install binary for checking - run: | - cargo install --path parquet --bin parquet-show-bloom-filter --features=cli - cargo install --path parquet --bin parquet-fromcsv --features=arrow,cli - - name: Run pytest - run: | - cd parquet/pytest - pytest -v - clippy: name: Clippy runs-on: ubuntu-latest diff --git a/.github/workflows/rust.yml b/.github/workflows/rust.yml index 80fef2674aae..271ccafc80e6 100644 --- a/.github/workflows/rust.yml +++ b/.github/workflows/rust.yml @@ -31,61 +31,6 @@ on: jobs: - # Check workspace wide compile and test with default features for - # mac - macos: - name: Test on Mac - runs-on: macos-latest - steps: - - uses: actions/checkout@v4 - with: - submodules: true - - name: Install protoc with brew - run: brew install protobuf - - name: Setup Rust toolchain - run: | - rustup toolchain install stable --no-self-update - rustup default stable - - name: Run tests - shell: bash - run: | - # do not produce debug symbols to keep memory usage down - export RUSTFLAGS="-C debuginfo=0" - cargo test - - - # Check workspace wide compile and test with default features for - # windows - windows: - name: Test on Windows - runs-on: windows-latest - steps: - - uses: actions/checkout@v4 - with: - submodules: true - - name: Install protobuf compiler in /d/protoc - shell: bash - run: | - mkdir /d/protoc - cd /d/protoc - curl -LO https://github.com/protocolbuffers/protobuf/releases/download/v21.4/protoc-21.4-win64.zip - unzip protoc-21.4-win64.zip - export PATH=$PATH:/d/protoc/bin - protoc --version - - - name: Setup Rust toolchain - run: | - rustup toolchain install stable --no-self-update - rustup default stable - - name: Run tests - shell: bash - run: | - # do not produce debug symbols to keep memory usage down - export RUSTFLAGS="-C debuginfo=0" - export PATH=$PATH:/d/protoc/bin - cargo test - - # Run cargo fmt for all crates lint: name: Lint (cargo fmt) @@ -109,31 +54,31 @@ jobs: # cargo fmt -p parquet -- --config skip_children=true `find . -name "*.rs" \! -name format.rs` cargo fmt -p parquet -- --check --config skip_children=true `find . -name "*.rs" \! -name format.rs` - msrv: - name: Verify MSRV (Minimum Supported Rust Version) - runs-on: ubuntu-latest - container: - image: amd64/rust - steps: - - uses: actions/checkout@v4 - - name: Setup Rust toolchain - uses: ./.github/actions/setup-builder - - name: Install cargo-msrv - run: cargo install cargo-msrv - - name: Downgrade arrow-pyarrow-integration-testing dependencies - working-directory: arrow-pyarrow-integration-testing - # Necessary because half 2.5 requires rust 1.81 or newer - run: | - cargo update -p half --precise 2.4.0 - - name: Downgrade workspace dependencies - # Necessary because half 2.5 requires rust 1.81 or newer - run: | - cargo update -p half --precise 2.4.0 - - name: Check all packages - run: | - # run `cargo msrv verify --manifest-path "path/to/Cargo.toml"` to see problematic dependencies - find . -mindepth 2 -name Cargo.toml | while read -r dir - do - echo "Checking package '$dir'" - cargo msrv verify --manifest-path "$dir" --output-format=json || exit 1 - done + # msrv: + # name: Verify MSRV (Minimum Supported Rust Version) + # runs-on: ubuntu-latest + # container: + # image: amd64/rust + # steps: + # - uses: actions/checkout@v4 + # - name: Setup Rust toolchain + # uses: ./.github/actions/setup-builder + # - name: Install cargo-msrv + # run: cargo install cargo-msrv + # - name: Downgrade arrow-pyarrow-integration-testing dependencies + # working-directory: arrow-pyarrow-integration-testing + # # Necessary because half 2.5 requires rust 1.81 or newer + # run: | + # cargo update -p half --precise 2.4.0 + # - name: Downgrade workspace dependencies + # # Necessary because half 2.5 requires rust 1.81 or newer + # run: | + # cargo update -p half --precise 2.4.0 + # - name: Check all packages + # run: | + # # run `cargo msrv verify --manifest-path "path/to/Cargo.toml"` to see problematic dependencies + # find . -mindepth 2 -name Cargo.toml | while read -r dir + # do + # echo "Checking package '$dir'" + # cargo msrv verify --manifest-path "$dir" --output-format=json || exit 1 + # done diff --git a/Cargo.toml b/Cargo.toml index aa9b7a8e39b0..5bd7b81120b8 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -100,4 +100,5 @@ chrono = { version = "0.4.40", default-features = false, features = ["clock"] } [profile.profiling] inherits = "release" debug = true -strip = false \ No newline at end of file +strip = false + diff --git a/arrow-array/Cargo.toml b/arrow-array/Cargo.toml index a65c0c9ca8e6..54553eec0616 100644 --- a/arrow-array/Cargo.toml +++ b/arrow-array/Cargo.toml @@ -45,7 +45,7 @@ arrow-data = { workspace = true } chrono = { workspace = true } chrono-tz = { version = "0.10", optional = true } num = { version = "0.4.1", default-features = false, features = ["std"] } -half = { version = "2.1", default-features = false, features = ["num-traits"] } +half = { version = ">=2.1.0, <=2.7.1", default-features = false, features = ["num-traits"] } hashbrown = { version = "0.15.1", default-features = false } [package.metadata.docs.rs] diff --git a/arrow-avro/Cargo.toml b/arrow-avro/Cargo.toml index d531bc18d04b..e149e93757c7 100644 --- a/arrow-avro/Cargo.toml +++ b/arrow-avro/Cargo.toml @@ -49,7 +49,7 @@ serde = { version = "1.0.188", features = ["derive"] } flate2 = { version = "1.0", default-features = false, features = ["rust_backend"], optional = true } snap = { version = "1.0", default-features = false, optional = true } zstd = { version = "0.13", default-features = false, optional = true } -crc = { version = "3.0", optional = true } +crc = { version = "=3.0.1", optional = true } [dev-dependencies] rand = { version = "0.9", default-features = false, features = ["std", "std_rng", "thread_rng"] } diff --git a/arrow-cast/Cargo.toml b/arrow-cast/Cargo.toml index 49145cf987f9..b0bd7686f74c 100644 --- a/arrow-cast/Cargo.toml +++ b/arrow-cast/Cargo.toml @@ -46,17 +46,17 @@ arrow-data = { workspace = true } arrow-schema = { workspace = true } arrow-select = { workspace = true } chrono = { workspace = true } -half = { version = "2.1", default-features = false } +half = { version = ">=2.1.0, <=2.7.1", default-features = false } num = { version = "0.4", default-features = false, features = ["std"] } lexical-core = { version = "1.0", default-features = false, features = ["write-integers", "write-floats", "parse-integers", "parse-floats"] } atoi = "2.0.0" -comfy-table = { version = "7.0", optional = true, default-features = false } +comfy-table = { version = "=7.1.0", optional = true, default-features = false } base64 = "0.22" ryu = "1.0.16" [dev-dependencies] criterion = { version = "0.5", default-features = false } -half = { version = "2.1", default-features = false } +half = { version = ">=2.1.0, <=2.7.1", default-features = false } rand = "0.9" [[bench]] diff --git a/arrow-data/Cargo.toml b/arrow-data/Cargo.toml index fbed24fea1fa..b876e1ee0a11 100644 --- a/arrow-data/Cargo.toml +++ b/arrow-data/Cargo.toml @@ -49,7 +49,7 @@ arrow-buffer = { workspace = true } arrow-schema = { workspace = true } num = { version = "0.4", default-features = false, features = ["std"] } -half = { version = "2.1", default-features = false } +half = { version = ">=2.1.0, <=2.7.1", default-features = false } [dev-dependencies] diff --git a/arrow-flight/Cargo.toml b/arrow-flight/Cargo.toml index 30ee7169702b..c84dbf044d46 100644 --- a/arrow-flight/Cargo.toml +++ b/arrow-flight/Cargo.toml @@ -48,7 +48,7 @@ prost = { version = "0.13.1", default-features = false, features = ["prost-deriv # For Timestamp type prost-types = { version = "0.13.1", default-features = false } tokio = { version = "1.0", default-features = false, features = ["macros", "rt", "rt-multi-thread"], optional = true } -tonic = { version = "0.12.3", default-features = false, features = ["transport", "codegen", "prost"] } +tonic = { version = "=0.12.3", default-features = false, features = ["transport", "codegen", "prost"] } # CLI-related dependencies anyhow = { version = "1.0", optional = true } @@ -68,10 +68,10 @@ cli = ["arrow-array/chrono-tz", "arrow-cast/prettyprint", "tonic/tls-webpki-root [dev-dependencies] arrow-cast = { workspace = true, features = ["prettyprint"] } -assert_cmd = "2.0.8" +assert_cmd = ">=2.0.8, <2.1.0" http = "1.1.0" http-body = "1.0.0" -hyper-util = "0.1" +hyper-util = "=0.1.4" pin-project-lite = "0.2" tempfile = "3.3" tracing-log = { version = "0.2" } diff --git a/arrow-ipc/Cargo.toml b/arrow-ipc/Cargo.toml index a1f826ef7d10..373093e2e1e5 100644 --- a/arrow-ipc/Cargo.toml +++ b/arrow-ipc/Cargo.toml @@ -41,7 +41,7 @@ arrow-buffer = { workspace = true } arrow-data = { workspace = true } arrow-schema = { workspace = true } flatbuffers = { version = "25.2.10", default-features = false } -lz4_flex = { version = "0.11", default-features = false, features = ["std", "frame"], optional = true } +lz4_flex = { version = ">=0.11.0, <=0.11.6", default-features = false, features = ["std", "frame"], optional = true } zstd = { version = "0.13.0", default-features = false, optional = true } [features] diff --git a/arrow-json/Cargo.toml b/arrow-json/Cargo.toml index cae0e173b445..9925480a6eba 100644 --- a/arrow-json/Cargo.toml +++ b/arrow-json/Cargo.toml @@ -41,8 +41,8 @@ arrow-buffer = { workspace = true } arrow-cast = { workspace = true } arrow-data = { workspace = true } arrow-schema = { workspace = true } -half = { version = "2.1", default-features = false } -indexmap = { version = "2.0", default-features = false, features = ["std"] } +half = { version = ">=2.1.0, <=2.7.1", default-features = false } +indexmap = { version = ">=2.0.2, <=2.12.1", default-features = false, features = ["std"] } num = { version = "0.4", default-features = false, features = ["std"] } serde = { version = "1.0", default-features = false } serde_json = { version = "1.0", default-features = false, features = ["std"] } diff --git a/arrow-ord/Cargo.toml b/arrow-ord/Cargo.toml index ae76841bda39..c246636a9844 100644 --- a/arrow-ord/Cargo.toml +++ b/arrow-ord/Cargo.toml @@ -43,5 +43,5 @@ arrow-schema = { workspace = true } arrow-select = { workspace = true } [dev-dependencies] -half = { version = "2.1", default-features = false, features = ["num-traits"] } +half = { version = ">=2.1.0, <=2.7.1", default-features = false, features = ["num-traits"] } rand = { version = "0.9", default-features = false, features = ["std", "std_rng"] } diff --git a/arrow-row/Cargo.toml b/arrow-row/Cargo.toml index 7d136939b05c..33e2b0bd7ebd 100644 --- a/arrow-row/Cargo.toml +++ b/arrow-row/Cargo.toml @@ -41,7 +41,7 @@ arrow-buffer = { workspace = true } arrow-data = { workspace = true } arrow-schema = { workspace = true } -half = { version = "2.1", default-features = false } +half = { version = ">=2.1.0, <=2.7.1", default-features = false } [dev-dependencies] arrow-cast = { workspace = true } diff --git a/arrow/Cargo.toml b/arrow/Cargo.toml index 31398b462eb3..91a36d920d4f 100644 --- a/arrow/Cargo.toml +++ b/arrow/Cargo.toml @@ -55,7 +55,7 @@ arrow-string = { workspace = true } rand = { version = "0.9", default-features = false, features = ["std", "std_rng", "thread_rng"], optional = true } pyo3 = { version = "0.24.1", default-features = false, optional = true } -half = { version = "2.1", default-features = false, optional = true } +half = { version = ">=2.1.0, <=2.7.1", default-features = false, optional = true } [package.metadata.docs.rs] all-features = true @@ -85,7 +85,7 @@ canonical_extension_types = ["arrow-schema/canonical_extension_types"] [dev-dependencies] chrono = { workspace = true } criterion = { version = "0.5", default-features = false } -half = { version = "2.1", default-features = false } +half = { version = ">=2.1.0, <=2.7.1", default-features = false } rand = { version = "0.9", default-features = false, features = ["std", "std_rng", "thread_rng"] } serde = { version = "1.0", default-features = false, features = ["derive"] } # used in examples diff --git a/parquet/Cargo.toml b/parquet/Cargo.toml index 9d21cc1d40a5..268fa40c00c9 100644 --- a/parquet/Cargo.toml +++ b/parquet/Cargo.toml @@ -52,7 +52,7 @@ thrift = { version = "0.17", default-features = false } snap = { version = "1.0", default-features = false, optional = true } brotli = { version = "8.0", default-features = false, features = ["std"], optional = true } flate2 = { version = "1.1", default-features = false, features = ["zlib-rs"], optional = true } -lz4_flex = { version = "0.11", default-features = false, features = ["std", "frame"], optional = true } +lz4_flex = { version = ">=0.11.0, <=0.11.6", default-features = false, features = ["std", "frame"], optional = true } zstd = { version = "0.13", optional = true, default-features = false } chrono = { workspace = true } num = { version = "0.4", default-features = false } @@ -67,7 +67,7 @@ tokio = { version = "1.0", optional = true, default-features = false, features = hashbrown = { version = "0.15", default-features = false } twox-hash = { version = "2.0", default-features = false, features = ["xxhash64"] } paste = { version = "1.0" } -half = { version = "2.1", default-features = false, features = ["num-traits"] } +half = { version = ">=2.1.0, <=2.7.1", default-features = false, features = ["num-traits"] } crc32fast = { version = "1.4.2", optional = true, default-features = false } simdutf8 = { version = "0.1.5", optional = true, default-features = false } ring = { version = "0.17", default-features = false, features = ["std"], optional = true } @@ -79,7 +79,7 @@ snap = { version = "1.0", default-features = false } tempfile = { version = "3.0", default-features = false } brotli = { version = "8.0", default-features = false, features = ["std"] } flate2 = { version = "1.0", default-features = false, features = ["rust_backend"] } -lz4_flex = { version = "0.11", default-features = false, features = ["std", "frame"] } +lz4_flex = { version = ">=0.11.0, <=0.11.6", default-features = false, features = ["std", "frame"] } zstd = { version = "0.13", default-features = false } serde_json = { version = "1.0", features = ["std"], default-features = false } arrow = { workspace = true, features = ["ipc", "test_utils", "prettyprint", "json"] } From b6bffcb8a22600faffe9580d160269df30501722 Mon Sep 17 00:00:00 2001 From: Dan Harris Date: Sat, 8 Mar 2025 15:24:32 -0500 Subject: [PATCH 2/7] Add coop yield in parquet reader (fork-only) --- parquet/src/arrow/arrow_reader/mod.rs | 59 +++++++++++++++++++++++++++ parquet/src/arrow/async_reader/mod.rs | 17 ++++---- 2 files changed, 69 insertions(+), 7 deletions(-) diff --git a/parquet/src/arrow/arrow_reader/mod.rs b/parquet/src/arrow/arrow_reader/mod.rs index 2f670a64e108..6cf0cfd2218f 100644 --- a/parquet/src/arrow/arrow_reader/mod.rs +++ b/parquet/src/arrow/arrow_reader/mod.rs @@ -1015,6 +1015,65 @@ pub(crate) fn evaluate_predicate( }) } +/// Maximum number of bytes that can be evaluated in a row filter +/// before yielding back to the scheduler +const DECODE_BUDGET: usize = 2 * 1024 * 1024; + +/// Evaluates an [`ArrowPredicate`], returning a [`RowSelection`] indicating +/// which rows to return. +/// +/// `input_selection`: Optional pre-existing selection. If `Some`, then the +/// final [`RowSelection`] will be the conjunction of it and the rows selected +/// by `predicate`. +/// +/// Note: A pre-existing selection may come from evaluating a previous predicate +/// or if the [`ParquetRecordBatchReader`] specified an explicit +/// [`RowSelection`] in addition to one or more predicates. +#[allow(dead_code)] +pub(crate) async fn evaluate_predicate_coop( + batch_size: usize, + array_reader: Box, + input_selection: Option, + predicate: &mut dyn ArrowPredicate, +) -> Result { + let mut budget = DECODE_BUDGET; + + let reader = ParquetRecordBatchReader::new(batch_size, array_reader, input_selection.clone()); + let mut filters = vec![]; + for maybe_batch in reader { + let maybe_batch = maybe_batch?; + budget = budget.saturating_sub(maybe_batch.get_array_memory_size()); + + let input_rows = maybe_batch.num_rows(); + let filter = predicate.evaluate(maybe_batch)?; + // Since user supplied predicate, check error here to catch bugs quickly + if filter.len() != input_rows { + return Err(arrow_err!( + "ArrowPredicate predicate returned {} rows, expected {input_rows}", + filter.len() + )); + } + match filter.null_count() { + 0 => filters.push(filter), + _ => filters.push(prep_null_mask_filter(&filter)), + }; + + if budget == 0 { + // If we have consumed our decode budget, reset the budget and yield + // back to the scheduler + budget = DECODE_BUDGET; + #[cfg(feature = "async")] + tokio::task::yield_now().await; + } + } + + let raw = RowSelection::from_filters(&filters); + Ok(match input_selection { + Some(selection) => selection.and_then(&raw), + None => raw, + }) +} + #[cfg(test)] mod tests { use std::cmp::min; diff --git a/parquet/src/arrow/async_reader/mod.rs b/parquet/src/arrow/async_reader/mod.rs index 45df68821ca8..76f173a4cc57 100644 --- a/parquet/src/arrow/async_reader/mod.rs +++ b/parquet/src/arrow/async_reader/mod.rs @@ -40,7 +40,7 @@ use arrow_schema::{DataType, Fields, Schema, SchemaRef}; use crate::arrow::array_reader::{build_array_reader, RowGroups}; use crate::arrow::arrow_reader::{ - apply_range, evaluate_predicate, selects_any, ArrowReaderBuilder, ArrowReaderMetadata, + apply_range, evaluate_predicate_coop, selects_any, ArrowReaderBuilder, ArrowReaderMetadata, ArrowReaderOptions, ParquetRecordBatchReader, RowFilter, RowSelection, }; use crate::arrow::ProjectionMask; @@ -600,12 +600,15 @@ where let array_reader = build_array_reader(self.fields.as_deref(), predicate_projection, &row_group)?; - selection = Some(evaluate_predicate( - batch_size, - array_reader, - selection, - predicate.as_mut(), - )?); + selection = Some( + evaluate_predicate_coop( + batch_size, + array_reader, + selection, + predicate.as_mut(), + ) + .await?, + ); } } From 67a7d1cb0cf482be22f5b04886c00e93545ffeb0 Mon Sep 17 00:00:00 2001 From: Thomas Peiselt Date: Mon, 23 Feb 2026 12:55:03 +0100 Subject: [PATCH 3/7] Tonic 0.13 (56.0.0 #7839) --- .github/workflows/arrow_flight.yml | 6 +++--- arrow-flight/Cargo.toml | 11 +++++------ arrow-flight/gen/Cargo.toml | 2 +- arrow-flight/src/arrow.flight.protocol.rs | 14 ++++++++------ arrow-integration-testing/Cargo.toml | 2 +- 5 files changed, 18 insertions(+), 17 deletions(-) diff --git a/.github/workflows/arrow_flight.yml b/.github/workflows/arrow_flight.yml index b50366b4f0b9..fd769bd1bab6 100644 --- a/.github/workflows/arrow_flight.yml +++ b/.github/workflows/arrow_flight.yml @@ -58,9 +58,9 @@ jobs: - name: Test --all-features run: | cargo test -p arrow-flight --all-features - - name: Test --examples - run: | - cargo test -p arrow-flight --features=flight-sql-experimental,tls --examples + # - name: Test --examples + # run: | + # cargo test -p arrow-flight --features=flight-sql-experimental --examples clippy: name: Clippy diff --git a/arrow-flight/Cargo.toml b/arrow-flight/Cargo.toml index c84dbf044d46..4fdcca96b808 100644 --- a/arrow-flight/Cargo.toml +++ b/arrow-flight/Cargo.toml @@ -48,7 +48,7 @@ prost = { version = "0.13.1", default-features = false, features = ["prost-deriv # For Timestamp type prost-types = { version = "0.13.1", default-features = false } tokio = { version = "1.0", default-features = false, features = ["macros", "rt", "rt-multi-thread"], optional = true } -tonic = { version = "=0.12.3", default-features = false, features = ["transport", "codegen", "prost"] } +tonic = { version = "=0.13.0", default-features = false, features = ["transport", "codegen", "prost", "router"] } # CLI-related dependencies anyhow = { version = "1.0", optional = true } @@ -62,7 +62,6 @@ all-features = true [features] default = [] flight-sql-experimental = ["dep:arrow-arith", "dep:arrow-data", "dep:arrow-ord", "dep:arrow-row", "dep:arrow-select", "dep:arrow-string", "dep:once_cell", "dep:paste"] -tls = ["tonic/tls"] # Enable CLI tools cli = ["arrow-array/chrono-tz", "arrow-cast/prettyprint", "tonic/tls-webpki-roots", "dep:anyhow", "dep:clap", "dep:tracing-log", "dep:tracing-subscriber"] @@ -83,18 +82,18 @@ uuid = { version = "1.10.0", features = ["v4"] } [[example]] name = "flight_sql_server" -required-features = ["flight-sql-experimental", "tls"] +required-features = ["flight-sql-experimental"] [[bin]] name = "flight_sql_client" -required-features = ["cli", "flight-sql-experimental", "tls"] +required-features = ["cli", "flight-sql-experimental"] [[test]] name = "flight_sql_client" path = "tests/flight_sql_client.rs" -required-features = ["flight-sql-experimental", "tls"] +required-features = ["flight-sql-experimental"] [[test]] name = "flight_sql_client_cli" path = "tests/flight_sql_client_cli.rs" -required-features = ["cli", "flight-sql-experimental", "tls"] +required-features = ["cli", "flight-sql-experimental"] diff --git a/arrow-flight/gen/Cargo.toml b/arrow-flight/gen/Cargo.toml index 79d46cd377fa..d23d3574c0b6 100644 --- a/arrow-flight/gen/Cargo.toml +++ b/arrow-flight/gen/Cargo.toml @@ -33,4 +33,4 @@ publish = false # Pin specific version of the tonic-build dependencies to avoid auto-generated # (and checked in) arrow.flight.protocol.rs from changing prost-build = { version = "=0.13.5", default-features = false } -tonic-build = { version = "=0.12.3", default-features = false, features = ["transport", "prost"] } +tonic-build = { version = "=0.13.0", default-features = false, features = ["transport", "prost"] } diff --git a/arrow-flight/src/arrow.flight.protocol.rs b/arrow-flight/src/arrow.flight.protocol.rs index 0cd4f6948b77..a08ea01105e5 100644 --- a/arrow-flight/src/arrow.flight.protocol.rs +++ b/arrow-flight/src/arrow.flight.protocol.rs @@ -448,7 +448,7 @@ pub mod flight_service_client { } impl FlightServiceClient where - T: tonic::client::GrpcService, + T: tonic::client::GrpcService, T::Error: Into, T::ResponseBody: Body + std::marker::Send + 'static, ::Error: Into + std::marker::Send, @@ -469,13 +469,13 @@ pub mod flight_service_client { F: tonic::service::Interceptor, T::ResponseBody: Default, T: tonic::codegen::Service< - http::Request, + http::Request, Response = http::Response< - >::ResponseBody, + >::ResponseBody, >, >, , + http::Request, >>::Error: Into + std::marker::Send + std::marker::Sync, { FlightServiceClient::new(InterceptedService::new(inner, interceptor)) @@ -1098,7 +1098,7 @@ pub mod flight_service_server { B: Body + std::marker::Send + 'static, B::Error: Into + std::marker::Send + 'static, { - type Response = http::Response; + type Response = http::Response; type Error = std::convert::Infallible; type Future = BoxFuture; fn poll_ready( @@ -1571,7 +1571,9 @@ pub mod flight_service_server { } _ => { Box::pin(async move { - let mut response = http::Response::new(empty_body()); + let mut response = http::Response::new( + tonic::body::Body::default(), + ); let headers = response.headers_mut(); headers .insert( diff --git a/arrow-integration-testing/Cargo.toml b/arrow-integration-testing/Cargo.toml index 8654b4b92734..3a717532b2b1 100644 --- a/arrow-integration-testing/Cargo.toml +++ b/arrow-integration-testing/Cargo.toml @@ -43,7 +43,7 @@ prost = { version = "0.13", default-features = false } serde = { version = "1.0", default-features = false, features = ["rc", "derive"] } serde_json = { version = "1.0", default-features = false, features = ["std"] } tokio = { version = "1.0", default-features = false, features = [ "rt-multi-thread"] } -tonic = { version = "0.12", default-features = false } +tonic = { version = "=0.13", default-features = false } tracing-subscriber = { version = "0.3.1", default-features = false, features = ["fmt"], optional = true } flate2 = { version = "1", default-features = false, features = ["rust_backend"] } From 4cbfabe1b5911dcbf9a6e50d047bdbf7d14f7d27 Mon Sep 17 00:00:00 2001 From: Thomas Peiselt Date: Wed, 25 Feb 2026 18:25:16 +0100 Subject: [PATCH 4/7] Add row_id column to parquet reader for index-based row lookups (fork-only) --- parquet/src/arrow/arrow_reader/mod.rs | 126 ++++++- parquet/src/arrow/async_reader/mod.rs | 501 +++++++++++++++++++++++++- 2 files changed, 610 insertions(+), 17 deletions(-) diff --git a/parquet/src/arrow/arrow_reader/mod.rs b/parquet/src/arrow/arrow_reader/mod.rs index 6cf0cfd2218f..9cdbea43f8fa 100644 --- a/parquet/src/arrow/arrow_reader/mod.rs +++ b/parquet/src/arrow/arrow_reader/mod.rs @@ -17,10 +17,11 @@ //! Contains reader which reads parquet data into arrow [`RecordBatch`] +use arrow_array::builder::UInt64Builder; use arrow_array::cast::AsArray; -use arrow_array::Array; +use arrow_array::{Array, ArrayRef}; use arrow_array::{RecordBatch, RecordBatchReader}; -use arrow_schema::{ArrowError, DataType as ArrowType, Schema, SchemaRef}; +use arrow_schema::{ArrowError, DataType as ArrowType, Field, FieldRef, Schema, SchemaRef}; use arrow_select::filter::prep_null_mask_filter; pub use filter::{ArrowPredicate, ArrowPredicateFn, RowFilter}; pub use selection::{RowSelection, RowSelector}; @@ -110,6 +111,9 @@ pub struct ArrowReaderBuilder { pub(crate) limit: Option, pub(crate) offset: Option, + + #[allow(unused)] + pub(crate) row_id: Option, } impl ArrowReaderBuilder { @@ -126,6 +130,7 @@ impl ArrowReaderBuilder { selection: None, limit: None, offset: None, + row_id: None, } } @@ -152,6 +157,15 @@ impl ArrowReaderBuilder { Self { batch_size, ..self } } + /// Project a column into the result with name `field_name` that will contain the row ID + /// for each row. The row ID will be the row offset of the row in the underlying file + pub fn with_row_id(self, field_name: impl Into) -> Self { + Self { + row_id: Some(RowId::field_ref(field_name)), + ..self + } + } + /// Only read data from the provided row group indexes /// /// This is also called row group filtering @@ -710,6 +724,8 @@ impl ParquetRecordBatchReaderBuilder { batch_size, array_reader, apply_range(selection, reader.num_rows(), self.offset, self.limit), + // TODO what do we do here? + None, )) } } @@ -786,6 +802,48 @@ impl Iterator for ReaderPageIterator { impl PageIterator for ReaderPageIterator {} +pub(crate) struct RowId { + offset: u64, + field: FieldRef, + buffer: UInt64Builder, +} + +impl RowId { + #[allow(unused)] + pub fn new(offset: u64, field: FieldRef, batch_size: usize) -> Self { + Self { + offset, + field, + buffer: UInt64Builder::with_capacity(batch_size), + } + } + + pub fn field_ref(name: impl Into) -> FieldRef { + Arc::new(Field::new(name, ArrowType::UInt64, false)) + } + + pub fn skip(&mut self, n: usize) { + self.offset += n as u64; + } + + pub fn field(&self) -> FieldRef { + self.field.clone() + } + + fn read(&mut self, n: usize) { + // SAFETY: We are appending a `Range` which has a trusted length + unsafe { + self.buffer + .append_trusted_len_iter(self.offset..self.offset + n as u64) + } + self.offset += n as u64; + } + + fn consume(&mut self) -> ArrayRef { + Arc::new(self.buffer.finish()) + } +} + /// An `Iterator>` that yields [`RecordBatch`] /// read from a parquet data source pub struct ParquetRecordBatchReader { @@ -793,6 +851,7 @@ pub struct ParquetRecordBatchReader { array_reader: Box, schema: SchemaRef, selection: Option>, + row_id: Option, } impl Iterator for ParquetRecordBatchReader { @@ -810,6 +869,10 @@ impl Iterator for ParquetRecordBatchReader { Err(e) => return Some(Err(e.into())), }; + if let Some(row_id) = self.row_id.as_mut() { + row_id.skip(skipped); + } + if skipped != front.row_count { return Some(Err(general_err!( "failed to skip rows, expected {}, got {}", @@ -840,16 +903,24 @@ impl Iterator for ParquetRecordBatchReader { }; match self.array_reader.read_records(to_read) { Ok(0) => break, - Ok(rec) => read_records += rec, + Ok(rec) => { + if let Some(rowid) = self.row_id.as_mut() { + rowid.read(rec); + } + read_records += rec + } Err(error) => return Some(Err(error.into())), } } } - None => { - if let Err(error) = self.array_reader.read_records(self.batch_size) { - return Some(Err(error.into())); + None => match self.array_reader.read_records(self.batch_size) { + Ok(n) => { + if let Some(rowid) = self.row_id.as_mut() { + rowid.read(n); + } } - } + Err(error) => return Some(Err(error.into())), + }, }; match self.array_reader.consume_batch() { @@ -863,7 +934,23 @@ impl Iterator for ParquetRecordBatchReader { match struct_array { Err(err) => Some(Err(err)), - Ok(e) => (e.len() > 0).then(|| Ok(RecordBatch::from(e))), + Ok(e) => { + if e.len() > 0 { + Some(Ok(match self.row_id.as_mut() { + Some(rowid) => { + let columns = std::iter::once(rowid.consume()) + .chain(e.columns().iter().cloned()) + .collect(); + + RecordBatch::try_new(self.schema.clone(), columns) + .expect("invalid schema") + } + None => RecordBatch::from(e), + })) + } else { + None + } + } } } } @@ -908,6 +995,7 @@ impl ParquetRecordBatchReader { array_reader, schema: Arc::new(Schema::new(levels.fields.clone())), selection: selection.map(|s| s.trim().into()), + row_id: None, }) } @@ -918,17 +1006,29 @@ impl ParquetRecordBatchReader { batch_size: usize, array_reader: Box, selection: Option, + rowid: Option, ) -> Self { - let schema = match array_reader.get_data_type() { - ArrowType::Struct(ref fields) => Schema::new(fields.clone()), + let struct_fields = match array_reader.get_data_type() { + ArrowType::Struct(ref fields) => fields.clone(), _ => unreachable!("Struct array reader's data type is not struct!"), }; + let schema = match rowid.as_ref() { + Some(rowid) => { + let fields: Vec<_> = std::iter::once(rowid.field()) + .chain(struct_fields.iter().cloned()) + .collect(); + Schema::new(fields) + } + None => Schema::new(struct_fields), + }; + Self { batch_size, array_reader, schema: Arc::new(schema), selection: selection.map(|s| s.trim().into()), + row_id: rowid, } } } @@ -989,7 +1089,8 @@ pub(crate) fn evaluate_predicate( input_selection: Option, predicate: &mut dyn ArrowPredicate, ) -> Result { - let reader = ParquetRecordBatchReader::new(batch_size, array_reader, input_selection.clone()); + let reader = + ParquetRecordBatchReader::new(batch_size, array_reader, input_selection.clone(), None); let mut filters = vec![]; for maybe_batch in reader { let maybe_batch = maybe_batch?; @@ -1038,7 +1139,8 @@ pub(crate) async fn evaluate_predicate_coop( ) -> Result { let mut budget = DECODE_BUDGET; - let reader = ParquetRecordBatchReader::new(batch_size, array_reader, input_selection.clone()); + let reader = + ParquetRecordBatchReader::new(batch_size, array_reader, input_selection.clone(), None); let mut filters = vec![]; for maybe_batch in reader { let maybe_batch = maybe_batch?; diff --git a/parquet/src/arrow/async_reader/mod.rs b/parquet/src/arrow/async_reader/mod.rs index 76f173a4cc57..3ff408e31ab7 100644 --- a/parquet/src/arrow/async_reader/mod.rs +++ b/parquet/src/arrow/async_reader/mod.rs @@ -36,12 +36,12 @@ use futures::stream::Stream; use tokio::io::{AsyncRead, AsyncReadExt, AsyncSeek, AsyncSeekExt}; use arrow_array::RecordBatch; -use arrow_schema::{DataType, Fields, Schema, SchemaRef}; +use arrow_schema::{DataType, FieldRef, Fields, Schema, SchemaRef}; use crate::arrow::array_reader::{build_array_reader, RowGroups}; use crate::arrow::arrow_reader::{ apply_range, evaluate_predicate_coop, selects_any, ArrowReaderBuilder, ArrowReaderMetadata, - ArrowReaderOptions, ParquetRecordBatchReader, RowFilter, RowSelection, + ArrowReaderOptions, ParquetRecordBatchReader, RowFilter, RowId, RowSelection, }; use crate::arrow::ProjectionMask; @@ -509,17 +509,27 @@ impl ParquetRecordBatchStreamBuilder { fields: self.fields, limit: self.limit, offset: self.offset, + rowid: self.row_id.clone(), }; // Ensure schema of ParquetRecordBatchStream respects projection, and does // not store metadata (same as for ParquetRecordBatchReader and emitted RecordBatches) - let projected_fields = match reader_factory.fields.as_deref().map(|pf| &pf.arrow_type) { + let mut projected_fields = match reader_factory.fields.as_deref().map(|pf| &pf.arrow_type) { Some(DataType::Struct(fields)) => { fields.filter_leaves(|idx, _| self.projection.leaf_included(idx)) } None => Fields::empty(), _ => unreachable!("Must be Struct for root type"), }; + + if let Some(field) = &self.row_id { + projected_fields = Fields::from( + std::iter::once(field.clone()) + .chain(projected_fields.iter().cloned()) + .collect::>(), + ); + } + let schema = Arc::new(Schema::new(projected_fields)); Ok(ParquetRecordBatchStream { @@ -551,6 +561,8 @@ struct ReaderFactory { limit: Option, offset: Option, + + rowid: Option, } impl ReaderFactory @@ -649,10 +661,19 @@ where .fetch(&mut self.input, &projection, selection.as_ref()) .await?; + let rowid = self.rowid.clone().map(|field| { + let offset = self.metadata.row_groups()[..row_group_idx] + .iter() + .map(|rg| rg.num_rows() as u64) + .sum::(); + RowId::new(offset, field, batch_size) + }); + let reader = ParquetRecordBatchReader::new( batch_size, build_array_reader(self.fields.as_deref(), &projection, &row_group)?, selection, + rowid, ); Ok((self, Some(reader))) @@ -1096,14 +1117,15 @@ mod tests { use crate::file::properties::WriterProperties; use arrow::compute::kernels::cmp::eq; use arrow::error::Result as ArrowResult; - use arrow_array::builder::{ListBuilder, StringBuilder}; + use arrow_array::builder::{ListBuilder, StringBuilder, UInt64Builder}; use arrow_array::cast::AsArray; - use arrow_array::types::Int32Type; + use arrow_array::types::{Int32Type, UInt64Type}; use arrow_array::{ Array, ArrayRef, Int32Array, Int8Array, RecordBatchReader, Scalar, StringArray, StructArray, UInt64Array, }; use arrow_schema::{DataType, Field, Schema}; + use arrow_select::concat::concat; use futures::{StreamExt, TryStreamExt}; use rand::{rng, Rng}; use std::collections::HashMap; @@ -1115,6 +1137,7 @@ mod tests { data: Bytes, metadata: Option>, requests: Arc>>>, + max_concurrent_requests: Arc>, } impl AsyncFileReader for TestReader { @@ -1130,6 +1153,27 @@ mod tests { .boxed() } + fn get_byte_ranges( + &mut self, + ranges: Vec>, + ) -> BoxFuture<'_, Result>> { + let usize_ranges: Vec> = ranges + .iter() + .map(|r| r.start as usize..r.end as usize) + .collect(); + self.requests.lock().unwrap().extend(usize_ranges.clone()); + let mut max = self.max_concurrent_requests.lock().unwrap(); + if ranges.len() > *max { + *max = ranges.len(); + } + + let mut results = Vec::with_capacity(ranges.len()); + for range in usize_ranges { + results.push(self.data.slice(range)); + } + futures::future::ready(Ok(results)).boxed() + } + fn get_metadata<'a>( &'a mut self, options: Option<&'a ArrowReaderOptions>, @@ -1153,6 +1197,7 @@ mod tests { data: data.clone(), metadata: Default::default(), requests: Default::default(), + max_concurrent_requests: Default::default(), }; let requests = async_reader.requests.clone(); @@ -1206,6 +1251,7 @@ mod tests { data: data.clone(), metadata: Default::default(), requests: Default::default(), + max_concurrent_requests: Default::default(), }; let requests = async_reader.requests.clone(); @@ -1257,6 +1303,84 @@ mod tests { ); } + #[tokio::test] + async fn test_async_reader_with_rowid() { + let testdata = arrow::util::test_util::parquet_test_data(); + let path = format!("{testdata}/alltypes_plain.parquet"); + let data = Bytes::from(std::fs::read(path).unwrap()); + + let metadata = ParquetMetaDataReader::new() + .parse_and_finish(&data) + .unwrap(); + let metadata = Arc::new(metadata); + + assert_eq!(metadata.num_row_groups(), 1); + + let async_reader = TestReader { + data: data.clone(), + metadata: Some(metadata.clone()), + requests: Default::default(), + max_concurrent_requests: Default::default(), + }; + + let requests = async_reader.requests.clone(); + let builder = ParquetRecordBatchStreamBuilder::new(async_reader) + .await + .unwrap(); + + let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![1, 2]); + let stream = builder + .with_projection(mask.clone()) + .with_batch_size(1024) + .with_row_id("_rowid") + .build() + .unwrap(); + + assert_eq!( + stream + .schema() + .fields() + .first() + .expect("no fields in schema") + .name(), + "_rowid" + ); + + let async_batches: Vec<_> = stream.try_collect().await.unwrap(); + + assert!(async_batches.iter().all(|batch| { + batch + .schema() + .fields() + .first() + .expect("no fields in schema") + .name() + == "_rowid" + })); + + let rowid_arrays = async_batches + .iter() + .map(|batch| batch.column(0).as_ref()) + .collect::>(); + let rowids = concat(&rowid_arrays).expect("concat rowids"); + + let expected_rowids = UInt64Array::from_iter_values(0..rowids.len() as u64); + + assert_eq!(rowids.as_primitive::(), &expected_rowids); + + let requests = requests.lock().unwrap(); + let (offset_1, length_1) = metadata.row_group(0).column(1).byte_range(); + let (offset_2, length_2) = metadata.row_group(0).column(2).byte_range(); + + assert_eq!( + &requests[..], + &[ + offset_1 as usize..(offset_1 + length_1) as usize, + offset_2 as usize..(offset_2 + length_2) as usize + ] + ); + } + #[tokio::test] async fn test_async_reader_with_index() { let testdata = arrow::util::test_util::parquet_test_data(); @@ -1267,6 +1391,7 @@ mod tests { data: data.clone(), metadata: Default::default(), requests: Default::default(), + max_concurrent_requests: Default::default(), }; let options = ArrowReaderOptions::new().with_page_index(true); @@ -1319,6 +1444,61 @@ mod tests { assert_eq!(async_batches, sync_batches); } + #[tokio::test] + async fn test_async_reader_with_rowid_offset() { + let testdata = arrow::util::test_util::parquet_test_data(); + let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet"); + let data = Bytes::from(std::fs::read(path).unwrap()); + + let metadata = ParquetMetaDataReader::new() + .parse_and_finish(&data) + .unwrap(); + let metadata = Arc::new(metadata); + + assert_eq!(metadata.num_row_groups(), 1); + + let async_reader = TestReader { + data: data.clone(), + metadata: Some(metadata.clone()), + requests: Default::default(), + max_concurrent_requests: Default::default(), + }; + + let builder = ParquetRecordBatchStreamBuilder::new(async_reader) + .await + .unwrap(); + + let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![1, 2]); + let stream = builder + .with_projection(mask.clone()) + .with_batch_size(1024) + .with_offset(3) + .with_row_id("_rowid") + .build() + .unwrap(); + + let async_batches: Vec<_> = stream.try_collect().await.unwrap(); + + assert!(async_batches.iter().all(|batch| { + batch + .schema() + .fields() + .first() + .expect("no fields in schema") + .name() + == "_rowid" + })); + + let rowid_arrays = async_batches + .iter() + .map(|batch| batch.column(0).as_ref()) + .collect::>(); + let rowids = concat(&rowid_arrays).expect("concat rowids"); + + let expected_rowids = UInt64Array::from_iter_values(3..rowids.len() as u64 + 3); + assert_eq!(rowids.as_primitive::(), &expected_rowids); + } + #[tokio::test] async fn test_async_reader_with_limit() { let testdata = arrow::util::test_util::parquet_test_data(); @@ -1336,6 +1516,7 @@ mod tests { data: data.clone(), metadata: Default::default(), requests: Default::default(), + max_concurrent_requests: Default::default(), }; let builder = ParquetRecordBatchStreamBuilder::new(async_reader) @@ -1377,6 +1558,7 @@ mod tests { data: data.clone(), metadata: Default::default(), requests: Default::default(), + max_concurrent_requests: Default::default(), }; let options = ArrowReaderOptions::new().with_page_index(true); @@ -1455,6 +1637,7 @@ mod tests { data: data.clone(), metadata: Default::default(), requests: Default::default(), + max_concurrent_requests: Default::default(), }; let options = ArrowReaderOptions::new().with_page_index(true); @@ -1481,6 +1664,104 @@ mod tests { } } + #[tokio::test] + async fn test_fuzz_async_reader_with_rowid_and_selection() { + let testdata = arrow::util::test_util::parquet_test_data(); + let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet"); + let data = Bytes::from(std::fs::read(path).unwrap()); + + let metadata = ParquetMetaDataReader::new() + .parse_and_finish(&data) + .unwrap(); + let metadata = Arc::new(metadata); + + assert_eq!(metadata.num_row_groups(), 1); + + let mut rand = rng(); + + for _ in 0..100 { + let mut expected_rowids_builder = UInt64Builder::new(); + let mut offset = 0; + + let mut expected_rows = 0; + let mut total_rows = 0; + let mut skip = false; + let mut selectors = vec![]; + + while total_rows < 7300 { + let row_count: usize = rand.random_range(1..100); + + let row_count = row_count.min(7300 - total_rows); + + selectors.push(RowSelector { row_count, skip }); + + total_rows += row_count; + if !skip { + expected_rowids_builder.append_slice( + (offset..offset + row_count as u64) + .collect::>() + .as_slice(), + ); + expected_rows += row_count; + } + + offset += row_count as u64; + + skip = !skip; + } + + let selection = RowSelection::from(selectors); + + let async_reader = TestReader { + data: data.clone(), + metadata: Some(metadata.clone()), + requests: Default::default(), + max_concurrent_requests: Default::default(), + }; + + let options = ArrowReaderOptions::new().with_page_index(true); + let builder = ParquetRecordBatchStreamBuilder::new_with_options(async_reader, options) + .await + .unwrap(); + + let col_idx: usize = rand.random_range(0..13); + let mask = ProjectionMask::leaves(builder.parquet_schema(), vec![col_idx]); + + let stream = builder + .with_projection(mask.clone()) + .with_row_selection(selection.clone()) + .with_row_id("_rowid") + .build() + .expect("building stream"); + + let async_batches: Vec<_> = stream.try_collect().await.unwrap(); + + let expected_rowids = expected_rowids_builder.finish(); + + assert!(async_batches.iter().all(|batch| { + batch + .schema() + .fields() + .first() + .expect("no fields in schema") + .name() + == "_rowid" + })); + + let rowid_arrays = async_batches + .iter() + .map(|batch| batch.column(0).as_ref()) + .collect::>(); + let rowids = concat(&rowid_arrays).expect("concat rowids"); + + assert_eq!(rowids.as_primitive::(), &expected_rowids); + + let actual_rows: usize = async_batches.into_iter().map(|b| b.num_rows()).sum(); + + assert_eq!(actual_rows, expected_rows); + } + } + #[tokio::test] async fn test_async_reader_zero_row_selector() { //See https://github.com/apache/arrow-rs/issues/2669 @@ -1521,6 +1802,7 @@ mod tests { data: data.clone(), metadata: Default::default(), requests: Default::default(), + max_concurrent_requests: Default::default(), }; let options = ArrowReaderOptions::new().with_page_index(true); @@ -1573,8 +1855,10 @@ mod tests { data, metadata: Default::default(), requests: Default::default(), + max_concurrent_requests: Default::default(), }; let requests = test.requests.clone(); + let max_concurrent_requests = test.max_concurrent_requests.clone(); let a_scalar = StringArray::from_iter_values(["b"]); let a_filter = ArrowPredicateFn::new( @@ -1615,6 +1899,86 @@ mod tests { let val = col.as_any().downcast_ref::().unwrap().value(0); assert_eq!(val, 3); + // Should only have made 3 requests + assert_eq!(requests.lock().unwrap().len(), 3); + assert_eq!(*max_concurrent_requests.lock().unwrap(), 1); + } + + #[tokio::test] + async fn test_async_reader_with_row_id_and_row_filter() { + let a = StringArray::from_iter_values(["a", "b", "b", "b", "c", "c"]); + let b = StringArray::from_iter_values(["1", "2", "3", "4", "5", "6"]); + let c = Int32Array::from_iter(0..6); + let data = RecordBatch::try_from_iter([ + ("a", Arc::new(a) as ArrayRef), + ("b", Arc::new(b) as ArrayRef), + ("c", Arc::new(c) as ArrayRef), + ]) + .unwrap(); + + let mut buf = Vec::with_capacity(1024); + let mut writer = ArrowWriter::try_new(&mut buf, data.schema(), None).unwrap(); + writer.write(&data).unwrap(); + writer.close().unwrap(); + + let data: Bytes = buf.into(); + let metadata = ParquetMetaDataReader::new() + .parse_and_finish(&data) + .unwrap(); + let parquet_schema = metadata.file_metadata().schema_descr_ptr(); + + let test = TestReader { + data, + metadata: Some(Arc::new(metadata)), + requests: Default::default(), + max_concurrent_requests: Default::default(), + }; + let requests = test.requests.clone(); + + let a_scalar = StringArray::from_iter_values(["b"]); + let a_filter = ArrowPredicateFn::new( + ProjectionMask::leaves(&parquet_schema, vec![0]), + move |batch| eq(batch.column(0), &Scalar::new(&a_scalar)), + ); + + let b_scalar = StringArray::from_iter_values(["4"]); + let b_filter = ArrowPredicateFn::new( + ProjectionMask::leaves(&parquet_schema, vec![1]), + move |batch| eq(batch.column(0), &Scalar::new(&b_scalar)), + ); + + let filter = RowFilter::new(vec![Box::new(a_filter), Box::new(b_filter)]); + + let mask = ProjectionMask::leaves(&parquet_schema, vec![0, 2]); + let stream = ParquetRecordBatchStreamBuilder::new(test) + .await + .unwrap() + .with_projection(mask.clone()) + .with_batch_size(1024) + .with_row_filter(filter) + .with_row_id("_rowid") + .build() + .unwrap(); + + let batches: Vec<_> = stream.try_collect().await.unwrap(); + assert_eq!(batches.len(), 1); + + let batch = &batches[0]; + assert_eq!(batch.num_rows(), 1); + assert_eq!(batch.num_columns(), 3); + + let col = batch.column(0); + let val = col.as_any().downcast_ref::().unwrap().value(0); + assert_eq!(val, 3); + + let col = batch.column(1); + let val = col.as_any().downcast_ref::().unwrap().value(0); + assert_eq!(val, "b"); + + let col = batch.column(2); + let val = col.as_any().downcast_ref::().unwrap().value(0); + assert_eq!(val, 3); + // Should only have made 3 requests assert_eq!(requests.lock().unwrap().len(), 3); } @@ -1650,6 +2014,7 @@ mod tests { data, metadata: Default::default(), requests: Default::default(), + max_concurrent_requests: Default::default(), }; let stream = ParquetRecordBatchStreamBuilder::new(test.clone()) @@ -1724,6 +2089,125 @@ mod tests { assert_eq!(col2.values(), &[4, 5]); } + #[tokio::test] + async fn test_async_reader_with_rowid_limit_multiple_row_groups() { + let a = StringArray::from_iter_values(["a", "b", "b", "b", "c", "c"]); + let b = StringArray::from_iter_values(["1", "2", "3", "4", "5", "6"]); + let c = Int32Array::from_iter(0..6); + let data = RecordBatch::try_from_iter([ + ("a", Arc::new(a) as ArrayRef), + ("b", Arc::new(b) as ArrayRef), + ("c", Arc::new(c) as ArrayRef), + ]) + .unwrap(); + + let mut buf = Vec::with_capacity(1024); + let props = WriterProperties::builder() + .set_max_row_group_size(3) + .build(); + let mut writer = ArrowWriter::try_new(&mut buf, data.schema(), Some(props)).unwrap(); + writer.write(&data).unwrap(); + writer.close().unwrap(); + + let data: Bytes = buf.into(); + let metadata = ParquetMetaDataReader::new() + .parse_and_finish(&data) + .unwrap(); + + assert_eq!(metadata.num_row_groups(), 2); + + let test = TestReader { + data, + metadata: Some(Arc::new(metadata)), + requests: Default::default(), + max_concurrent_requests: Default::default(), + }; + + let stream = ParquetRecordBatchStreamBuilder::new(test.clone()) + .await + .unwrap() + .with_batch_size(1024) + .with_limit(4) + .with_row_id("_rowid") + .build() + .unwrap(); + + let batches: Vec<_> = stream.try_collect().await.unwrap(); + // Expect one batch for each row group + assert_eq!(batches.len(), 2); + + let batch = &batches[0]; + // First batch should contain all rows + assert_eq!(batch.num_rows(), 3); + assert_eq!(batch.num_columns(), 4); + let rowids = batch.column(0).as_primitive::(); + assert_eq!(rowids.values(), &[0, 1, 2]); + let col3 = batch.column(3).as_primitive::(); + assert_eq!(col3.values(), &[0, 1, 2]); + + let batch = &batches[1]; + // Second batch should trigger the limit and only have one row + assert_eq!(batch.num_rows(), 1); + assert_eq!(batch.num_columns(), 4); + let rowids = batch.column(0).as_primitive::(); + assert_eq!(rowids.values(), &[3]); + let col3 = batch.column(3).as_primitive::(); + assert_eq!(col3.values(), &[3]); + + let stream = ParquetRecordBatchStreamBuilder::new(test.clone()) + .await + .unwrap() + .with_offset(2) + .with_limit(3) + .with_row_id("_rowid") + .build() + .unwrap(); + + let batches: Vec<_> = stream.try_collect().await.unwrap(); + // Expect one batch for each row group + assert_eq!(batches.len(), 2); + + let batch = &batches[0]; + // First batch should contain one row + assert_eq!(batch.num_rows(), 1); + assert_eq!(batch.num_columns(), 4); + let rowids = batch.column(0).as_primitive::(); + assert_eq!(rowids.values(), &[2]); + let col3 = batch.column(3).as_primitive::(); + assert_eq!(col3.values(), &[2]); + + let batch = &batches[1]; + // Second batch should contain two rows + assert_eq!(batch.num_rows(), 2); + assert_eq!(batch.num_columns(), 4); + let rowids = batch.column(0).as_primitive::(); + assert_eq!(rowids.values(), &[3, 4]); + let col3 = batch.column(3).as_primitive::(); + assert_eq!(col3.values(), &[3, 4]); + + let stream = ParquetRecordBatchStreamBuilder::new(test.clone()) + .await + .unwrap() + .with_offset(4) + .with_limit(20) + .with_row_id("_rowid") + .build() + .unwrap(); + + let batches: Vec<_> = stream.try_collect().await.unwrap(); + // Should skip first row group + assert_eq!(batches.len(), 1); + + let batch = &batches[0]; + // First batch should contain two rows + assert_eq!(batch.num_rows(), 2); + assert_eq!(batch.num_columns(), 4); + let rowids = batch.column(0).as_primitive::(); + assert_eq!(rowids.values(), &[4, 5]); + let col3 = batch.column(3).as_primitive::(); + assert_eq!(col3.values(), &[4, 5]); + } + #[tokio::test] async fn test_row_filter_with_index() { let testdata = arrow::util::test_util::parquet_test_data(); @@ -1741,6 +2225,7 @@ mod tests { data: data.clone(), metadata: Default::default(), requests: Default::default(), + max_concurrent_requests: Default::default(), }; let a_filter = @@ -1809,6 +2294,7 @@ mod tests { data: data.clone(), metadata: Default::default(), requests: Default::default(), + max_concurrent_requests: Default::default(), }; let requests = async_reader.requests.clone(); @@ -1830,6 +2316,7 @@ mod tests { filter: None, limit: None, offset: None, + rowid: None, }; let mut skip = true; @@ -1879,6 +2366,7 @@ mod tests { data: data.clone(), metadata: Default::default(), requests: Default::default(), + max_concurrent_requests: Default::default(), }; let builder = ParquetRecordBatchStreamBuilder::new(async_reader) @@ -2022,6 +2510,7 @@ mod tests { data: data.clone(), metadata: Default::default(), requests: Default::default(), + max_concurrent_requests: Default::default(), }; let builder = ParquetRecordBatchStreamBuilder::new(async_reader) .await @@ -2049,6 +2538,7 @@ mod tests { data: data.clone(), metadata: Default::default(), requests: Default::default(), + max_concurrent_requests: Default::default(), }; let mut builder = ParquetRecordBatchStreamBuilder::new(async_reader) @@ -2192,6 +2682,7 @@ mod tests { data, metadata: Default::default(), requests: Default::default(), + max_concurrent_requests: Default::default(), }; let requests = test.requests.clone(); From 2362ba0198eb97936211a5f31be205e820b4175f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Dani=C3=ABl=20Heres?= Date: Thu, 3 Jul 2025 04:56:03 -0600 Subject: [PATCH 5/7] Ignore dict_id in try_merge (56.0.0 #7968) --- arrow-schema/src/field.rs | 13 ++++++------- 1 file changed, 6 insertions(+), 7 deletions(-) diff --git a/arrow-schema/src/field.rs b/arrow-schema/src/field.rs index dbd671a62a3a..077e0b6ac91f 100644 --- a/arrow-schema/src/field.rs +++ b/arrow-schema/src/field.rs @@ -640,13 +640,12 @@ impl Field { /// assert!(field.is_nullable()); /// ``` pub fn try_merge(&mut self, from: &Field) -> Result<(), ArrowError> { - #[allow(deprecated)] - if from.dict_id != self.dict_id { - return Err(ArrowError::SchemaError(format!( - "Fail to merge schema field '{}' because from dict_id = {} does not match {}", - self.name, from.dict_id, self.dict_id - ))); - } + // if from.dict_id != self.dict_id { + // return Err(ArrowError::SchemaError(format!( + // "Fail to merge schema field '{}' because from dict_id = {} does not match {}", + // self.name, from.dict_id, self.dict_id + // ))); + // } if from.dict_is_ordered != self.dict_is_ordered { return Err(ArrowError::SchemaError(format!( "Fail to merge schema field '{}' because from dict_is_ordered = {} does not match {}", From 1acce1e2e599318dbd6ac349f54fd0db24e0780d Mon Sep 17 00:00:00 2001 From: Georgi Krastev Date: Thu, 12 Mar 2026 17:17:01 +0200 Subject: [PATCH 6/7] Bulk RLE encode (fork-only) --- parquet/src/arrow/arrow_writer/levels.rs | 32 +++++ parquet/src/column/writer/mod.rs | 144 ++++++++++++++++++++++- parquet/src/encodings/levels.rs | 22 ++++ parquet/src/encodings/rle.rs | 90 ++++++++++++++ 4 files changed, 285 insertions(+), 3 deletions(-) diff --git a/parquet/src/arrow/arrow_writer/levels.rs b/parquet/src/arrow/arrow_writer/levels.rs index e4662b8f316c..2b118ced0b4a 100644 --- a/parquet/src/arrow/arrow_writer/levels.rs +++ b/parquet/src/arrow/arrow_writer/levels.rs @@ -566,6 +566,9 @@ pub(crate) struct ArrayLevels { /// The arrow array array: ArrayRef, + + #[allow(dead_code)] + def_levels_runs: Option>, } impl PartialEq for ArrayLevels { @@ -595,6 +598,7 @@ impl ArrayLevels { max_def_level, max_rep_level, array, + def_levels_runs: (max_def_level != 0).then(Vec::new), } } @@ -668,6 +672,7 @@ mod tests { max_def_level: 2, max_rep_level: 2, array: Arc::new(primitives), + def_levels_runs: None, }; assert_eq!(&levels[0], &expected); } @@ -688,6 +693,7 @@ mod tests { max_def_level: 0, max_rep_level: 0, array, + def_levels_runs: None, }; assert_eq!(&levels[0], &expected_levels); } @@ -714,6 +720,7 @@ mod tests { max_def_level: 1, max_rep_level: 0, array, + def_levels_runs: None, }; assert_eq!(&levels[0], &expected_levels); } @@ -748,6 +755,7 @@ mod tests { max_def_level: 1, max_rep_level: 1, array: Arc::new(leaf_array), + def_levels_runs: None, }; assert_eq!(&levels[0], &expected_levels); @@ -781,6 +789,7 @@ mod tests { max_def_level: 2, max_rep_level: 1, array: Arc::new(leaf_array), + def_levels_runs: None, }; assert_eq!(&levels[0], &expected_levels); } @@ -830,6 +839,7 @@ mod tests { max_def_level: 3, max_rep_level: 1, array: Arc::new(leaf), + def_levels_runs: None, }; assert_eq!(&levels[0], &expected_levels); @@ -880,6 +890,7 @@ mod tests { max_def_level: 5, max_rep_level: 2, array: Arc::new(leaf), + def_levels_runs: None, }; assert_eq!(&levels[0], &expected_levels); @@ -917,6 +928,7 @@ mod tests { max_def_level: 1, max_rep_level: 1, array: Arc::new(leaf), + def_levels_runs: None, }; assert_eq!(&levels[0], &expected_levels); @@ -949,6 +961,7 @@ mod tests { max_def_level: 3, max_rep_level: 1, array: Arc::new(leaf), + def_levels_runs: None, }; assert_eq!(&levels[0], &expected_levels); @@ -997,6 +1010,7 @@ mod tests { max_def_level: 5, max_rep_level: 2, array: Arc::new(leaf), + def_levels_runs: None, }; assert_eq!(&levels[0], &expected_levels); } @@ -1036,6 +1050,7 @@ mod tests { max_def_level: 3, max_rep_level: 0, array: leaf, + def_levels_runs: None, }; assert_eq!(&levels[0], &expected_levels); } @@ -1075,6 +1090,7 @@ mod tests { max_def_level: 3, max_rep_level: 1, array: Arc::new(a_values), + def_levels_runs: None, }; assert_eq!(list_level, &expected_level); } @@ -1167,6 +1183,7 @@ mod tests { max_def_level: 0, max_rep_level: 0, array: Arc::new(a), + def_levels_runs: None, }; assert_eq!(list_level, &expected_level); @@ -1180,6 +1197,7 @@ mod tests { max_def_level: 1, max_rep_level: 0, array: Arc::new(b), + def_levels_runs: None, }; assert_eq!(list_level, &expected_level); @@ -1193,6 +1211,7 @@ mod tests { max_def_level: 2, max_rep_level: 0, array: Arc::new(d), + def_levels_runs: None, }; assert_eq!(list_level, &expected_level); @@ -1206,6 +1225,7 @@ mod tests { max_def_level: 3, max_rep_level: 0, array: Arc::new(f), + def_levels_runs: None, }; assert_eq!(list_level, &expected_level); } @@ -1312,6 +1332,7 @@ mod tests { max_def_level: 1, max_rep_level: 1, array: map.keys().clone(), + def_levels_runs: None, }; assert_eq!(list_level, &expected_level); @@ -1325,6 +1346,7 @@ mod tests { max_def_level: 2, max_rep_level: 1, array: map.values().clone(), + def_levels_runs: None, }; assert_eq!(list_level, &expected_level); } @@ -1410,6 +1432,7 @@ mod tests { max_def_level: 4, max_rep_level: 1, array: values, + def_levels_runs: None, }; assert_eq!(list_level, &expected_level); @@ -1450,6 +1473,7 @@ mod tests { max_def_level: 4, max_rep_level: 1, array: values, + def_levels_runs: None, }; assert_eq!(&levels[0], &expected_level); @@ -1535,6 +1559,7 @@ mod tests { max_def_level: 6, max_rep_level: 2, array: a1_values, + def_levels_runs: None, }; assert_eq!(&levels[0], &expected_level); @@ -1546,6 +1571,7 @@ mod tests { max_def_level: 4, max_rep_level: 1, array: a2_values, + def_levels_runs: None, }; assert_eq!(&levels[1], &expected_level); @@ -1584,6 +1610,7 @@ mod tests { max_def_level: 3, max_rep_level: 1, array: values, + def_levels_runs: None, }; assert_eq!(list_level, &expected_level); } @@ -1734,6 +1761,7 @@ mod tests { max_def_level: 4, max_rep_level: 1, array: values_a, + def_levels_runs: None, }; // [[{b: 2}, null], null, [null, null], [{b: 3}, {b: 4}]] let expected_b = ArrayLevels { @@ -1743,6 +1771,7 @@ mod tests { max_def_level: 3, max_rep_level: 1, array: values_b, + def_levels_runs: None, }; assert_eq!(a_levels, &expected_a); @@ -1774,6 +1803,7 @@ mod tests { max_def_level: 3, max_rep_level: 1, array: values, + def_levels_runs: None, }; assert_eq!(list_level, &expected_level); } @@ -1809,6 +1839,7 @@ mod tests { max_def_level: 5, max_rep_level: 2, array: values, + def_levels_runs: None, }; assert_eq!(levels[0], expected_level); @@ -1839,6 +1870,7 @@ mod tests { max_def_level: 1, max_rep_level: 0, array: Arc::new(dict), + def_levels_runs: None, }; assert_eq!(levels[0], expected_level); } diff --git a/parquet/src/column/writer/mod.rs b/parquet/src/column/writer/mod.rs index 02570d3f3c69..1ea22ef29172 100644 --- a/parquet/src/column/writer/mod.rs +++ b/parquet/src/column/writer/mod.rs @@ -360,6 +360,9 @@ pub struct GenericColumnWriter<'a, E: ColumnValueEncoder> { data_page_boundary_descending: bool, /// (min, max) last_non_null_data_page_min_max: Option<(E::T, E::T)>, + + def_levels_runs_sink: Vec<(i16, usize)>, + num_levels: usize, } impl<'a, E: ColumnValueEncoder> GenericColumnWriter<'a, E> { @@ -425,6 +428,8 @@ impl<'a, E: ColumnValueEncoder> GenericColumnWriter<'a, E> { data_page_boundary_ascending: true, data_page_boundary_descending: true, last_non_null_data_page_min_max: None, + def_levels_runs_sink: vec![], + num_levels: 0, } } @@ -439,6 +444,9 @@ impl<'a, E: ColumnValueEncoder> GenericColumnWriter<'a, E> { max: Option<&E::T>, distinct_count: Option, ) -> Result { + if !self.def_levels_runs_sink.is_empty() { + self.add_data_page()?; + } // Check if number of definition levels is the same as number of repetition levels. if let (Some(def), Some(rep)) = (def_levels, rep_levels) { if def.len() != rep.len() { @@ -555,6 +563,112 @@ impl<'a, E: ColumnValueEncoder> GenericColumnWriter<'a, E> { ) } + /// Write a batch of values with a uniform definition level (bulk RLE optimization). + pub fn write_def_level_range_batch( + &mut self, + values: &E::Values, + def_level: i16, + num_levels: usize, + min: Option<&E::T>, + max: Option<&E::T>, + distinct_count: Option, + ) -> Result { + if self.descr.max_rep_level() > 0 { + return Err(general_err!( + "Cannot write def level range when rep level > 0", + )); + } + if self.statistics_enabled == EnabledStatistics::Chunk { + match (min, max) { + (Some(min), Some(max)) => { + update_min(&self.descr, min, &mut self.column_metrics.min_column_value); + update_max(&self.descr, max, &mut self.column_metrics.max_column_value); + } + (None, Some(_)) | (Some(_), None) => { + panic!("min/max should be both set or both None") + } + (None, None) => {} + }; + } + + // We can only set the distinct count if there are no other writes + if self.encoder.num_values() == 0 { + self.column_metrics.column_distinct_count = distinct_count; + } else { + self.column_metrics.column_distinct_count = None; + } + + let mut values_offset = 0; + let mut levels_offset = 0; + let base_batch_size = self.props.write_batch_size(); + while levels_offset < num_levels { + let end_offset = num_levels.min(levels_offset + base_batch_size); + + values_offset += self.write_mini_batch_with_def_level( + values, + values_offset, + None, + end_offset - levels_offset, + def_level, + )?; + levels_offset = end_offset; + } + + // Return total number of values processed. + Ok(values_offset) + } + + fn write_mini_batch_with_def_level( + &mut self, + values: &E::Values, + values_offset: usize, + value_indices: Option<&[usize]>, + num_levels: usize, + def_level: i16, + ) -> Result { + if !self.def_levels_sink.is_empty() { + self.add_data_page()?; + } + if let Some((last_def_level, last_count)) = self.def_levels_runs_sink.last_mut() { + if *last_def_level == def_level { + *last_count += num_levels; + } else { + self.def_levels_runs_sink.push((def_level, num_levels)); + } + } else { + self.def_levels_runs_sink.push((def_level, num_levels)); + } + self.num_levels += num_levels; + let values_to_write = if def_level == self.descr.max_def_level() { + num_levels + } else { + self.page_metrics.num_page_nulls += num_levels as u64; + 0 + }; + + self.page_metrics.num_buffered_rows += num_levels as u32; + + match value_indices { + Some(indices) => { + let indices = &indices[values_offset..values_offset + values_to_write]; + self.encoder.write_gather(values, indices)?; + } + None => self.encoder.write(values, values_offset, values_to_write)?, + } + + self.page_metrics.num_buffered_values += num_levels as u32; + + if self.should_add_data_page() { + self.add_data_page()?; + } + + if self.should_dict_fallback() { + self.dict_fallback()?; + } + + Ok(values_to_write) + } + /// Returns the estimated total memory usage. /// /// Unlike [`Self::get_estimated_total_bytes`] this is an estimate @@ -1008,13 +1122,21 @@ impl<'a, E: ColumnValueEncoder> GenericColumnWriter<'a, E> { } if max_def_level > 0 { - buffer.extend_from_slice( + let encoded_levels = if self.def_levels_runs_sink.is_empty() { &self.encode_levels_v1( Encoding::RLE, &self.def_levels_sink[..], max_def_level, - )[..], - ); + )[..] + } else { + &self.encode_levels_v1_bulk( + Encoding::RLE, + &self.def_levels_runs_sink[..], + self.num_levels, + max_def_level, + )[..] + }; + buffer.extend_from_slice(encoded_levels); } buffer.extend_from_slice(&values_data.buf); @@ -1095,6 +1217,8 @@ impl<'a, E: ColumnValueEncoder> GenericColumnWriter<'a, E> { self.rep_levels_sink.clear(); self.def_levels_sink.clear(); self.page_metrics.new_page(); + self.def_levels_runs_sink.clear(); + self.num_levels = 0; Ok(()) } @@ -1220,6 +1344,20 @@ impl<'a, E: ColumnValueEncoder> GenericColumnWriter<'a, E> { encoder.consume() } + /// Encodes definition or repetition levels for Data Page v1. + #[inline] + fn encode_levels_v1_bulk( + &self, + encoding: Encoding, + level_runs: &[(i16, usize)], + num_levels: usize, + max_level: i16, + ) -> Vec { + let mut encoder = LevelEncoder::v1(encoding, max_level, num_levels); + encoder.put_bulk(level_runs); + encoder.consume() + } + /// Encodes definition or repetition levels for Data Page v2. /// Encoding is always RLE. #[inline] diff --git a/parquet/src/encodings/levels.rs b/parquet/src/encodings/levels.rs index 6f662b614fca..cc8b99a75c2b 100644 --- a/parquet/src/encodings/levels.rs +++ b/parquet/src/encodings/levels.rs @@ -108,6 +108,28 @@ impl LevelEncoder { num_encoded } + /// Put/encode levels vector into this level encoder. + /// Returns number of encoded values that are less than or equal to length of the + /// input buffer. + /// Write multiple level runs as `(level, count)` pairs. + #[inline] + pub fn put_bulk(&mut self, runs: &[(i16, usize)]) -> usize { + let mut num_encoded = 0; + match *self { + LevelEncoder::Rle(ref mut encoder) | LevelEncoder::RleV2(ref mut encoder) => { + for &(value, count) in runs { + encoder.put_bulk(value as u64, count); + num_encoded += count; + } + encoder.flush(); + } + LevelEncoder::BitPacked(_, _) => { + unimplemented!() + } + } + num_encoded + } + /// Finalizes level encoder, flush all intermediate buffers and return resulting /// encoded buffer. Returned buffer is already truncated to encoded bytes only. #[inline] diff --git a/parquet/src/encodings/rle.rs b/parquet/src/encodings/rle.rs index d6e32600d321..4f34285f0cb1 100644 --- a/parquet/src/encodings/rle.rs +++ b/parquet/src/encodings/rle.rs @@ -151,6 +151,96 @@ impl RleEncoder { } } + /// Encode a single value repeated `count` times. + #[inline] + pub fn put_bulk(&mut self, value: u64, count: usize) { + assert!(count > 0, "Count must be positive"); + let remaining = 8 - self.num_buffered_values; + if self.current_value == value { + if self.repeat_count >= 8 { + self.repeat_count += count; + return; + } + let n = remaining.min(count); + for _ in 0..n { + self.buffered_values[self.num_buffered_values] = value; + self.num_buffered_values += 1; + self.repeat_count += 1; + } + if self.num_buffered_values == 8 { + assert_eq!(self.bit_packed_count % 8, 0); + self.flush_buffered_values(); + if count > n { + let mut remaining_run = count - n; + // Fill buffer + for _ in 0..remaining_run.min(8) { + self.repeat_count += 1; + remaining_run -= 1; + if self.repeat_count > 8 { + break; + } + self.buffered_values[self.num_buffered_values] = value; + self.num_buffered_values += 1; + } + if self.num_buffered_values == 8 && self.repeat_count <= 8 { + // Buffered values are full. Flush them. + assert_eq!(self.bit_packed_count % 8, 0); + self.flush_buffered_values(); + } + + self.repeat_count += remaining_run; + } + } + } else { + if self.repeat_count >= 8 { + // The current RLE run has ended and we've gathered enough. Flush first. + assert_eq!( + self.bit_packed_count, 0, + "rc = {}, value = {}, num_buffered = {}, c = {count}, {:?}", + self.repeat_count, value, self.num_buffered_values, self.buffered_values + ); + self.flush_rle_run(); + } + self.current_value = value; + let n = remaining.min(count); + self.repeat_count = 0; + for _ in 0..n { + self.buffered_values[self.num_buffered_values] = value; + self.num_buffered_values += 1; + self.repeat_count += 1; + } + if self.num_buffered_values == 8 { + let mut new_count = count; + if self.repeat_count < 8 { + new_count = count - self.repeat_count; + } + // Buffered values are full. Flush them. + assert_eq!(self.bit_packed_count % 8, 0); + self.flush_buffered_values(); + if count > n { + if self.repeat_count <= 8 { + for _ in 0..(count - n).min(8) { + self.repeat_count += 1; + if self.repeat_count > 8 { + break; + } + self.buffered_values[self.num_buffered_values] = value; + self.num_buffered_values += 1; + } + if self.num_buffered_values == 8 && self.repeat_count <= 8 { + // Buffered values are full. Flush them. + assert_eq!(self.bit_packed_count % 8, 0); + self.flush_buffered_values(); + } + } + self.repeat_count = new_count; + } + } else if count > 8 { + self.repeat_count = count; + } + } + } + #[inline] #[allow(unused)] pub fn buffer(&self) -> &[u8] { From 147b6d00c24718fb0721ff4ecffa8a4a6d98c24e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Sava=20Vrane=C5=A1evi=C4=87?= <20240220+svranesevic@users.noreply.github.com> Date: Tue, 17 Mar 2026 11:19:42 +0100 Subject: [PATCH 7/7] Expose option to set line terminator (upstream issue #9571) --- arrow-csv/src/writer.rs | 108 ++++++++++++++++++++++++++++++++++++++-- 1 file changed, 104 insertions(+), 4 deletions(-) diff --git a/arrow-csv/src/writer.rs b/arrow-csv/src/writer.rs index c5a0a0b76d59..4aaa2efed7c9 100644 --- a/arrow-csv/src/writer.rs +++ b/arrow-csv/src/writer.rs @@ -197,6 +197,8 @@ pub struct WriterBuilder { quote: u8, /// Optional escape character. Defaults to `b'\\'` escape: u8, + /// Optional line terminator. Defaults to `LF` (`\n`) + terminator: Terminator, /// Enable double quote escapes. Defaults to `true` double_quote: bool, /// Optional date format for date arrays @@ -213,6 +215,15 @@ pub struct WriterBuilder { null_value: Option, } +/// The line terminator to use when writing CSV files. +#[derive(Clone, Debug)] +pub enum Terminator { + /// Use CRLF (`\r\n`) as the line terminator + CRLF, + /// Use the specified byte character as the line terminator + Any(u8), +} + impl Default for WriterBuilder { fn default() -> Self { WriterBuilder { @@ -220,6 +231,7 @@ impl Default for WriterBuilder { has_header: true, quote: b'"', escape: b'\\', + terminator: Terminator::Any(b'\n'), double_quote: true, date_format: None, datetime_format: None, @@ -389,14 +401,32 @@ impl WriterBuilder { self.null_value.as_deref().unwrap_or(DEFAULT_NULL_VALUE) } + /// Set the CSV file's line terminator + pub fn with_line_terminator(mut self, terminator: Terminator) -> Self { + self.terminator = terminator; + self + } + + /// Get the CSV file's line terminator, defaults to `LF` (`\n`) + pub fn line_terminator(&self) -> &Terminator { + &self.terminator + } + /// Create a new `Writer` pub fn build(self, writer: W) -> Writer { let mut builder = csv::WriterBuilder::new(); + + let terminator = match self.terminator { + Terminator::CRLF => csv::Terminator::CRLF, + Terminator::Any(byte) => csv::Terminator::Any(byte), + }; + let writer = builder .delimiter(self.delimiter) .quote(self.quote) .double_quote(self.double_quote) .escape(self.escape) + .terminator(terminator) .from_writer(writer); Writer { writer, @@ -788,10 +818,80 @@ sed do eiusmod tempor,-556132.25,1,,2019-04-18T02:45:55.555,23:46:03,foo let mut buffer: Vec = vec![]; file.read_to_end(&mut buffer).unwrap(); - assert_eq!( - "c1,c2\n00:02,46:17\n00:02,\n", - String::from_utf8(buffer).unwrap() - ); + let output = String::from_utf8(buffer).unwrap(); + assert_eq!(output, "c1,c2\n00:02,46:17\n00:02,\n"); + } + + #[test] + fn test_write_csv_with_lf_terminator() { + let schema = Schema::new(vec![ + Field::new("c1", DataType::Utf8, false), + Field::new("c2", DataType::UInt32, false), + ]); + + let c1 = StringArray::from(vec!["hello", "world"]); + let c2 = PrimitiveArray::::from(vec![1, 2]); + + let batch = + RecordBatch::try_new(Arc::new(schema), vec![Arc::new(c1), Arc::new(c2)]).unwrap(); + + let mut buf = Vec::new(); + let mut writer = WriterBuilder::new() + .with_line_terminator(Terminator::Any(b'\n')) + .build(&mut buf); + writer.write(&batch).unwrap(); + drop(writer); + + let output = String::from_utf8(buf).unwrap(); + assert_eq!(output, "c1,c2\nhello,1\nworld,2\n"); + } + + #[test] + fn test_write_csv_with_crlf_terminator() { + let schema = Schema::new(vec![ + Field::new("c1", DataType::Utf8, false), + Field::new("c2", DataType::UInt32, false), + ]); + + let c1 = StringArray::from(vec!["hello", "world"]); + let c2 = PrimitiveArray::::from(vec![1, 2]); + + let batch = + RecordBatch::try_new(Arc::new(schema), vec![Arc::new(c1), Arc::new(c2)]).unwrap(); + + let mut buf = Vec::new(); + let mut writer = WriterBuilder::new() + .with_line_terminator(Terminator::CRLF) + .build(&mut buf); + writer.write(&batch).unwrap(); + drop(writer); + + let output = String::from_utf8(buf).unwrap(); + assert_eq!(output, "c1,c2\r\nhello,1\r\nworld,2\r\n"); + } + + #[test] + fn test_write_csv_with_any_terminator() { + let schema = Schema::new(vec![ + Field::new("c1", DataType::Utf8, false), + Field::new("c2", DataType::UInt32, false), + ]); + + let c1 = StringArray::from(vec!["hello", "world"]); + let c2 = PrimitiveArray::::from(vec![1, 2]); + + let batch = + RecordBatch::try_new(Arc::new(schema), vec![Arc::new(c1), Arc::new(c2)]).unwrap(); + + let mut buf = Vec::new(); + let mut writer = WriterBuilder::new() + .with_line_terminator(Terminator::Any(b'|')) + .build(&mut buf); + writer.write(&batch).unwrap(); + drop(writer); + + let output = String::from_utf8(buf).unwrap(); + assert_eq!(output, "c1,c2|hello,1|world,2|"); } #[test]