diff --git a/.github/workflows/build.yml b/.github/workflows/build.yml index 8f69d98f9..c35801b11 100644 --- a/.github/workflows/build.yml +++ b/.github/workflows/build.yml @@ -48,7 +48,7 @@ jobs: - uses: actions/checkout@v6 - name: Setup Rust - uses: dtolnay/rust-toolchain@29eef336d9b2848a0b548edc03f92a220660cdb8 + uses: dtolnay/rust-toolchain@4cda84d5c5c54efe2404f9d843567869ab1699d4 with: toolchain: "nightly" components: rustfmt @@ -152,7 +152,7 @@ jobs: path: . - name: Setup Rust - uses: dtolnay/rust-toolchain@29eef336d9b2848a0b548edc03f92a220660cdb8 + uses: dtolnay/rust-toolchain@4cda84d5c5c54efe2404f9d843567869ab1699d4 - name: Cache Cargo uses: Swatinem/rust-cache@v2 @@ -231,7 +231,7 @@ jobs: path: . - name: Setup Rust - uses: dtolnay/rust-toolchain@29eef336d9b2848a0b548edc03f92a220660cdb8 + uses: dtolnay/rust-toolchain@4cda84d5c5c54efe2404f9d843567869ab1699d4 - name: Cache Cargo uses: Swatinem/rust-cache@v2 @@ -286,7 +286,8 @@ jobs: steps: - uses: actions/checkout@v6 - - uses: dtolnay/rust-toolchain@29eef336d9b2848a0b548edc03f92a220660cdb8 + - name: Setup Rust + uses: dtolnay/rust-toolchain@4cda84d5c5c54efe2404f9d843567869ab1699d4 - run: rm LICENSE.txt - name: Download LICENSE.txt @@ -368,7 +369,8 @@ jobs: steps: - uses: actions/checkout@v6 - - uses: dtolnay/rust-toolchain@29eef336d9b2848a0b548edc03f92a220660cdb8 + - name: Setup Rust + uses: dtolnay/rust-toolchain@4cda84d5c5c54efe2404f9d843567869ab1699d4 - run: rm LICENSE.txt - name: Download LICENSE.txt diff --git a/Cargo.lock b/Cargo.lock index 15895ae19..ab222b177 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -99,9 +99,9 @@ checksum = "7c02d123df017efcdfbd739ef81735b36c5ba83ec3c59c80a9d7ecc718f92e50" [[package]] name = "arrow" -version = "58.3.0" +version = "59.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "378530e55cd479eda3c14eb345310799717e6f76d0c332041e8487022166b471" +checksum = "61d285d16bce7d0be61912f7928342b673067b6b7d7ef6cc179258ba7de1fecf" dependencies = [ "arrow-arith", "arrow-array", @@ -121,9 +121,9 @@ dependencies = [ [[package]] name = "arrow-arith" -version = "58.3.0" +version = "59.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a0ab212d2c1886e802f51c5212d78ebbcbb0bec980fff9dadc1eb8d45cd0b738" +checksum = "757ef1836251e88222542a7da2623bc1c9cb9e20afefa6db2c41e79991cd91d4" dependencies = [ "arrow-array", "arrow-buffer", @@ -135,9 +135,9 @@ dependencies = [ [[package]] name = "arrow-array" -version = "58.3.0" +version = "59.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "cfd33d3e92f207444098c75b42de99d329562be0cf686b307b097cc52b4e999e" +checksum = "bc9a4a4b2b5ecd0e04df03471661cb61f28bed3c7fd50994715129b01b2edb97" dependencies = [ "ahash", "arrow-buffer", @@ -147,6 +147,7 @@ dependencies = [ "chrono-tz", "half", "hashbrown 0.17.1", + "libc", "num-complex", "num-integer", "num-traits", @@ -154,9 +155,9 @@ dependencies = [ [[package]] name = "arrow-avro" -version = "58.3.0" +version = "59.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "049230728cd6e093088c8d231b4beede184e35cad7777c1505c0d5a8571f4376" +checksum = "9fb45cd6bd2b25c0965793b83200eaca82214273a8030fbbc2d783e4c7c65a61" dependencies = [ "arrow-array", "arrow-buffer", @@ -178,21 +179,21 @@ dependencies = [ [[package]] name = "arrow-buffer" -version = "58.3.0" +version = "59.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0c6cd424c2693bcdbc150d843dc9d4d137dd2de4782ce6df491ad11a3a0416c0" +checksum = "c12b576ef18c1deb80925a248b25ad84f419198d791b8e293fc6aaa60441fe90" dependencies = [ "bytes", "half", - "num-bigint", + "num-bigint 0.5.1", "num-traits", ] [[package]] name = "arrow-cast" -version = "58.3.0" +version = "59.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4c5aefb56a2c02e9e2b30746241058b85f8983f0fcff2ba0c6d09006e1cded7f" +checksum = "68338a9096a5dc9bc11927c58c43a8526d96bf6abd2012ef6c0c9f505991cc79" dependencies = [ "arrow-array", "arrow-buffer", @@ -201,7 +202,7 @@ dependencies = [ "arrow-schema", "arrow-select", "atoi", - "base64", + "base64 0.23.1", "chrono", "comfy-table", "half", @@ -212,9 +213,9 @@ dependencies = [ [[package]] name = "arrow-csv" -version = "58.3.0" +version = "59.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e94e8cf7e517657a52b91ea1263acf38c4ca62a84655d72458a3359b12ab97de" +checksum = "25011b52b346407d497ef0030e12b45e4f2d0cc279efc09c4f3d09106db30e36" dependencies = [ "arrow-array", "arrow-cast", @@ -227,9 +228,9 @@ dependencies = [ [[package]] name = "arrow-data" -version = "58.3.0" +version = "59.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3c88210023a2bfee1896af366309a3028fc3bcbd6515fa29a7990ee1baa08ee0" +checksum = "723fe4aeed7604e00b9883a465af4ff0a0e6c44c03e41a68c3d1cbc403e0e44d" dependencies = [ "arrow-buffer", "arrow-schema", @@ -240,9 +241,9 @@ dependencies = [ [[package]] name = "arrow-ipc" -version = "58.3.0" +version = "59.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "238438f0834483703d88896db6fe5a7138b2230debc31b34c0336c2996e3c64f" +checksum = "149437b14371f5b9ec60f5ddc751483ae99d7a7072653c0075e5e469156eea7b" dependencies = [ "arrow-array", "arrow-buffer", @@ -256,9 +257,9 @@ dependencies = [ [[package]] name = "arrow-json" -version = "58.3.0" +version = "59.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "205ca2119e6d679d5c133c6f30e68f027738d95ed948cf77677ea69c7800036b" +checksum = "f18b9123ccfec418a663f821c9a034af339711678c11ffe00d3ec07da5ff9f7e" dependencies = [ "arrow-array", "arrow-buffer", @@ -281,9 +282,9 @@ dependencies = [ [[package]] name = "arrow-ord" -version = "58.3.0" +version = "59.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1bffd8fd2579286a5d63bac898159873e5094a79009940bcb42bbfce4f19f1d0" +checksum = "e6c08dff0686cf23ca4f562803f191ccbeb726dbae6309cd4b4aaf65e0f2c979" dependencies = [ "arrow-array", "arrow-buffer", @@ -294,9 +295,9 @@ dependencies = [ [[package]] name = "arrow-pyarrow" -version = "58.3.0" +version = "59.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d29abdf672a81c1aeb57fd2661457f9918964d49aed0e9f18932535f2a9e49ce" +checksum = "c196ecc25b3a8dcbc1d842f2619cee653dcfa2fb8b56a291bc0481c3cf5c3821" dependencies = [ "arrow-array", "arrow-data", @@ -306,9 +307,9 @@ dependencies = [ [[package]] name = "arrow-row" -version = "58.3.0" +version = "59.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "bab5994731204603c73ba69267616c50f80780774c6bb0476f1f830625115e0c" +checksum = "bbec439386df71ad570e6758a946111322b9e9dc8db83b5527321f0b4c9119c2" dependencies = [ "arrow-array", "arrow-buffer", @@ -319,9 +320,9 @@ dependencies = [ [[package]] name = "arrow-schema" -version = "58.3.0" +version = "59.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f633dbfdf39c039ada1bf9e34c694816eb71fbb7dc78f613993b7245e078a1ed" +checksum = "e6fed2ca0d1eade57e811cbe73b98ad50cc08a1183e13b2d2aa43a7df593f40e" dependencies = [ "bitflags", "serde_core", @@ -330,9 +331,9 @@ dependencies = [ [[package]] name = "arrow-select" -version = "58.3.0" +version = "59.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8cd065c54172ac787cf3f2f8d4107e0d3fdc26edba76fdf4f4cc170258942222" +checksum = "466b19cf75130b891dc1b23a84b343c714c62c64c9c62e365c76aa0ff90a53fb" dependencies = [ "ahash", "arrow-array", @@ -344,9 +345,9 @@ dependencies = [ [[package]] name = "arrow-string" -version = "58.3.0" +version = "59.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "29dd7cda3ab9692f43a2e4acc444d760cc17b12bb6d8232ddf64e9bab7c06b42" +checksum = "c838a25bb3691e919e0f617616ac51a4ff8517a952e29ca133cf0c22b2ce65b1" dependencies = [ "arrow-array", "arrow-buffer", @@ -385,7 +386,7 @@ checksum = "3b43422f69d8ff38f95f1b2bb76517c91589a924d1559a0e935d7c8ce0274c11" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.118", ] [[package]] @@ -396,7 +397,7 @@ checksum = "9035ad2d096bed7955a320ee7e2230574d28fd3c3a0f186cbea1ff3c7eed5dbb" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.118", ] [[package]] @@ -426,6 +427,12 @@ version = "0.22.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "72b3254f16251a8381aa12e40e3c4d2f0199f8c6508fbecb9d91f575e0fbb8c6" +[[package]] +name = "base64" +version = "0.23.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ac07cdecf99051d9a5238b80f35af32cdeba5b336e55d957b318b50137e18da5" + [[package]] name = "bigdecimal" version = "0.4.10" @@ -434,7 +441,7 @@ checksum = "4d6867f1565b3aad85681f1015055b087fcfd840d6aeee6eee7f2da317603695" dependencies = [ "autocfg", "libm", - "num-bigint", + "num-bigint 0.4.6", "num-integer", "num-traits", ] @@ -513,12 +520,6 @@ version = "3.20.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "72f5acc6cb2ba439de613abc23857ec3d78374d8ed5ac84e9d11336e87da8649" -[[package]] -name = "byteorder" -version = "1.5.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1fd0f2584146f6f2ef48085050886acf353beff7305ebd1ae69500e27c67f64b" - [[package]] name = "bytes" version = "1.11.1" @@ -790,9 +791,8 @@ dependencies = [ [[package]] name = "datafusion" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "997a31e15872606a49478e670c58302094c97cb96abb0a7d60720f8e92170040" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "arrow", "arrow-schema", @@ -828,7 +828,7 @@ dependencies = [ "flate2", "futures", "indexmap", - "itertools", + "itertools 0.15.0", "liblzma", "log", "object_store", @@ -844,9 +844,8 @@ dependencies = [ [[package]] name = "datafusion-catalog" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f7dd61161508f8f5fa1107774ea687bd753c22d83a32eebf963549f89de14139" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "arrow", "async-trait", @@ -860,7 +859,7 @@ dependencies = [ "datafusion-physical-plan", "datafusion-session", "futures", - "itertools", + "itertools 0.15.0", "log", "object_store", "parking_lot", @@ -869,9 +868,8 @@ dependencies = [ [[package]] name = "datafusion-catalog-listing" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "897c70f871277f9ce99aa38347be0d679bbe3e617156c4d2a8378cec8a2a0891" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "arrow", "async-trait", @@ -885,16 +883,16 @@ dependencies = [ "datafusion-physical-expr-common", "datafusion-physical-plan", "futures", - "itertools", + "itertools 0.15.0", "log", "object_store", + "percent-encoding", ] [[package]] name = "datafusion-common" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "121c9ded5d87d9172319e006f2afdb9928d72dbacd6a90a458d8acb1e3b43a65" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "arrow", "arrow-ipc", @@ -904,9 +902,10 @@ dependencies = [ "half", "hashbrown 0.17.1", "indexmap", - "itertools", + "itertools 0.15.0", "libc", "log", + "num-traits", "object_store", "parquet", "recursive", @@ -918,9 +917,8 @@ dependencies = [ [[package]] name = "datafusion-common-runtime" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "981b9dae74f78ee3d9f714fb49b01919eab975461b56149510c3ba9ea11287d1" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "futures", "log", @@ -929,9 +927,8 @@ dependencies = [ [[package]] name = "datafusion-datasource" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "ffd7d295b2ec7c00d8a56562f41ed41062cf0af75549ed891c12a0a09eddfefe" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "arrow", "async-compression", @@ -947,11 +944,12 @@ dependencies = [ "datafusion-physical-expr-adapter", "datafusion-physical-expr-common", "datafusion-physical-plan", + "datafusion-proto-models", "datafusion-session", "flate2", "futures", "glob", - "itertools", + "itertools 0.15.0", "liblzma", "log", "object_store", @@ -965,9 +963,8 @@ dependencies = [ [[package]] name = "datafusion-datasource-arrow" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "552b0b3f342f7ec41b3fbd70f6339dc82a30cfd0349e7f280e7852528085349f" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "arrow", "arrow-ipc", @@ -980,18 +977,18 @@ dependencies = [ "datafusion-expr", "datafusion-physical-expr-common", "datafusion-physical-plan", + "datafusion-proto-models", "datafusion-session", "futures", - "itertools", + "itertools 0.15.0", "object_store", "tokio", ] [[package]] name = "datafusion-datasource-avro" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "fb517d08967d536284ce70afb5fe8583133779249f2d7b90587d339741a7f195" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "arrow", "arrow-avro", @@ -1008,9 +1005,8 @@ dependencies = [ [[package]] name = "datafusion-datasource-csv" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "68850aa426b897e879c8b87e512ea8124f1d0a2869a4e51808ddaaddf1bc0ada" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "arrow", "async-trait", @@ -1022,6 +1018,7 @@ dependencies = [ "datafusion-expr", "datafusion-physical-expr-common", "datafusion-physical-plan", + "datafusion-proto-models", "datafusion-session", "futures", "object_store", @@ -1031,9 +1028,8 @@ dependencies = [ [[package]] name = "datafusion-datasource-json" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "402f93242ae08ef99139ee2c528a49d087efe88d5c7b2c3ff5480855a40ce54f" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "arrow", "async-trait", @@ -1045,6 +1041,7 @@ dependencies = [ "datafusion-expr", "datafusion-physical-expr-common", "datafusion-physical-plan", + "datafusion-proto-models", "datafusion-session", "futures", "object_store", @@ -1054,11 +1051,11 @@ dependencies = [ [[package]] name = "datafusion-datasource-parquet" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "ffd2499c1bee0eeccf6a57156105700eeeb17bc701899ac719183c4e74231450" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "arrow", + "arrow-schema", "async-trait", "bytes", "datafusion-common", @@ -1072,10 +1069,11 @@ dependencies = [ "datafusion-physical-expr-adapter", "datafusion-physical-expr-common", "datafusion-physical-plan", + "datafusion-proto-models", "datafusion-pruning", "datafusion-session", "futures", - "itertools", + "itertools 0.15.0", "log", "object_store", "parking_lot", @@ -1085,19 +1083,18 @@ dependencies = [ [[package]] name = "datafusion-doc" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "cb9e7e5d11130c48c8bd4e80c79a9772dd28ce6dc330baca9246205d245b9e2e" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" [[package]] name = "datafusion-execution" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "37a8643ab852eb68864e1b72ae789e8066282dce48eea6347ffb0aee33d1ccc0" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "arrow", "arrow-buffer", "async-trait", + "bytes", "dashmap", "datafusion-common", "datafusion-expr", @@ -1106,16 +1103,18 @@ dependencies = [ "log", "object_store", "parking_lot", + "pin-project-lite", "rand 0.9.4", "tempfile", + "tokio", + "tokio-util", "url", ] [[package]] name = "datafusion-expr" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6932f4d71eed9c8d9341476a2b845aadfabde5495d08dbcd8fc23881f49fa7a0" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "arrow", "arrow-schema", @@ -1127,8 +1126,10 @@ dependencies = [ "datafusion-functions-aggregate-common", "datafusion-functions-window-common", "datafusion-physical-expr-common", + "datafusion-proto-common", + "datafusion-proto-models", "indexmap", - "itertools", + "itertools 0.15.0", "recursive", "serde_json", "sqlparser", @@ -1136,21 +1137,19 @@ dependencies = [ [[package]] name = "datafusion-expr-common" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0225491839a31b1f7d2cb8092c2d50792e2fe1c1724e4e6d08e011f5feaf4ed2" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "arrow", "datafusion-common", "indexmap", - "itertools", + "itertools 0.15.0", ] [[package]] name = "datafusion-ffi" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e5660e8fa79fd51e29ce46f3026b67317ef738ebd633e106beb1a1907a406152" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "arrow", "arrow-schema", @@ -1203,13 +1202,12 @@ dependencies = [ [[package]] name = "datafusion-functions" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "14872c47bfc3d21e53ec82f57074e6987a15941c1e2f43cde4ac6ae2746634e3" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "arrow", "arrow-buffer", - "base64", + "base64 0.23.1", "blake2", "blake3", "chrono", @@ -1222,7 +1220,7 @@ dependencies = [ "datafusion-macros", "datafusion-physical-expr-common", "hex", - "itertools", + "itertools 0.15.0", "log", "md-5 0.11.0", "memchr", @@ -1235,9 +1233,8 @@ dependencies = [ [[package]] name = "datafusion-functions-aggregate" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "75a2ca14e1b609be21e657e2d3130b2f446456b08393b377bb721a33952d2e09" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "arrow", "datafusion-common", @@ -1248,17 +1245,16 @@ dependencies = [ "datafusion-macros", "datafusion-physical-expr", "datafusion-physical-expr-common", - "foldhash 0.2.0", "half", + "hashbrown 0.17.1", "log", "num-traits", ] [[package]] name = "datafusion-functions-aggregate-common" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1ece74ba09092d2ef9c9b54a38445450aea292a1f8b04faf531936b723a24b3c" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "arrow", "datafusion-common", @@ -1268,9 +1264,8 @@ dependencies = [ [[package]] name = "datafusion-functions-nested" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3f3e3f9ee8ca59bf70518802107de6f1b88a9509efdc629fadc5de9d6b2d5ef5" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "arrow", "arrow-ord", @@ -1285,7 +1280,7 @@ dependencies = [ "datafusion-macros", "datafusion-physical-expr-common", "hashbrown 0.17.1", - "itertools", + "itertools 0.15.0", "itoa", "log", "memchr", @@ -1293,9 +1288,8 @@ dependencies = [ [[package]] name = "datafusion-functions-table" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "89161dffc22cf2b50f9f4b1bee83b5221d3b4ed7c2e37fd7aa2b22a5297b3a26" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "arrow", "async-trait", @@ -1309,9 +1303,8 @@ dependencies = [ [[package]] name = "datafusion-functions-window" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d7339345b226b3874037708bf5023ba1c2de705128f8457a095aae5ae9cb9c78" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "arrow", "datafusion-common", @@ -1326,9 +1319,8 @@ dependencies = [ [[package]] name = "datafusion-functions-window-common" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "fa84836dc2392df6f43d6a29d37fb56a8ebdc8b3f4e10ae8dc15861fd20278fb" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "datafusion-common", "datafusion-physical-expr-common", @@ -1336,20 +1328,18 @@ dependencies = [ [[package]] name = "datafusion-macros" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "587164e03ad68732aa9e7bfe5686e3f25970d4c64fd4bd80790749840892dae5" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "datafusion-doc", "quote", - "syn", + "syn 3.0.3", ] [[package]] name = "datafusion-optimizer" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "77f20e8cf9e8654d92f4c16b24c487353ee5bf153ffc12d5772cd399ab8cd281" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "arrow", "chrono", @@ -1358,7 +1348,7 @@ dependencies = [ "datafusion-expr-common", "datafusion-physical-expr", "indexmap", - "itertools", + "itertools 0.15.0", "log", "recursive", "regex", @@ -1367,9 +1357,8 @@ dependencies = [ [[package]] name = "datafusion-physical-expr" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f015a4a82f6f7ff7e1d8d4bf3870a936752fa38b17705dfcc14adef95aa8922c" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "arrow", "datafusion-common", @@ -1377,10 +1366,11 @@ dependencies = [ "datafusion-expr-common", "datafusion-functions-aggregate-common", "datafusion-physical-expr-common", + "datafusion-proto-models", "half", "hashbrown 0.17.1", "indexmap", - "itertools", + "itertools 0.15.0", "parking_lot", "petgraph", "recursive", @@ -1389,9 +1379,8 @@ dependencies = [ [[package]] name = "datafusion-physical-expr-adapter" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "51e6ffff8acdfe54e0ea15ccf38115c4a9184433b0439f42907637928d00a235" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "arrow", "datafusion-common", @@ -1399,31 +1388,30 @@ dependencies = [ "datafusion-functions", "datafusion-physical-expr", "datafusion-physical-expr-common", - "itertools", + "itertools 0.15.0", ] [[package]] name = "datafusion-physical-expr-common" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7967a3e171c6a4bf09474b3f7a14f1a3db13ed1714ba12156f33fcce2bba54e8" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "arrow", "chrono", "datafusion-common", "datafusion-expr-common", + "datafusion-proto-models", "hashbrown 0.17.1", "indexmap", - "itertools", + "itertools 0.15.0", "parking_lot", "pin-project", ] [[package]] name = "datafusion-physical-optimizer" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "59ff803e2a96054cb6d83f35f9e60fd4f42eac515e1932bd1b2dbc91d5fcbf36" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "arrow", "datafusion-common", @@ -1434,15 +1422,15 @@ dependencies = [ "datafusion-physical-expr-common", "datafusion-physical-plan", "datafusion-pruning", - "itertools", + "datafusion-session", + "itertools 0.15.0", "recursive", ] [[package]] name = "datafusion-physical-plan" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "776ee54d47d15bdb126452f9ca17b03761e3b004682914beaedd3f86eb507fbc" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "arrow", "arrow-data", @@ -1450,6 +1438,7 @@ dependencies = [ "arrow-ord", "arrow-schema", "async-trait", + "bytes", "datafusion-common", "datafusion-common-runtime", "datafusion-execution", @@ -1459,26 +1448,27 @@ dependencies = [ "datafusion-functions-window-common", "datafusion-physical-expr", "datafusion-physical-expr-common", + "datafusion-proto-common", + "datafusion-proto-models", "futures", "half", "hashbrown 0.17.1", "indexmap", - "itertools", + "itertools 0.15.0", "log", "num-traits", "parking_lot", "pin-project-lite", + "serde_json", "tokio", ] [[package]] name = "datafusion-proto" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9dd15a1ba5d3af93808241065c6c44dbca8296a189845e8a587c45c07bf0ffae" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "arrow", - "chrono", "datafusion-catalog", "datafusion-catalog-listing", "datafusion-common", @@ -1494,26 +1484,35 @@ dependencies = [ "datafusion-physical-expr-common", "datafusion-physical-plan", "datafusion-proto-common", + "datafusion-proto-models", "object_store", "prost", ] [[package]] name = "datafusion-proto-common" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "90042982cf9462eb06a0b81f92efa4188dae871e7ea3ab8dc61aa9c9349b2530" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "arrow", "datafusion-common", "prost", ] +[[package]] +name = "datafusion-proto-models" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" +dependencies = [ + "datafusion-common", + "datafusion-proto-common", + "prost", +] + [[package]] name = "datafusion-pruning" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d5fb9e5774660aa69c3ba93c610f175f75b65cb8c3776edb3626de8f3a4f4ee3" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "arrow", "datafusion-common", @@ -1572,10 +1571,10 @@ dependencies = [ [[package]] name = "datafusion-session" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "15ce715fa2a61f4623cc234bcc14a3ef6a91f189128d5b14b468a6a17cdfc417" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ + "arrow-schema", "async-trait", "datafusion-common", "datafusion-execution", @@ -1586,14 +1585,14 @@ dependencies = [ [[package]] name = "datafusion-spark" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "390bb0bf37cb2b95ffd65eacd66f60df50793d1f94097799e416f39477a51957" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "arrow", "bigdecimal", "chrono", "crc32fast", + "datafusion", "datafusion-catalog", "datafusion-common", "datafusion-execution", @@ -1615,9 +1614,8 @@ dependencies = [ [[package]] name = "datafusion-sql" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6094ad36a3ed6d7ac87b20b479b2d0b118250f66cf997603829fdc65b44a7099" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "arrow", "bigdecimal", @@ -1630,20 +1628,20 @@ dependencies = [ "recursive", "regex", "sqlparser", + "stacker", ] [[package]] name = "datafusion-substrait" -version = "54.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3b22c8f8c72d317e54fad6f85c0ef6d1e1da53cc7faadc7eea8daf0f8d86d4f2" +version = "55.0.0" +source = "git+https://github.com/apache/datafusion?rev=55.0.0-rc3#d5552342012888b7d1a3ab88d92e3d292fc0cde0" dependencies = [ "async-recursion", "async-trait", "chrono", "datafusion", "half", - "itertools", + "itertools 0.15.0", "object_store", "pbjson-types", "prost", @@ -1682,7 +1680,7 @@ checksum = "1ac70aa55017e108007fbaf5aa0f54b021c98f92ff8af59d42eda9da96e3dd4f" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.118", ] [[package]] @@ -1835,7 +1833,7 @@ checksum = "e835b70203e41293343137df5c0664546da5745f82ec9b84d40be8336958447b" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.118", ] [[package]] @@ -2099,7 +2097,7 @@ version = "0.1.20" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "96547c2556ec9d12fb1578c4eaf448b04993e7fb79cbaad930a656880a6bdfa0" dependencies = [ - "base64", + "base64 0.22.1", "bytes", "futures-channel", "futures-util", @@ -2255,12 +2253,6 @@ dependencies = [ "serde_core", ] -[[package]] -name = "integer-encoding" -version = "3.0.4" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8bb03732005da905c88227371639bf1ad885cc712789c011c31c5fb3ab3ccf02" - [[package]] name = "ipnet" version = "2.12.0" @@ -2276,6 +2268,15 @@ dependencies = [ "either", ] +[[package]] +name = "itertools" +version = "0.15.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8b4baf93f58d4425749ca49a51c50ebab072c5df6994d08fed93541c331481dc" +dependencies = [ + "either", +] + [[package]] name = "itoa" version = "1.0.18" @@ -2452,9 +2453,9 @@ checksum = "112b39cec0b298b6c1999fee3e31427f74f676e4cb9879ed1a121b43661a4154" [[package]] name = "lz4_flex" -version = "0.13.1" +version = "0.14.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7ef0d4ed8669f8f8826eb00dc878084aa8f253506c4fd5e8f58f5bce72ddb97e" +checksum = "ecbdfe44b1bd960b68170b417450a628c43f7cf56bb3c5317e61cb230ee7f226" dependencies = [ "twox-hash", ] @@ -2531,6 +2532,16 @@ dependencies = [ "num-traits", ] +[[package]] +name = "num-bigint" +version = "0.5.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "93e7820bc0a80a0238e650327316f929ba18d5be054b647490a3a6a339f3e7c0" +dependencies = [ + "num-integer", + "num-traits", +] + [[package]] name = "num-complex" version = "0.4.6" @@ -2575,7 +2586,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "622acbc9100d3c10e2ee15804b0caa40e55c933d5aa53814cd520805b7958a49" dependencies = [ "async-trait", - "base64", + "base64 0.22.1", "bytes", "chrono", "form_urlencoded", @@ -2587,7 +2598,7 @@ dependencies = [ "httparse", "humantime", "hyper", - "itertools", + "itertools 0.14.0", "md-5 0.10.6", "parking_lot", "percent-encoding", @@ -2620,15 +2631,6 @@ version = "0.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7c87def4c32ab89d880effc9e097653c8da5d6ef28e6b539d313baaacfbafcbe" -[[package]] -name = "ordered-float" -version = "2.10.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "68f19d67e5a2795c94e73e0bb1cc1a7edeb2e28efd39e2e1c9b7a40c1108b11c" -dependencies = [ - "num-traits", -] - [[package]] name = "parking_lot" version = "0.12.5" @@ -2654,9 +2656,9 @@ dependencies = [ [[package]] name = "parquet" -version = "58.3.0" +version = "59.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5dafa7d01085b62a47dd0c1829550a0a36710ea9c4fe358a05a85477cec8a908" +checksum = "7065842956a20c2a536924ce8e4d9955f7422451511b9eb7500d7bfe5077e59c" dependencies = [ "ahash", "arrow-array", @@ -2665,7 +2667,7 @@ dependencies = [ "arrow-ipc", "arrow-schema", "arrow-select", - "base64", + "base64 0.23.1", "brotli", "bytes", "chrono", @@ -2674,33 +2676,25 @@ dependencies = [ "half", "hashbrown 0.17.1", "lz4_flex", - "num-bigint", + "num-bigint 0.5.1", "num-integer", "num-traits", "object_store", - "paste", "seq-macro", "simdutf8", "snap", - "thrift", "tokio", "twox-hash", "zstd", ] -[[package]] -name = "paste" -version = "1.0.15" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "57c0d7b74b563b49d38dae00a0c37d4d6de9b432382b2892f0574ddcae73fd0a" - [[package]] name = "pbjson" version = "0.8.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "898bac3fa00d0ba57a4e8289837e965baa2dee8c3749f3b11d45a64b4223d9c3" dependencies = [ - "base64", + "base64 0.22.1", "serde", ] @@ -2711,7 +2705,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "af22d08a625a2213a78dbb0ffa253318c5c79ce3133d32d296655a7bdfb02095" dependencies = [ "heck", - "itertools", + "itertools 0.14.0", "prost", "prost-types", ] @@ -2784,7 +2778,7 @@ checksum = "c96395f0a926bc13b1c17622aaddda1ecb55d49c8f1bf9777e4d877800a43f8b" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.118", ] [[package]] @@ -2830,7 +2824,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "479ca8adacdd7ce8f1fb39ce9ecccbfe93a3f1344b3d0d97f20bc0196208f62b" dependencies = [ "proc-macro2", - "syn", + "syn 2.0.118", ] [[package]] @@ -2868,7 +2862,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "03da047801ff44bb6a4d407d4860c05fd70bb81714e6b2f3812603d5b145b042" dependencies = [ "heck", - "itertools", + "itertools 0.14.0", "log", "multimap", "petgraph", @@ -2876,7 +2870,7 @@ dependencies = [ "prost", "prost-types", "regex", - "syn", + "syn 2.0.118", "tempfile", ] @@ -2887,10 +2881,10 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b570b25f7617e43d59005d0990ccb79e950a423952cea19671b7a876da390adf" dependencies = [ "anyhow", - "itertools", + "itertools 0.14.0", "proc-macro2", "quote", - "syn", + "syn 2.0.118", ] [[package]] @@ -2923,9 +2917,9 @@ dependencies = [ [[package]] name = "pyo3" -version = "0.28.3" +version = "0.29.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "91fd8e38a3b50ed1167fb981cd6fd60147e091784c427b8f7183a7ee32c31c12" +checksum = "cd274650b21d4bfc26a0a47587962c1edb425f69287324355cd040c3ea66071c" dependencies = [ "libc", "once_cell", @@ -2937,9 +2931,9 @@ dependencies = [ [[package]] name = "pyo3-async-runtimes" -version = "0.28.0" +version = "0.29.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9e7364a95bf00e8377bbf9b0f09d7ff9715a29d8fcf93b47d1a967363b973178" +checksum = "b3ef68daa7316a3fac65e5e18b2203f010346de1c1c53456811a2624673ab046" dependencies = [ "futures-channel", "futures-util", @@ -2951,19 +2945,18 @@ dependencies = [ [[package]] name = "pyo3-build-config" -version = "0.28.3" +version = "0.29.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e368e7ddfdeb98c9bca7f8383be1648fd84ab466bf2bc015e94008db6d35611e" +checksum = "c5e2a7d2f0d013342f295c048ad19237add5154a55b1c5a254c0ec93d4109078" dependencies = [ - "python3-dll-a", "target-lexicon", ] [[package]] name = "pyo3-ffi" -version = "0.28.3" +version = "0.29.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7f29e10af80b1f7ccaf7f69eace800a03ecd13e883acfacc1e5d0988605f651e" +checksum = "ca85c467da1bbc8d866eea5deff9cf29ea5f7785054a17da36e65bda9c05845b" dependencies = [ "libc", "pyo3-build-config", @@ -2982,36 +2975,26 @@ dependencies = [ [[package]] name = "pyo3-macros" -version = "0.28.3" +version = "0.29.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "df6e520eff47c45997d2fc7dd8214b25dd1310918bbb2642156ef66a67f29813" +checksum = "9ac53762fd065daa3194dd09337a38bd793a188100fd1a9304c4ab312d901771" dependencies = [ "proc-macro2", "pyo3-macros-backend", "quote", - "syn", + "syn 2.0.118", ] [[package]] name = "pyo3-macros-backend" -version = "0.28.3" +version = "0.29.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c4cdc218d835738f81c2338f822078af45b4afdf8b2e33cbb5916f108b813acb" +checksum = "4ca3a1557399783172dc5bf39cfca835157732532cba56b71d2292161e53b362" dependencies = [ "heck", "proc-macro2", - "pyo3-build-config", "quote", - "syn", -] - -[[package]] -name = "python3-dll-a" -version = "0.2.15" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d80ba7540edb18890d444c5aa8e1f1f99b1bdf26fb26ae383135325f4a36042b" -dependencies = [ - "cc", + "syn 2.0.118", ] [[package]] @@ -3163,7 +3146,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "76009fbe0614077fc1a2ce255e3a1881a2e3a3527097d5dc6d8212c585e7e38b" dependencies = [ "quote", - "syn", + "syn 2.0.118", ] [[package]] @@ -3220,7 +3203,7 @@ version = "0.12.28" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "eddd3ca559203180a307f12d114c268abf583f59b03cb906fd0b3ff8646c1147" dependencies = [ - "base64", + "base64 0.22.1", "bytes", "futures-core", "futures-util", @@ -3396,7 +3379,7 @@ dependencies = [ "proc-macro2", "quote", "serde_derive_internals", - "syn", + "syn 2.0.118", ] [[package]] @@ -3471,7 +3454,7 @@ checksum = "d540f220d3187173da220f885ab66608367b6574e925011a9353e4badda91d79" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.118", ] [[package]] @@ -3482,7 +3465,7 @@ checksum = "18d26a20a969b9e3fdf2fc2d9f21eda6c40e2de84c9408bb5d3b05d499aae711" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.118", ] [[package]] @@ -3508,7 +3491,7 @@ dependencies = [ "proc-macro2", "quote", "serde", - "syn", + "syn 2.0.118", ] [[package]] @@ -3635,7 +3618,7 @@ checksum = "a6dd45d8fc1c79299bfbb7190e42ccbbdf6a5f52e4a6ad98d92357ea965bd289" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.118", ] [[package]] @@ -3669,7 +3652,7 @@ dependencies = [ "proc-macro-crate", "proc-macro2", "quote", - "syn", + "syn 2.0.118", ] [[package]] @@ -3700,7 +3683,7 @@ dependencies = [ "heck", "proc-macro2", "quote", - "syn", + "syn 2.0.118", ] [[package]] @@ -3725,7 +3708,7 @@ dependencies = [ "serde", "serde_json", "serde_yaml", - "syn", + "syn 2.0.118", "typify", "walkdir", ] @@ -3747,6 +3730,17 @@ dependencies = [ "unicode-ident", ] +[[package]] +name = "syn" +version = "3.0.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "53e9bae58849f64dfa4f5d5ae372c8341f7305f82a3868709269343628b659a3" +dependencies = [ + "proc-macro2", + "quote", + "unicode-ident", +] + [[package]] name = "sync_wrapper" version = "1.0.2" @@ -3764,7 +3758,7 @@ checksum = "728a70f3dbaf5bab7f0c4b1ac8d7ae5ea60a4b5549c8a5914361c99147a709d2" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.118", ] [[package]] @@ -3803,18 +3797,7 @@ checksum = "ebc4ee7f67670e9b64d05fa4253e753e016c6c95ff35b89b7941d6b856dec1d5" dependencies = [ "proc-macro2", "quote", - "syn", -] - -[[package]] -name = "thrift" -version = "0.17.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7e54bc85fc7faa8bc175c4bab5b92ba8d9a3ce893d0e9f42cc455c8ab16a9e09" -dependencies = [ - "byteorder", - "integer-encoding", - "ordered-float", + "syn 2.0.118", ] [[package]] @@ -3874,7 +3857,7 @@ checksum = "385a6cb71ab9ab790c5fe8d67f1645e6c450a7ce006a33de03daa956cf70a496" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.118", ] [[package]] @@ -4006,7 +3989,7 @@ checksum = "7490cfa5ec963746568740651ac6781f701c9c5ea257c58e057f3ba8cf69e8da" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.118", ] [[package]] @@ -4064,7 +4047,7 @@ dependencies = [ "semver", "serde", "serde_json", - "syn", + "syn 2.0.118", "thiserror", "unicode-ident", ] @@ -4082,7 +4065,7 @@ dependencies = [ "serde", "serde_json", "serde_tokenstream", - "syn", + "syn 2.0.118", "typify-impl", ] @@ -4227,7 +4210,7 @@ dependencies = [ "bumpalo", "proc-macro2", "quote", - "syn", + "syn 2.0.118", "wasm-bindgen-shared", ] @@ -4303,7 +4286,7 @@ checksum = "053e2e040ab57b9dc951b72c264860db7eb3b0200ba345b4e4c3b14f67855ddf" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.118", ] [[package]] @@ -4314,7 +4297,7 @@ checksum = "3f316c4a2570ba26bbec722032c4099d8c8bc095efccdc15688708623367e358" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.118", ] [[package]] @@ -4537,7 +4520,7 @@ checksum = "de844c262c8848816172cef550288e7dc6c7b7814b4ee56b3e1553f275f1858e" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.118", "synstructure", ] @@ -4558,7 +4541,7 @@ checksum = "1ae7f38b72ec2a254e2b87ef277cf2cd4fb97cbebf944faa6f33354da0867930" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.118", ] [[package]] @@ -4578,7 +4561,7 @@ checksum = "11532158c46691caf0f2593ea8358fed6bbf68a0315e80aae9bd41fbade684a1" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.118", "synstructure", ] @@ -4618,7 +4601,7 @@ checksum = "625dc425cab0dca6dc3c3319506e6593dcb08a9f387ea3b284dbd52a92c40555" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.118", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index b37d65331..a9e15d7e9 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -32,24 +32,24 @@ resolver = "3" [workspace.dependencies] tokio = { version = "1.52" } -pyo3 = { version = "0.28" } -pyo3-async-runtimes = { version = "0.28" } +pyo3 = { version = "0.29" } +pyo3-async-runtimes = { version = "0.29" } pyo3-log = "0.13.3" chrono = { version = "0.4", default-features = false } -arrow = { version = "58" } -arrow-array = { version = "58" } -arrow-schema = { version = "58" } -arrow-select = { version = "58" } -datafusion = { version = "54" } -datafusion-substrait = { version = "54" } -datafusion-proto = { version = "54" } -datafusion-ffi = { version = "54" } -datafusion-catalog = { version = "54", default-features = false } -datafusion-common = { version = "54", default-features = false } -datafusion-functions-aggregate = { version = "54" } -datafusion-functions-window = { version = "54" } -datafusion-spark = { version = "54" } -datafusion-expr = { version = "54" } +arrow = { version = "59" } +arrow-array = { version = "59" } +arrow-schema = { version = "59" } +arrow-select = { version = "59" } +datafusion = { version = "55.0.0" } +datafusion-substrait = { version = "55.0.0" } +datafusion-proto = { version = "55.0.0" } +datafusion-ffi = { version = "55.0.0" } +datafusion-catalog = { version = "55.0.0", default-features = false } +datafusion-common = { version = "55.0.0", default-features = false } +datafusion-functions-aggregate = { version = "55.0.0" } +datafusion-functions-window = { version = "55.0.0" } +datafusion-spark = { version = "55.0.0" } +datafusion-expr = { version = "55.0.0" } prost = "0.14.3" serde_json = "1" uuid = { version = "1.23" } @@ -62,7 +62,7 @@ url = "2" log = "0.4.29" parking_lot = "0.12" prost-types = "0.14.3" # keep in line with `datafusion-substrait` -pyo3-build-config = "0.28" +pyo3-build-config = "0.29" datafusion-python-util = { path = "crates/util", version = "54.0.0" } [profile.release] @@ -72,3 +72,13 @@ codegen-units = 2 # We cannot publish to crates.io with any patches in the below section. Developers # must remove any entries in this section before creating a release candidate. [patch.crates-io] +datafusion = { git = "https://github.com/apache/datafusion", rev = "55.0.0-rc3" } +datafusion-substrait = { git = "https://github.com/apache/datafusion", rev = "55.0.0-rc3" } +datafusion-proto = { git = "https://github.com/apache/datafusion", rev = "55.0.0-rc3" } +datafusion-ffi = { git = "https://github.com/apache/datafusion", rev = "55.0.0-rc3" } +datafusion-catalog = { git = "https://github.com/apache/datafusion", rev = "55.0.0-rc3" } +datafusion-common = { git = "https://github.com/apache/datafusion", rev = "55.0.0-rc3" } +datafusion-functions-aggregate = { git = "https://github.com/apache/datafusion", rev = "55.0.0-rc3" } +datafusion-functions-window = { git = "https://github.com/apache/datafusion", rev = "55.0.0-rc3" } +datafusion-spark = { git = "https://github.com/apache/datafusion", rev = "55.0.0-rc3" } +datafusion-expr = { git = "https://github.com/apache/datafusion", rev = "55.0.0-rc3" } diff --git a/crates/core/Cargo.toml b/crates/core/Cargo.toml index e2e922a82..c5f1e0167 100644 --- a/crates/core/Cargo.toml +++ b/crates/core/Cargo.toml @@ -53,7 +53,7 @@ datafusion = { workspace = true, features = ["avro", "unicode_expressions"] } datafusion-substrait = { workspace = true, optional = true } datafusion-proto = { workspace = true } datafusion-ffi = { workspace = true } -datafusion-spark = { workspace = true } +datafusion-spark = { workspace = true, features = ["core"] } prost = { workspace = true } # keep in line with `datafusion-substrait` serde_json = { workspace = true } uuid = { workspace = true, features = ["v4"] } diff --git a/crates/core/src/array.rs b/crates/core/src/array.rs index f284fa9de..dfe963183 100644 --- a/crates/core/src/array.rs +++ b/crates/core/src/array.rs @@ -63,10 +63,10 @@ impl PyArrowArrayExportable { }; let ffi_schema = FFI_ArrowSchema::try_from(&field)?; - let schema_capsule = PyCapsule::new(py, ffi_schema, Some(cr"arrow_schema".into()))?; + let schema_capsule = PyCapsule::new_with_value(py, ffi_schema, cr"arrow_schema")?; let ffi_array = FFI_ArrowArray::new(&self.array.to_data()); - let array_capsule = PyCapsule::new(py, ffi_array, Some(cr"arrow_array".into()))?; + let array_capsule = PyCapsule::new_with_value(py, ffi_array, cr"arrow_array")?; Ok((schema_capsule, array_capsule)) } diff --git a/crates/core/src/codec.rs b/crates/core/src/codec.rs index a6ea6671c..26853e69f 100644 --- a/crates/core/src/codec.rs +++ b/crates/core/src/codec.rs @@ -100,9 +100,13 @@ use datafusion::logical_expr::{ TypeSignature, Volatility, WindowUDF, WindowUDFImpl, }; use datafusion::physical_expr::PhysicalExpr; +use datafusion::physical_expr_common::physical_expr::proto_decode::PhysicalExprDecodeCtx; +use datafusion::physical_expr_common::physical_expr::proto_encode::PhysicalExprEncodeCtx; use datafusion::physical_plan::ExecutionPlan; use datafusion_proto::logical_plan::{DefaultLogicalExtensionCodec, LogicalExtensionCodec}; -use datafusion_proto::physical_plan::{DefaultPhysicalExtensionCodec, PhysicalExtensionCodec}; +use datafusion_proto::physical_plan::{ + DefaultPhysicalExtensionCodec, PhysicalExtensionCodec, PhysicalProtoConverterExtension, +}; use pyo3::prelude::*; use pyo3::sync::PyOnceLock; use pyo3::types::{PyBytes, PyTuple}; @@ -483,12 +487,18 @@ impl PhysicalExtensionCodec for PythonPhysicalCodec { buf: &[u8], inputs: &[Arc], ctx: &TaskContext, + proto_converter: &dyn PhysicalProtoConverterExtension, ) -> Result> { - self.inner.try_decode(buf, inputs, ctx) + self.inner.try_decode(buf, inputs, ctx, proto_converter) } - fn try_encode(&self, node: Arc, buf: &mut Vec) -> Result<()> { - self.inner.try_encode(node, buf) + fn try_encode( + &self, + node: Arc, + buf: &mut Vec, + proto_converter: &dyn PhysicalProtoConverterExtension, + ) -> Result<()> { + self.inner.try_encode(node, buf, proto_converter) } fn try_encode_udf(&self, node: &ScalarUDF, buf: &mut Vec) -> Result<()> { @@ -509,16 +519,22 @@ impl PhysicalExtensionCodec for PythonPhysicalCodec { self.inner.try_decode_udf(name, buf) } - fn try_encode_expr(&self, node: &Arc, buf: &mut Vec) -> Result<()> { - self.inner.try_encode_expr(node, buf) + fn try_encode_expr( + &self, + node: &Arc, + buf: &mut Vec, + ctx: &PhysicalExprEncodeCtx<'_>, + ) -> Result<()> { + self.inner.try_encode_expr(node, buf, ctx) } fn try_decode_expr( &self, buf: &[u8], inputs: &[Arc], + ctx: &PhysicalExprDecodeCtx<'_>, ) -> Result> { - self.inner.try_decode_expr(buf, inputs) + self.inner.try_decode_expr(buf, inputs, ctx) } fn try_encode_udaf(&self, node: &AggregateUDF, buf: &mut Vec) -> Result<()> { diff --git a/crates/core/src/context.rs b/crates/core/src/context.rs index 0613a96dc..7bbeed2f1 100644 --- a/crates/core/src/context.rs +++ b/crates/core/src/context.rs @@ -832,7 +832,19 @@ impl PySessionContext { name: &str, partitions: PyArrowType>>, ) -> PyDataFusionResult<()> { - let schema = partitions.0[0][0].schema(); + // Take the schema from the first available batch; error instead of + // panicking when the partitions hold no batches. + let schema = partitions + .0 + .iter() + .find_map(|partition| partition.first()) + .ok_or_else(|| { + PyValueError::new_err( + "Cannot register record batches without a schema: the \ + provided partitions contain no record batches.", + ) + })? + .schema(); let table = MemTable::try_new(schema, partitions.0)?; self.ctx.register_table(name, Arc::new(table))?; Ok(()) @@ -1359,12 +1371,10 @@ impl PySessionContext { &self, py: Python<'py>, ) -> PyResult> { - let name = cr"datafusion_task_context_provider".into(); - let ctx_provider = Arc::clone(&self.ctx) as Arc; let ffi_ctx_provider = FFI_TaskContextProvider::from(&ctx_provider); - PyCapsule::new(py, ffi_ctx_provider, Some(name)) + PyCapsule::new_with_value(py, ffi_ctx_provider, cr"datafusion_task_context_provider") } pub fn __datafusion_logical_extension_codec__<'py>( diff --git a/crates/core/src/dataframe.rs b/crates/core/src/dataframe.rs index 8f1a20d0d..b1f305551 100644 --- a/crates/core/src/dataframe.rs +++ b/crates/core/src/dataframe.rs @@ -16,7 +16,7 @@ // under the License. use std::collections::HashMap; -use std::ffi::{CStr, CString}; +use std::ffi::CStr; use std::ptr::NonNull; use std::str::FromStr; use std::sync::Arc; @@ -1237,8 +1237,7 @@ impl PyDataFrame { // destructor provided by PyO3 will drop the stream unless ownership is // transferred to PyArrow during import. let stream = FFI_ArrowArrayStream::new(reader); - let name = CString::new(ARROW_ARRAY_STREAM_NAME.to_bytes()).unwrap(); - let capsule = PyCapsule::new(py, stream, Some(name))?; + let capsule = PyCapsule::new_with_value(py, stream, ARROW_ARRAY_STREAM_NAME)?; Ok(capsule) } @@ -1317,7 +1316,8 @@ impl PyDataFrame { None => Vec::new(), // Empty vector means fill null for all columns }; - let df = self.df.as_ref().clone().fill_null(scalar_value.0, cols)?; + let cols = cols.iter().map(String::as_str).collect::>(); + let df = self.df.as_ref().fill_null(&scalar_value.0, &cols)?; Ok(Self::new(df)) } } diff --git a/crates/core/src/dataset_exec.rs b/crates/core/src/dataset_exec.rs index 32c030b00..963603c8b 100644 --- a/crates/core/src/dataset_exec.rs +++ b/crates/core/src/dataset_exec.rs @@ -21,11 +21,12 @@ use datafusion::arrow::datatypes::SchemaRef; use datafusion::arrow::error::{ArrowError, Result as ArrowResult}; use datafusion::arrow::pyarrow::PyArrowType; use datafusion::arrow::record_batch::RecordBatch; +use datafusion::common::tree_node::TreeNodeRecursion; use datafusion::error::{DataFusionError as InnerDataFusionError, Result as DFResult}; use datafusion::execution::context::TaskContext; use datafusion::logical_expr::Expr; use datafusion::logical_expr::utils::conjunction; -use datafusion::physical_expr::{EquivalenceProperties, LexOrdering}; +use datafusion::physical_expr::{EquivalenceProperties, LexOrdering, PhysicalExpr}; use datafusion::physical_plan::execution_plan::{Boundedness, EmissionType}; use datafusion::physical_plan::stream::RecordBatchStreamAdapter; use datafusion::physical_plan::{ @@ -165,6 +166,15 @@ impl ExecutionPlan for DatasetExec { vec![] } + fn apply_expressions( + &self, + _f: &mut dyn FnMut(&Arc) -> DFResult, + ) -> DFResult { + // Any filters are pushed down into the PyArrow dataset scanner as pyarrow + // expressions, so this node owns no physical expressions to visit. + Ok(TreeNodeRecursion::Continue) + } + fn with_new_children( self: Arc, _: Vec>, diff --git a/crates/core/src/expr.rs b/crates/core/src/expr.rs index 432c4cd23..cab997d7a 100644 --- a/crates/core/src/expr.rs +++ b/crates/core/src/expr.rs @@ -433,7 +433,7 @@ impl PyExpr { Expr::Literal(scalar_value, _) => scalar_to_pyarrow(scalar_value, py), _ => Err(py_type_err(format!( "Non Expr::Literal encountered in types: {:?}", - &self.expr + self.expr ))), } } @@ -612,7 +612,7 @@ impl PyExpr { _ => { return Err(py_type_err(format!( "Catch all triggered in get_operator_name: {:?}", - &self.expr + self.expr ))); } }) diff --git a/crates/core/src/expr/aggregate.rs b/crates/core/src/expr/aggregate.rs index 5a6a771a7..7177fb469 100644 --- a/crates/core/src/expr/aggregate.rs +++ b/crates/core/src/expr/aggregate.rs @@ -65,8 +65,8 @@ impl Display for PyAggregate { \nAggregates(s): {:?} \nInput: {:?} \nProjected Schema: {:?}", - &self.aggregate.group_expr, - &self.aggregate.aggr_expr, + self.aggregate.group_expr, + self.aggregate.aggr_expr, self.aggregate.input, self.aggregate.schema ) diff --git a/crates/core/src/expr/alias.rs b/crates/core/src/expr/alias.rs index b76e82e22..391c94dbc 100644 --- a/crates/core/src/expr/alias.rs +++ b/crates/core/src/expr/alias.rs @@ -53,7 +53,7 @@ impl Display for PyAlias { "Alias \nExpr: `{:?}` \nAlias Name: `{}`", - &self.alias.expr, &self.alias.name + self.alias.expr, self.alias.name ) } } diff --git a/crates/core/src/expr/between.rs b/crates/core/src/expr/between.rs index 6943b6c3b..80f6c70da 100644 --- a/crates/core/src/expr/between.rs +++ b/crates/core/src/expr/between.rs @@ -55,7 +55,7 @@ impl Display for PyBetween { Negated: {:?} Low: {:?} High: {:?}", - &self.between.expr, &self.between.negated, &self.between.low, &self.between.high + self.between.expr, self.between.negated, self.between.low, self.between.high ) } } diff --git a/crates/core/src/expr/bool_expr.rs b/crates/core/src/expr/bool_expr.rs index 9e374c7e2..d1cd7bcf7 100644 --- a/crates/core/src/expr/bool_expr.rs +++ b/crates/core/src/expr/bool_expr.rs @@ -46,7 +46,7 @@ impl Display for PyNot { f, "Not Expr: {}", - &self.expr + self.expr ) } } @@ -82,7 +82,7 @@ impl Display for PyIsNotNull { f, "IsNotNull Expr: {}", - &self.expr + self.expr ) } } @@ -118,7 +118,7 @@ impl Display for PyIsNull { f, "IsNull Expr: {}", - &self.expr + self.expr ) } } @@ -154,7 +154,7 @@ impl Display for PyIsTrue { f, "IsTrue Expr: {}", - &self.expr + self.expr ) } } @@ -190,7 +190,7 @@ impl Display for PyIsFalse { f, "IsFalse Expr: {}", - &self.expr + self.expr ) } } @@ -226,7 +226,7 @@ impl Display for PyIsUnknown { f, "IsUnknown Expr: {}", - &self.expr + self.expr ) } } @@ -262,7 +262,7 @@ impl Display for PyIsNotTrue { f, "IsNotTrue Expr: {}", - &self.expr + self.expr ) } } @@ -298,7 +298,7 @@ impl Display for PyIsNotFalse { f, "IsNotFalse Expr: {}", - &self.expr + self.expr ) } } @@ -334,7 +334,7 @@ impl Display for PyIsNotUnknown { f, "IsNotUnknown Expr: {}", - &self.expr + self.expr ) } } @@ -370,7 +370,7 @@ impl Display for PyNegative { f, "Negative Expr: {}", - &self.expr + self.expr ) } } diff --git a/crates/core/src/expr/create_external_table.rs b/crates/core/src/expr/create_external_table.rs index 980eea131..e78b836bf 100644 --- a/crates/core/src/expr/create_external_table.rs +++ b/crates/core/src/expr/create_external_table.rs @@ -88,7 +88,7 @@ impl PyCreateExternalTable { let create = CreateExternalTable { schema: Arc::new(schema.into()), name: name.into(), - location, + locations: vec![location], file_type, table_partition_cols, if_not_exists, @@ -118,8 +118,8 @@ impl PyCreateExternalTable { Ok(self.create.name.to_string()) } - pub fn location(&self) -> String { - self.create.location.clone() + pub fn locations(&self) -> Vec { + self.create.locations.clone() } pub fn file_type(&self) -> String { diff --git a/crates/core/src/expr/create_memory_table.rs b/crates/core/src/expr/create_memory_table.rs index 3214dab0e..c27c11834 100644 --- a/crates/core/src/expr/create_memory_table.rs +++ b/crates/core/src/expr/create_memory_table.rs @@ -57,10 +57,7 @@ impl Display for PyCreateMemoryTable { Input: {:?} if_not_exists: {:?} or_replace: {:?}", - &self.create.name, - &self.create.input, - &self.create.if_not_exists, - &self.create.or_replace, + self.create.name, self.create.input, self.create.if_not_exists, self.create.or_replace, ) } } diff --git a/crates/core/src/expr/create_view.rs b/crates/core/src/expr/create_view.rs index 6941ef769..2f0315c40 100644 --- a/crates/core/src/expr/create_view.rs +++ b/crates/core/src/expr/create_view.rs @@ -58,7 +58,7 @@ impl Display for PyCreateView { input: {:?} or_replace: {:?} definition: {:?}", - &self.create.name, &self.create.input, &self.create.or_replace, &self.create.definition, + self.create.name, self.create.input, self.create.or_replace, self.create.definition, ) } } diff --git a/crates/core/src/expr/dml.rs b/crates/core/src/expr/dml.rs index 26f975820..5967d181e 100644 --- a/crates/core/src/expr/dml.rs +++ b/crates/core/src/expr/dml.rs @@ -18,6 +18,7 @@ use datafusion::logical_expr::dml::InsertOp; use datafusion::logical_expr::{DmlStatement, WriteOp}; use pyo3::IntoPyObjectExt; +use pyo3::exceptions::PyNotImplementedError; use pyo3::prelude::*; use super::logical_node::LogicalNode; @@ -71,8 +72,8 @@ impl PyDmlStatement { }) } - pub fn op(&self) -> PyWriteOp { - self.dml.op.clone().into() + pub fn op(&self) -> PyResult { + self.dml.op.clone().try_into() } pub fn input(&self) -> PyLogicalPlan { @@ -112,16 +113,21 @@ pub enum PyWriteOp { Truncate, } -impl From for PyWriteOp { - fn from(write_op: WriteOp) -> Self { +impl TryFrom for PyWriteOp { + type Error = PyErr; + + fn try_from(write_op: WriteOp) -> Result { match write_op { - WriteOp::Insert(InsertOp::Append) => PyWriteOp::Append, - WriteOp::Insert(InsertOp::Overwrite) => PyWriteOp::Overwrite, - WriteOp::Insert(InsertOp::Replace) => PyWriteOp::Replace, - WriteOp::Update => PyWriteOp::Update, - WriteOp::Delete => PyWriteOp::Delete, - WriteOp::Ctas => PyWriteOp::Ctas, - WriteOp::Truncate => PyWriteOp::Truncate, + WriteOp::Insert(InsertOp::Append) => Ok(PyWriteOp::Append), + WriteOp::Insert(InsertOp::Overwrite) => Ok(PyWriteOp::Overwrite), + WriteOp::Insert(InsertOp::Replace) => Ok(PyWriteOp::Replace), + WriteOp::Update => Ok(PyWriteOp::Update), + WriteOp::Delete => Ok(PyWriteOp::Delete), + WriteOp::Ctas => Ok(PyWriteOp::Ctas), + WriteOp::Truncate => Ok(PyWriteOp::Truncate), + unsupported => Err(PyNotImplementedError::new_err(format!( + "DataFusion write operation {unsupported:?} is not supported" + ))), } } } diff --git a/crates/core/src/expr/drop_catalog_schema.rs b/crates/core/src/expr/drop_catalog_schema.rs index fd5105332..f349098b7 100644 --- a/crates/core/src/expr/drop_catalog_schema.rs +++ b/crates/core/src/expr/drop_catalog_schema.rs @@ -18,9 +18,8 @@ use std::fmt::{self, Display, Formatter}; use std::sync::Arc; -use datafusion::common::SchemaReference; +use datafusion::common::{SchemaReference, TableReference}; use datafusion::logical_expr::DropCatalogSchema; -use datafusion::sql::TableReference; use pyo3::IntoPyObjectExt; use pyo3::exceptions::PyValueError; use pyo3::prelude::*; diff --git a/crates/core/src/expr/drop_table.rs b/crates/core/src/expr/drop_table.rs index 46fe67465..156d3bcc2 100644 --- a/crates/core/src/expr/drop_table.rs +++ b/crates/core/src/expr/drop_table.rs @@ -56,7 +56,7 @@ impl Display for PyDropTable { name: {:?} if_exists: {:?} schema: {:?}", - &self.drop.name, &self.drop.if_exists, &self.drop.schema, + self.drop.name, self.drop.if_exists, self.drop.schema, ) } } diff --git a/crates/core/src/expr/empty_relation.rs b/crates/core/src/expr/empty_relation.rs index f3c237731..5af4aafeb 100644 --- a/crates/core/src/expr/empty_relation.rs +++ b/crates/core/src/expr/empty_relation.rs @@ -56,7 +56,7 @@ impl Display for PyEmptyRelation { "Empty Relation Produce One Row: {:?} Schema: {:?}", - &self.empty.produce_one_row, &self.empty.schema + self.empty.produce_one_row, self.empty.schema ) } } diff --git a/crates/core/src/expr/explain.rs b/crates/core/src/expr/explain.rs index 6259951de..d6ba1c25c 100644 --- a/crates/core/src/expr/explain.rs +++ b/crates/core/src/expr/explain.rs @@ -61,11 +61,11 @@ impl Display for PyExplain { stringified_plans: {:?} schema: {:?} logical_optimization_succeeded: {:?}", - &self.explain.verbose, - &self.explain.plan, - &self.explain.stringified_plans, - &self.explain.schema, - &self.explain.logical_optimization_succeeded + self.explain.verbose, + self.explain.plan, + self.explain.stringified_plans, + self.explain.schema, + self.explain.logical_optimization_succeeded ) } } diff --git a/crates/core/src/expr/filter.rs b/crates/core/src/expr/filter.rs index 67426806d..1fe5f2c7f 100644 --- a/crates/core/src/expr/filter.rs +++ b/crates/core/src/expr/filter.rs @@ -57,7 +57,7 @@ impl Display for PyFilter { "Filter Predicate: {:?} Input: {:?}", - &self.filter.predicate, &self.filter.input + self.filter.predicate, self.filter.input ) } } diff --git a/crates/core/src/expr/higher_order_function.rs b/crates/core/src/expr/higher_order_function.rs index 91a94de2c..5ba64052a 100644 --- a/crates/core/src/expr/higher_order_function.rs +++ b/crates/core/src/expr/higher_order_function.rs @@ -52,7 +52,7 @@ impl Display for PyHigherOrderFunction { f, "HigherOrderFunction(name={}, args={:?})", self.higher_order.name(), - &self.higher_order.args, + self.higher_order.args, ) } } diff --git a/crates/core/src/expr/join.rs b/crates/core/src/expr/join.rs index b90f2f57d..634e13734 100644 --- a/crates/core/src/expr/join.rs +++ b/crates/core/src/expr/join.rs @@ -132,14 +132,14 @@ impl Display for PyJoin { JoinConstraint: {:?} Schema: {:?} NullEquality: {:?}", - &self.join.left, - &self.join.right, - &self.join.on, - &self.join.filter, - &self.join.join_type, - &self.join.join_constraint, - &self.join.schema, - &self.join.null_equality, + self.join.left, + self.join.right, + self.join.on, + self.join.filter, + self.join.join_type, + self.join.join_constraint, + self.join.schema, + self.join.null_equality, ) } } diff --git a/crates/core/src/expr/lambda.rs b/crates/core/src/expr/lambda.rs index 3ebc6e61c..7190521d6 100644 --- a/crates/core/src/expr/lambda.rs +++ b/crates/core/src/expr/lambda.rs @@ -51,7 +51,7 @@ impl Display for PyLambda { write!( f, "Lambda(params={:?}, body={:?})", - &self.lambda.params, &self.lambda.body, + self.lambda.params, self.lambda.body, ) } } diff --git a/crates/core/src/expr/lambda_variable.rs b/crates/core/src/expr/lambda_variable.rs index 2ef554e17..7baf5f21c 100644 --- a/crates/core/src/expr/lambda_variable.rs +++ b/crates/core/src/expr/lambda_variable.rs @@ -46,7 +46,7 @@ impl From for LambdaVariable { impl Display for PyLambdaVariable { fn fmt(&self, f: &mut Formatter) -> fmt::Result { - write!(f, "LambdaVariable({})", &self.variable.name) + write!(f, "LambdaVariable({})", self.variable.name) } } diff --git a/crates/core/src/expr/like.rs b/crates/core/src/expr/like.rs index 417dc9182..900551a20 100644 --- a/crates/core/src/expr/like.rs +++ b/crates/core/src/expr/like.rs @@ -55,10 +55,10 @@ impl Display for PyLike { Expr: {:?} Pattern: {:?} Escape_Char: {:?}", - &self.negated(), - &self.expr(), - &self.pattern(), - &self.escape_char() + self.negated(), + self.expr(), + self.pattern(), + self.escape_char() ) } } @@ -119,10 +119,10 @@ impl Display for PyILike { Expr: {:?} Pattern: {:?} Escape_Char: {:?}", - &self.negated(), - &self.expr(), - &self.pattern(), - &self.escape_char() + self.negated(), + self.expr(), + self.pattern(), + self.escape_char() ) } } @@ -183,10 +183,10 @@ impl Display for PySimilarTo { Expr: {:?} Pattern: {:?} Escape_Char: {:?}", - &self.negated(), - &self.expr(), - &self.pattern(), - &self.escape_char() + self.negated(), + self.expr(), + self.pattern(), + self.escape_char() ) } } diff --git a/crates/core/src/expr/limit.rs b/crates/core/src/expr/limit.rs index c04b8bfa8..f53737b5e 100644 --- a/crates/core/src/expr/limit.rs +++ b/crates/core/src/expr/limit.rs @@ -22,6 +22,7 @@ use pyo3::IntoPyObjectExt; use pyo3::prelude::*; use crate::common::df_schema::PyDFSchema; +use crate::expr::PyExpr; use crate::expr::logical_node::LogicalNode; use crate::sql::logical::PyLogicalPlan; @@ -57,26 +58,30 @@ impl Display for PyLimit { Skip: {:?} Fetch: {:?} Input: {:?}", - &self.limit.skip, &self.limit.fetch, &self.limit.input + self.limit.skip, self.limit.fetch, self.limit.input ) } } #[pymethods] impl PyLimit { - // NOTE: Upstream now has expressions for skip and fetch - // TODO: Do we still want to expose these? - // REF: https://github.com/apache/datafusion/pull/12836 - - // /// Retrieves the skip value for this `Limit` - // fn skip(&self) -> usize { - // self.limit.skip - // } + // Retrieves the skip expression for this `Limit`, if any. + // + // `LIMIT`/`OFFSET` were changed upstream to support arbitrary + // expressions (not just constants), see + // https://github.com/apache/datafusion/pull/13028. Callers that expect + // a simple literal (the common case, e.g. `OFFSET 5`) should evaluate + // the returned `PyExpr` via `Expr.python_value()`. + fn skip(&self) -> PyResult> { + Ok(self.limit.skip.as_deref().cloned().map(PyExpr::from)) + } - // /// Retrieves the fetch value for this `Limit` - // fn fetch(&self) -> Option { - // self.limit.fetch - // } + // Retrieves the fetch expression for this `Limit`, if any. + // + // See the note on `skip` above regarding expression-based limits. + fn fetch(&self) -> PyResult> { + Ok(self.limit.fetch.as_deref().cloned().map(PyExpr::from)) + } /// Retrieves the input `LogicalPlan` to this `Limit` node fn input(&self) -> PyResult> { diff --git a/crates/core/src/expr/projection.rs b/crates/core/src/expr/projection.rs index 456e06412..7e22e1e7b 100644 --- a/crates/core/src/expr/projection.rs +++ b/crates/core/src/expr/projection.rs @@ -65,7 +65,7 @@ impl Display for PyProjection { \nExpr(s): {:?} \nInput: {:?} \nProjected Schema: {:?}", - &self.projection.expr, &self.projection.input, &self.projection.schema, + self.projection.expr, self.projection.input, self.projection.schema, ) } } diff --git a/crates/core/src/expr/recursive_query.rs b/crates/core/src/expr/recursive_query.rs index e03137b80..0b198a191 100644 --- a/crates/core/src/expr/recursive_query.rs +++ b/crates/core/src/expr/recursive_query.rs @@ -22,6 +22,7 @@ use pyo3::IntoPyObjectExt; use pyo3::prelude::*; use super::logical_node::LogicalNode; +use crate::errors::PyDataFusionResult; use crate::sql::logical::PyLogicalPlan; #[pyclass( @@ -67,15 +68,15 @@ impl PyRecursiveQuery { static_term: PyLogicalPlan, recursive_term: PyLogicalPlan, is_distinct: bool, - ) -> Self { - Self { - query: RecursiveQuery { + ) -> PyDataFusionResult { + Ok(Self { + query: RecursiveQuery::try_new( name, - static_term: static_term.plan(), - recursive_term: recursive_term.plan(), + static_term.plan(), + recursive_term.plan(), is_distinct, - }, - } + )?, + }) } fn name(&self) -> PyResult { diff --git a/crates/core/src/expr/repartition.rs b/crates/core/src/expr/repartition.rs index be39b9978..cbc8a97bb 100644 --- a/crates/core/src/expr/repartition.rs +++ b/crates/core/src/expr/repartition.rs @@ -82,7 +82,7 @@ impl Display for PyRepartition { "Repartition input: {:?} partitioning_scheme: {:?}", - &self.repartition.input, &self.repartition.partitioning_scheme, + self.repartition.input, self.repartition.partitioning_scheme, ) } } diff --git a/crates/core/src/expr/sort.rs b/crates/core/src/expr/sort.rs index 7c1e654c5..1b1065011 100644 --- a/crates/core/src/expr/sort.rs +++ b/crates/core/src/expr/sort.rs @@ -61,7 +61,7 @@ impl Display for PySort { \nExpr(s): {:?} \nInput: {:?} \nSchema: {:?}", - &self.sort.expr, + self.sort.expr, self.sort.input, self.sort.input.schema() ) diff --git a/crates/core/src/expr/sort_expr.rs b/crates/core/src/expr/sort_expr.rs index 3c3c86bc1..93faffcec 100644 --- a/crates/core/src/expr/sort_expr.rs +++ b/crates/core/src/expr/sort_expr.rs @@ -54,7 +54,7 @@ impl Display for PySortExpr { Expr: {:?} Asc: {:?} NullsFirst: {:?}", - &self.sort.expr, &self.sort.asc, &self.sort.nulls_first + self.sort.expr, self.sort.asc, self.sort.nulls_first ) } } diff --git a/crates/core/src/expr/table_scan.rs b/crates/core/src/expr/table_scan.rs index 8ba7e4a69..46ce551d2 100644 --- a/crates/core/src/expr/table_scan.rs +++ b/crates/core/src/expr/table_scan.rs @@ -65,10 +65,10 @@ impl Display for PyTableScan { Projections: {:?} Projected Schema: {:?} Filters: {:?}", - &self.table_scan.table_name, - &self.py_projections(), - &self.py_schema(), - &self.py_filters(), + self.table_scan.table_name, + self.py_projections(), + self.py_schema(), + self.py_filters(), ) } } diff --git a/crates/core/src/expr/union.rs b/crates/core/src/expr/union.rs index a3b9efe91..bd5770e0a 100644 --- a/crates/core/src/expr/union.rs +++ b/crates/core/src/expr/union.rs @@ -56,7 +56,7 @@ impl Display for PyUnion { "Union Inputs: {:?} Schema: {:?}", - &self.union_.inputs, &self.union_.schema, + self.union_.inputs, self.union_.schema, ) } } diff --git a/crates/core/src/expr/unnest.rs b/crates/core/src/expr/unnest.rs index 880d0a279..540667824 100644 --- a/crates/core/src/expr/unnest.rs +++ b/crates/core/src/expr/unnest.rs @@ -56,7 +56,7 @@ impl Display for PyUnnest { "Unnest Inputs: {:?} Schema: {:?}", - &self.unnest_.input, &self.unnest_.schema, + self.unnest_.input, self.unnest_.schema, ) } } diff --git a/crates/core/src/expr/unnest_expr.rs b/crates/core/src/expr/unnest_expr.rs index 97feef1d1..549257b86 100644 --- a/crates/core/src/expr/unnest_expr.rs +++ b/crates/core/src/expr/unnest_expr.rs @@ -52,7 +52,7 @@ impl Display for PyUnnestExpr { f, "Unnest Expr: {:?}", - &self.unnest.expr, + self.unnest.expr, ) } } diff --git a/crates/core/src/expr/window.rs b/crates/core/src/expr/window.rs index 92d909bfc..e0050f671 100644 --- a/crates/core/src/expr/window.rs +++ b/crates/core/src/expr/window.rs @@ -105,7 +105,7 @@ impl Display for PyWindowExpr { "Over\n Window Expr: {:?} Schema: {:?}", - &self.window.window_expr, &self.window.schema + self.window.window_expr, self.window.schema ) } } diff --git a/crates/core/src/sql/logical.rs b/crates/core/src/sql/logical.rs index 647c3fa7e..12fc43bdc 100644 --- a/crates/core/src/sql/logical.rs +++ b/crates/core/src/sql/logical.rs @@ -135,7 +135,7 @@ impl PyLogicalPlan { LogicalPlan::Dml(plan) => PyDmlStatement::from(plan.clone()).to_variant(py), LogicalPlan::Ddl(plan) => match plan { DdlStatement::CreateExternalTable(plan) => { - PyCreateExternalTable::from(plan.clone()).to_variant(py) + PyCreateExternalTable::from(plan.as_ref().clone()).to_variant(py) } DdlStatement::CreateMemoryTable(plan) => { PyCreateMemoryTable::from(plan.clone()).to_variant(py) @@ -154,7 +154,7 @@ impl PyLogicalPlan { PyDropCatalogSchema::from(plan.clone()).to_variant(py) } DdlStatement::CreateFunction(plan) => { - PyCreateFunction::from(plan.clone()).to_variant(py) + PyCreateFunction::from(plan.as_ref().clone()).to_variant(py) } DdlStatement::DropFunction(plan) => { PyDropFunction::from(plan.clone()).to_variant(py) diff --git a/crates/core/src/substrait.rs b/crates/core/src/substrait.rs index 27e446f48..dcb5587cf 100644 --- a/crates/core/src/substrait.rs +++ b/crates/core/src/substrait.rs @@ -109,7 +109,7 @@ impl PySubstraitSerializer { ) -> PyDataFusionResult { PySubstraitSerializer::serialize_bytes(sql, ctx, py).and_then(|proto_bytes| { let proto_bytes = proto_bytes.bind(py).cast::().unwrap(); - PySubstraitSerializer::deserialize_bytes(proto_bytes.as_bytes().to_vec(), py) + PySubstraitSerializer::deserialize_bytes(proto_bytes.as_bytes().to_vec()) }) } @@ -131,8 +131,8 @@ impl PySubstraitSerializer { } #[staticmethod] - pub fn deserialize_bytes(proto_bytes: Vec, py: Python) -> PyDataFusionResult { - let plan = wait_for_future(py, serializer::deserialize_bytes(proto_bytes))??; + pub fn deserialize_bytes(proto_bytes: Vec) -> PyDataFusionResult { + let plan = serializer::deserialize_bytes(&proto_bytes)?; Ok(PyPlan { plan: *plan }) } } diff --git a/crates/util/src/lib.rs b/crates/util/src/lib.rs index 07aa0a2d5..9327d7f2f 100644 --- a/crates/util/src/lib.rs +++ b/crates/util/src/lib.rs @@ -164,7 +164,8 @@ pub fn validate_pycapsule(capsule: &Bound, name: &str) -> PyResult<() ))); } - let capsule_name = unsafe { capsule_name.unwrap().as_cstr().to_str()? }; + let capsule_name = unsafe { capsule_name.unwrap().as_cstr().to_str() } + .map_err(|err| PyValueError::new_err(err.to_string()))?; if capsule_name != name { return Err(PyValueError::new_err(format!( "Expected name '{name}' in PyCapsule, instead got '{capsule_name}'" @@ -208,10 +209,9 @@ pub fn create_logical_extension_capsule<'py>( py: Python<'py>, codec: &FFI_LogicalExtensionCodec, ) -> PyResult> { - let name = cr"datafusion_logical_extension_codec".into(); let codec = codec.clone(); - PyCapsule::new(py, codec, Some(name)) + PyCapsule::new_with_value(py, codec, cr"datafusion_logical_extension_codec") } pub fn ffi_logical_codec_from_pycapsule(obj: Bound) -> PyResult { @@ -235,10 +235,9 @@ pub fn create_physical_extension_capsule<'py>( py: Python<'py>, codec: &FFI_PhysicalExtensionCodec, ) -> PyResult> { - let name = cr"datafusion_physical_extension_codec".into(); let codec = codec.clone(); - PyCapsule::new(py, codec, Some(name)) + PyCapsule::new_with_value(py, codec, cr"datafusion_physical_extension_codec") } /// Define a `(obj) -> PyResult>` extractor that diff --git a/docs/source/contributor-guide/introduction.md b/docs/source/contributor-guide/introduction.md index 683b27bf2..24d56ed90 100644 --- a/docs/source/contributor-guide/introduction.md +++ b/docs/source/contributor-guide/introduction.md @@ -47,6 +47,8 @@ Bootstrap: ```shell # fetch this repo git clone git@github.com:apache/datafusion-python.git +# cd to the repo root +cd datafusion-python/ # create the virtual environment uv sync --dev --no-install-package datafusion # activate the environment @@ -64,7 +66,7 @@ Whenever rust code changes (your changes or via `git pull`): ```shell # make sure you activate the venv using "source .venv/bin/activate" first -maturin develop -uv +maturin develop --uv python -m pytest ``` diff --git a/examples/datafusion-ffi-example/python/tests/_test_table_provider_factory.py b/examples/datafusion-ffi-example/python/tests/_test_table_provider_factory.py index b1e94ec73..0c0f1a575 100644 --- a/examples/datafusion-ffi-example/python/tests/_test_table_provider_factory.py +++ b/examples/datafusion-ffi-example/python/tests/_test_table_provider_factory.py @@ -32,7 +32,7 @@ def test_table_provider_factory_ffi() -> None: CREATE EXTERNAL TABLE foo STORED AS my_format - LOCATION ''; + LOCATION 'unused'; """).collect() # Query the pre-populated table diff --git a/examples/datafusion-ffi-example/src/aggregate_udf.rs b/examples/datafusion-ffi-example/src/aggregate_udf.rs index 86737778f..ea1518365 100644 --- a/examples/datafusion-ffi-example/src/aggregate_udf.rs +++ b/examples/datafusion-ffi-example/src/aggregate_udf.rs @@ -50,12 +50,10 @@ impl MySumUDF { &self, py: Python<'py>, ) -> PyResult> { - let name = cr"datafusion_aggregate_udf".into(); - let func = Arc::new(AggregateUDF::from(self.clone())); let provider = FFI_AggregateUDF::from(func); - PyCapsule::new(py, provider, Some(name)) + PyCapsule::new_with_value(py, provider, cr"datafusion_aggregate_udf") } } diff --git a/examples/datafusion-ffi-example/src/catalog_provider.rs b/examples/datafusion-ffi-example/src/catalog_provider.rs index 6131ab0f0..a56b5855c 100644 --- a/examples/datafusion-ffi-example/src/catalog_provider.rs +++ b/examples/datafusion-ffi-example/src/catalog_provider.rs @@ -33,8 +33,8 @@ use pyo3::types::PyCapsule; use pyo3::{Bound, PyAny, PyResult, Python, pyclass, pymethods}; pub fn my_table() -> Arc { + use arrow::array::record_batch; use arrow::datatypes::{DataType, Field}; - use datafusion_common::record_batch; let schema = Arc::new(Schema::new(vec![ Field::new("units", DataType::Int32, true), @@ -92,14 +92,12 @@ impl FixedSchemaProvider { py: Python<'py>, session: Bound, ) -> PyResult> { - let name = cr"datafusion_schema_provider".into(); - let provider = Arc::clone(&self.inner) as Arc; let codec = ffi_logical_codec_from_pycapsule(session)?; let provider = FFI_SchemaProvider::new_with_ffi_codec(provider, None, codec); - PyCapsule::new(py, provider, Some(name)) + PyCapsule::new_with_value(py, provider, cr"datafusion_schema_provider") } } @@ -186,14 +184,12 @@ impl MyCatalogProvider { py: Python<'py>, session: Bound, ) -> PyResult> { - let name = cr"datafusion_catalog_provider".into(); - let provider = Arc::clone(&self.inner) as Arc; let codec = ffi_logical_codec_from_pycapsule(session)?; let provider = FFI_CatalogProvider::new_with_ffi_codec(provider, None, codec); - PyCapsule::new(py, provider, Some(name)) + PyCapsule::new_with_value(py, provider, cr"datafusion_catalog_provider") } } @@ -247,13 +243,11 @@ impl MyCatalogProviderList { py: Python<'py>, session: Bound, ) -> PyResult> { - let name = cr"datafusion_catalog_provider_list".into(); - let provider = Arc::clone(&self.inner) as Arc; let codec = ffi_logical_codec_from_pycapsule(session)?; let provider = FFI_CatalogProviderList::new_with_ffi_codec(provider, None, codec); - PyCapsule::new(py, provider, Some(name)) + PyCapsule::new_with_value(py, provider, cr"datafusion_catalog_provider_list") } } diff --git a/examples/datafusion-ffi-example/src/config.rs b/examples/datafusion-ffi-example/src/config.rs index 6cdb8aa83..9cf704b94 100644 --- a/examples/datafusion-ffi-example/src/config.rs +++ b/examples/datafusion-ffi-example/src/config.rs @@ -53,14 +53,12 @@ impl MyConfig { &self, py: Python<'py>, ) -> PyResult> { - let name = cr"datafusion_extension_options".into(); - let mut config = FFI_ExtensionOptions::default(); config .add_config(self) .map_err(|e| PyRuntimeError::new_err(e.to_string()))?; - PyCapsule::new(py, config, Some(name)) + PyCapsule::new_with_value(py, config, cr"datafusion_extension_options") } } diff --git a/examples/datafusion-ffi-example/src/logical_extension_codec.rs b/examples/datafusion-ffi-example/src/logical_extension_codec.rs index da9efb297..8c3976d37 100644 --- a/examples/datafusion-ffi-example/src/logical_extension_codec.rs +++ b/examples/datafusion-ffi-example/src/logical_extension_codec.rs @@ -147,7 +147,6 @@ impl MyLogicalExtensionCodec { let ctx_provider = bare_session as Arc; let ffi = FFI_LogicalExtensionCodec::new(inner, Some(runtime), &ctx_provider); - let name = cr"datafusion_logical_extension_codec".into(); - PyCapsule::new(py, ffi, Some(name)) + PyCapsule::new_with_value(py, ffi, cr"datafusion_logical_extension_codec") } } diff --git a/examples/datafusion-ffi-example/src/physical_extension_codec.rs b/examples/datafusion-ffi-example/src/physical_extension_codec.rs index b1a586d9e..35ef77f6b 100644 --- a/examples/datafusion-ffi-example/src/physical_extension_codec.rs +++ b/examples/datafusion-ffi-example/src/physical_extension_codec.rs @@ -24,7 +24,9 @@ use datafusion::logical_expr::ScalarUDF; use datafusion::physical_plan::ExecutionPlan; use datafusion::prelude::SessionContext; use datafusion_ffi::proto::physical_extension_codec::FFI_PhysicalExtensionCodec; -use datafusion_proto::physical_plan::{DefaultPhysicalExtensionCodec, PhysicalExtensionCodec}; +use datafusion_proto::physical_plan::{ + DefaultPhysicalExtensionCodec, PhysicalExtensionCodec, PhysicalProtoConverterExtension, +}; use datafusion_python_util::get_tokio_runtime; use pyo3::prelude::*; use pyo3::types::PyCapsule; @@ -51,12 +53,18 @@ impl PhysicalExtensionCodec for CountingPhysicalExtensionCodec { buf: &[u8], inputs: &[Arc], ctx: &TaskContext, + proto_converter: &dyn PhysicalProtoConverterExtension, ) -> Result> { - self.inner.try_decode(buf, inputs, ctx) + self.inner.try_decode(buf, inputs, ctx, proto_converter) } - fn try_encode(&self, node: Arc, buf: &mut Vec) -> Result<()> { - self.inner.try_encode(node, buf) + fn try_encode( + &self, + node: Arc, + buf: &mut Vec, + proto_converter: &dyn PhysicalProtoConverterExtension, + ) -> Result<()> { + self.inner.try_encode(node, buf, proto_converter) } fn try_decode_udf(&self, name: &str, buf: &[u8]) -> Result> { @@ -113,7 +121,6 @@ impl MyPhysicalExtensionCodec { let ctx_provider = bare_session as Arc; let ffi = FFI_PhysicalExtensionCodec::new(inner, Some(runtime), &ctx_provider); - let name = cr"datafusion_physical_extension_codec".into(); - PyCapsule::new(py, ffi, Some(name)) + PyCapsule::new_with_value(py, ffi, cr"datafusion_physical_extension_codec") } } diff --git a/examples/datafusion-ffi-example/src/physical_optimizer.rs b/examples/datafusion-ffi-example/src/physical_optimizer.rs index 0acd1bb4a..e17510495 100644 --- a/examples/datafusion-ffi-example/src/physical_optimizer.rs +++ b/examples/datafusion-ffi-example/src/physical_optimizer.rs @@ -92,7 +92,6 @@ impl MyPhysicalOptimizerRule { let runtime = get_tokio_runtime().handle().clone(); let ffi = FFI_PhysicalOptimizerRule::new(rule, Some(runtime)); - let name = cr"datafusion_physical_optimizer_rule".into(); - PyCapsule::new(py, ffi, Some(name)) + PyCapsule::new_with_value(py, ffi, cr"datafusion_physical_optimizer_rule") } } diff --git a/examples/datafusion-ffi-example/src/scalar_udf.rs b/examples/datafusion-ffi-example/src/scalar_udf.rs index a3c65e875..85b884ec6 100644 --- a/examples/datafusion-ffi-example/src/scalar_udf.rs +++ b/examples/datafusion-ffi-example/src/scalar_udf.rs @@ -50,12 +50,10 @@ impl IsNullUDF { } fn __datafusion_scalar_udf__<'py>(&self, py: Python<'py>) -> PyResult> { - let name = cr"datafusion_scalar_udf".into(); - let func = Arc::new(ScalarUDF::from(self.clone())); let provider = FFI_ScalarUDF::from(func); - PyCapsule::new(py, provider, Some(name)) + PyCapsule::new_with_value(py, provider, cr"datafusion_scalar_udf") } } diff --git a/examples/datafusion-ffi-example/src/table_function.rs b/examples/datafusion-ffi-example/src/table_function.rs index ed3ef142b..55543cb59 100644 --- a/examples/datafusion-ffi-example/src/table_function.rs +++ b/examples/datafusion-ffi-example/src/table_function.rs @@ -47,13 +47,11 @@ impl MyTableFunction { py: Python<'py>, session: Bound, ) -> PyResult> { - let name = cr"datafusion_table_function".into(); - let func = self.clone(); let codec = ffi_logical_codec_from_pycapsule(session)?; let provider = FFI_TableFunction::new_with_ffi_codec(Arc::new(func), None, codec); - PyCapsule::new(py, provider, Some(name)) + PyCapsule::new_with_value(py, provider, cr"datafusion_table_function") } } diff --git a/examples/datafusion-ffi-example/src/table_provider.rs b/examples/datafusion-ffi-example/src/table_provider.rs index 358ef7402..5756e6d02 100644 --- a/examples/datafusion-ffi-example/src/table_provider.rs +++ b/examples/datafusion-ffi-example/src/table_provider.rs @@ -99,8 +99,6 @@ impl MyTableProvider { py: Python<'py>, session: Bound, ) -> PyResult> { - let name = cr"datafusion_table_provider".into(); - let provider = self .create_table() .map_err(|e: DataFusionError| PyRuntimeError::new_err(e.to_string()))?; @@ -109,6 +107,6 @@ impl MyTableProvider { let provider = FFI_TableProvider::new_with_ffi_codec(Arc::new(provider), false, None, codec); - PyCapsule::new(py, provider, Some(name)) + PyCapsule::new_with_value(py, provider, cr"datafusion_table_provider") } } diff --git a/examples/datafusion-ffi-example/src/table_provider_factory.rs b/examples/datafusion-ffi-example/src/table_provider_factory.rs index 53248a905..71dfd73ca 100644 --- a/examples/datafusion-ffi-example/src/table_provider_factory.rs +++ b/examples/datafusion-ffi-example/src/table_provider_factory.rs @@ -77,11 +77,10 @@ impl MyTableProviderFactory { py: Python<'py>, codec: Bound, ) -> PyResult> { - let name = cr"datafusion_table_provider_factory".into(); let codec = ffi_logical_codec_from_pycapsule(codec)?; let factory = Arc::clone(&self.inner) as Arc; let factory = FFI_TableProviderFactory::new_with_ffi_codec(factory, None, codec); - PyCapsule::new(py, factory, Some(name)) + PyCapsule::new_with_value(py, factory, cr"datafusion_table_provider_factory") } } diff --git a/examples/datafusion-ffi-example/src/window_udf.rs b/examples/datafusion-ffi-example/src/window_udf.rs index f33a166ed..2956ad64c 100644 --- a/examples/datafusion-ffi-example/src/window_udf.rs +++ b/examples/datafusion-ffi-example/src/window_udf.rs @@ -45,12 +45,10 @@ impl MyRankUDF { } fn __datafusion_window_udf__<'py>(&self, py: Python<'py>) -> PyResult> { - let name = cr"datafusion_window_udf".into(); - let func = Arc::new(WindowUDF::from(self.clone())); let provider = FFI_WindowUDF::from(func); - PyCapsule::new(py, provider, Some(name)) + PyCapsule::new_with_value(py, provider, cr"datafusion_window_udf") } } diff --git a/python/datafusion/context.py b/python/datafusion/context.py index 0bfc59bfe..94b2bb1c6 100644 --- a/python/datafusion/context.py +++ b/python/datafusion/context.py @@ -54,6 +54,8 @@ from typing_extensions import deprecated # Python 3.12 +from urllib.parse import urlparse + import pyarrow as pa from datafusion.catalog import ( @@ -611,6 +613,45 @@ def deregister_object_store(self, schema: str, host: str | None = None) -> None: """ self.ctx.deregister_object_store(schema, host) + def _register_object_store_for_path( + self, path: str | pathlib.Path, store: Any + ) -> None: + """Parse a URL path and register the given object store for its scheme and host. + + This is a convenience helper used by methods like + :py:meth:`register_parquet` and :py:meth:`read_parquet` to + automatically register an object store when an ``object_store`` + parameter is provided. + + Args: + path: A URL-style path (e.g. ``"s3://bucket/key.parquet"`` or + ``"file:///tmp/data.parquet"``). + store: An object store instance to register. + + Raises: + ValueError: If the path does not contain a URL scheme, or if + a non-file scheme is missing a host/bucket component. + """ + parsed = urlparse(str(path)) + if not parsed.scheme: + msg = ( + f"Cannot determine object store URL from path {path!r}. " + "The path must use a URL scheme (e.g. 's3://bucket/key')." + ) + raise ValueError(msg) + # file:// URLs typically have an empty netloc (e.g. file:///tmp/a.parquet) + # For other schemes (s3, gs, az, https) the netloc (bucket/host) is required. + if parsed.scheme != "file" and not parsed.netloc: + msg = ( + f"Cannot determine object store URL from path {path!r}. " + "The path must include a host or bucket " + "(e.g. 's3://bucket/key')." + ) + raise ValueError(msg) + scheme = f"{parsed.scheme}://" + host = parsed.netloc or None + self.register_object_store(scheme, store, host=host) + def register_listing_table( self, name: str, @@ -1028,6 +1069,7 @@ def register_parquet( skip_metadata: bool = True, schema: pa.Schema | None = None, file_sort_order: Sequence[Sequence[SortKey]] | None = None, + object_store: Any | None = None, ) -> None: """Register a Parquet file as a table. @@ -1049,7 +1091,41 @@ def register_parquet( file_sort_order: Sort order for the file. Each sort key can be specified as a column name (``str``), an expression (``Expr``), or a ``SortExpr``. - """ + object_store: A pre-configured object store instance (e.g. + :py:class:`~datafusion.object_store.AmazonS3`, + :py:class:`~datafusion.object_store.GoogleCloud`, + :py:class:`~datafusion.object_store.MicrosoftAzure`) to use + for accessing the file. When provided, the store is + automatically registered for the URL scheme and host parsed + from ``path``, removing the need to call + :py:meth:`register_object_store` separately. This is + especially useful in multi-threaded environments where + setting credentials via ``os.environ`` is not thread-safe. + + Examples: + Register a local Parquet file: + + >>> import datafusion + >>> ctx = datafusion.SessionContext() + >>> ctx.register_parquet("my_table", "data.parquet") # doctest: +SKIP + + Register from S3 with inline credentials (thread-safe): + + >>> from datafusion.object_store import AmazonS3 # doctest: +SKIP + >>> store = AmazonS3( + ... bucket_name="my-bucket", + ... region="us-east-1", + ... access_key_id="...", + ... secret_access_key="...", + ... ) # doctest: +SKIP + >>> ctx.register_parquet( + ... "my_table", + ... "s3://my-bucket/data.parquet", + ... object_store=store, + ... ) # doctest: +SKIP + """ + if object_store is not None: + self._register_object_store_for_path(path, object_store) if table_partition_cols is None: table_partition_cols = [] table_partition_cols = _convert_table_partition_cols(table_partition_cols) @@ -1075,6 +1151,7 @@ def register_csv( file_extension: str = ".csv", file_compression_type: str | None = None, options: CsvReadOptions | None = None, + object_store: Any | None = None, ) -> None: """Register a CSV file as a table. @@ -1096,7 +1173,14 @@ def register_csv( file_compression_type: File compression type. options: Set advanced options for CSV reading. This cannot be combined with any of the other options in this method. - """ + object_store: A pre-configured object store instance to use for + accessing the file. When provided, the store is automatically + registered for the URL scheme and host parsed from ``path``. + """ + if object_store is not None: + # For list paths, register from the first entry + register_path = path[0] if isinstance(path, list) else path + self._register_object_store_for_path(register_path, object_store) if options is not None and ( schema is not None or not has_header @@ -1143,6 +1227,7 @@ def register_json( file_extension: str = ".json", table_partition_cols: list[tuple[str, str | pa.DataType]] | None = None, file_compression_type: str | None = None, + object_store: Any | None = None, ) -> None: """Register a JSON file as a table. @@ -1159,7 +1244,12 @@ def register_json( selected for data input. table_partition_cols: Partition columns. file_compression_type: File compression type. + object_store: A pre-configured object store instance to use for + accessing the file. When provided, the store is automatically + registered for the URL scheme and host parsed from ``path``. """ + if object_store is not None: + self._register_object_store_for_path(path, object_store) if table_partition_cols is None: table_partition_cols = [] table_partition_cols = _convert_table_partition_cols(table_partition_cols) @@ -1180,6 +1270,7 @@ def register_avro( schema: pa.Schema | None = None, file_extension: str = ".avro", table_partition_cols: list[tuple[str, str | pa.DataType]] | None = None, + object_store: Any | None = None, ) -> None: """Register an Avro file as a table. @@ -1192,7 +1283,12 @@ def register_avro( schema: The data source schema. file_extension: File extension to select. table_partition_cols: Partition columns. + object_store: A pre-configured object store instance to use for + accessing the file. When provided, the store is automatically + registered for the URL scheme and host parsed from ``path``. """ + if object_store is not None: + self._register_object_store_for_path(path, object_store) if table_partition_cols is None: table_partition_cols = [] table_partition_cols = _convert_table_partition_cols(table_partition_cols) @@ -1205,6 +1301,7 @@ def register_arrow( schema: pa.Schema | None = None, file_extension: str = ".arrow", table_partition_cols: list[tuple[str, str | pa.DataType]] | None = None, + object_store: Any | None = None, ) -> None: """Register an Arrow IPC file as a table. @@ -1217,6 +1314,9 @@ def register_arrow( schema: The data source schema. file_extension: File extension to select. table_partition_cols: Partition columns. + object_store: A pre-configured object store instance to use for + accessing the file. When provided, the store is automatically + registered for the URL scheme and host parsed from ``path``. Examples: >>> import tempfile, os @@ -1271,6 +1371,8 @@ def register_arrow( 30 ] """ + if object_store is not None: + self._register_object_store_for_path(path, object_store) if table_partition_cols is None: table_partition_cols = [] table_partition_cols = _convert_table_partition_cols(table_partition_cols) @@ -1690,6 +1792,7 @@ def read_json( file_extension: str = ".json", table_partition_cols: list[tuple[str, str | pa.DataType]] | None = None, file_compression_type: str | None = None, + object_store: Any | None = None, ) -> DataFrame: """Read a line-delimited JSON data source. @@ -1702,10 +1805,15 @@ def read_json( selected for data input. table_partition_cols: Partition columns. file_compression_type: File compression type. + object_store: A pre-configured object store instance to use for + accessing the file. When provided, the store is automatically + registered for the URL scheme and host parsed from ``path``. Returns: DataFrame representation of the read JSON files. """ + if object_store is not None: + self._register_object_store_for_path(path, object_store) if table_partition_cols is None: table_partition_cols = [] table_partition_cols = _convert_table_partition_cols(table_partition_cols) @@ -1731,6 +1839,7 @@ def read_csv( table_partition_cols: list[tuple[str, str | pa.DataType]] | None = None, file_compression_type: str | None = None, options: CsvReadOptions | None = None, + object_store: Any | None = None, ) -> DataFrame: """Read a CSV data source. @@ -1750,10 +1859,16 @@ def read_csv( file_compression_type: File compression type. options: Set advanced options for CSV reading. This cannot be combined with any of the other options in this method. + object_store: A pre-configured object store instance to use for + accessing the file. When provided, the store is automatically + registered for the URL scheme and host parsed from ``path``. Returns: DataFrame representation of the read CSV files """ + if object_store is not None: + register_path = path[0] if isinstance(path, list) else path + self._register_object_store_for_path(register_path, object_store) if options is not None and ( schema is not None or not has_header @@ -1803,6 +1918,7 @@ def read_parquet( skip_metadata: bool = True, schema: pa.Schema | None = None, file_sort_order: Sequence[Sequence[SortKey]] | None = None, + object_store: Any | None = None, ) -> DataFrame: """Read a Parquet source into a :py:class:`~datafusion.dataframe.Dataframe`. @@ -1822,10 +1938,43 @@ def read_parquet( file_sort_order: Sort order for the file. Each sort key can be specified as a column name (``str``), an expression (``Expr``), or a ``SortExpr``. + object_store: A pre-configured object store instance (e.g. + :py:class:`~datafusion.object_store.AmazonS3`, + :py:class:`~datafusion.object_store.GoogleCloud`, + :py:class:`~datafusion.object_store.MicrosoftAzure`) to use + for accessing the file. When provided, the store is + automatically registered for the URL scheme and host parsed + from ``path``, removing the need to call + :py:meth:`register_object_store` separately. This is + especially useful in multi-threaded environments where + setting credentials via ``os.environ`` is not thread-safe. Returns: DataFrame representation of the read Parquet files - """ + + Examples: + Read a local Parquet file: + + >>> import datafusion + >>> ctx = datafusion.SessionContext() + >>> df = ctx.read_parquet("data.parquet") # doctest: +SKIP + + Read from S3 with inline credentials (thread-safe): + + >>> from datafusion.object_store import AmazonS3 # doctest: +SKIP + >>> store = AmazonS3( + ... bucket_name="my-bucket", + ... region="us-east-1", + ... access_key_id="...", + ... secret_access_key="...", + ... ) # doctest: +SKIP + >>> df = ctx.read_parquet( + ... "s3://my-bucket/data.parquet", + ... object_store=store, + ... ) # doctest: +SKIP + """ + if object_store is not None: + self._register_object_store_for_path(path, object_store) if table_partition_cols is None: table_partition_cols = [] table_partition_cols = _convert_table_partition_cols(table_partition_cols) @@ -1848,6 +1997,7 @@ def read_avro( schema: pa.Schema | None = None, file_partition_cols: list[tuple[str, str | pa.DataType]] | None = None, file_extension: str = ".avro", + object_store: Any | None = None, ) -> DataFrame: """Create a :py:class:`DataFrame` for reading Avro data source. @@ -1856,10 +2006,15 @@ def read_avro( schema: The data source schema. file_partition_cols: Partition columns. file_extension: File extension to select. + object_store: A pre-configured object store instance to use for + accessing the file. When provided, the store is automatically + registered for the URL scheme and host parsed from ``path``. Returns: DataFrame representation of the read Avro file """ + if object_store is not None: + self._register_object_store_for_path(path, object_store) if file_partition_cols is None: file_partition_cols = [] file_partition_cols = _convert_table_partition_cols(file_partition_cols) @@ -1873,6 +2028,7 @@ def read_arrow( schema: pa.Schema | None = None, file_extension: str = ".arrow", file_partition_cols: list[tuple[str, str | pa.DataType]] | None = None, + object_store: Any | None = None, ) -> DataFrame: """Create a :py:class:`DataFrame` for reading an Arrow IPC data source. @@ -1881,6 +2037,9 @@ def read_arrow( schema: The data source schema. file_extension: File extension to select. file_partition_cols: Partition columns. + object_store: A pre-configured object store instance to use for + accessing the file. When provided, the store is automatically + registered for the URL scheme and host parsed from ``path``. Returns: DataFrame representation of the read Arrow IPC file. @@ -1932,6 +2091,8 @@ def read_arrow( 3 ] """ + if object_store is not None: + self._register_object_store_for_path(path, object_store) if file_partition_cols is None: file_partition_cols = [] file_partition_cols = _convert_table_partition_cols(file_partition_cols) diff --git a/python/datafusion/expr.py b/python/datafusion/expr.py index d2560fbc4..c198b646c 100644 --- a/python/datafusion/expr.py +++ b/python/datafusion/expr.py @@ -49,6 +49,11 @@ from collections.abc import Callable, Iterable, Sequence from typing import TYPE_CHECKING, Any, ClassVar +try: + from warnings import deprecated # Python 3.13+ +except ImportError: + from typing_extensions import deprecated # Python 3.12 + import pyarrow as pa from ._internal import expr as expr_internal @@ -90,6 +95,27 @@ CreateCatalog = expr_internal.CreateCatalog CreateCatalogSchema = expr_internal.CreateCatalogSchema CreateExternalTable = expr_internal.CreateExternalTable + + +@deprecated("CreateExternalTable.location() is deprecated; use locations() instead.") +def _create_external_table_location(self: Any) -> str: + """Return the first external table location. + + Use :meth:`CreateExternalTable.locations` instead. + + Examples: + >>> class Command: + ... def locations(self) -> list[str]: + ... return ["data.csv"] + >>> Command().locations() + ['data.csv'] + """ + locations = self.locations() + return locations[0] if locations else "" + + +CreateExternalTable.location = _create_external_table_location + CreateFunction = expr_internal.CreateFunction CreateFunctionBody = expr_internal.CreateFunctionBody CreateIndex = expr_internal.CreateIndex @@ -1226,6 +1252,12 @@ def isnan(self) -> Expr: return F.isnan(self) + def is_nan(self) -> Expr: + """Returns true if a given number is +NaN or -NaN otherwise returns false.""" + from . import functions as F + + return F.is_nan(self) + def degrees(self) -> Expr: """Converts the argument from radians to degrees.""" from . import functions as F diff --git a/python/datafusion/functions/__init__.py b/python/datafusion/functions/__init__.py index 9fb20dc15..291957490 100644 --- a/python/datafusion/functions/__init__.py +++ b/python/datafusion/functions/__init__.py @@ -217,6 +217,7 @@ def _warn_if_expr_for_literal_arg( "initcap", "inner_product", "instr", + "is_nan", "isnan", "iszero", "lag", @@ -413,6 +414,11 @@ def isnan(expr: Expr) -> Expr: return Expr(f.isnan(expr.expr)) +def is_nan(expr: Expr) -> Expr: + """Alias for :func:`isnan`.""" + return isnan(expr) + + def nullif(expr1: Expr, expr2: Expr) -> Expr: """Returns NULL if expr1 equals expr2; otherwise it returns expr1. diff --git a/python/datafusion/functions/spark.py b/python/datafusion/functions/spark.py index 5163a6701..0a0f41400 100644 --- a/python/datafusion/functions/spark.py +++ b/python/datafusion/functions/spark.py @@ -273,12 +273,10 @@ def slice(x: Expr, start: Expr | int, length: Expr | int) -> Expr: Examples: >>> ctx = dfn.SessionContext() - >>> df = ctx.from_pydict({"x": [1]}) + >>> df = ctx.from_pydict({"x": [[1, 2, 3, 4]]}) >>> r = df.select( ... dfn.functions.spark.slice( - ... dfn.functions.spark.array( - ... dfn.lit(1), dfn.lit(2), dfn.lit(3), dfn.lit(4)), - ... 2, 2, + ... dfn.col("x"), 2, 2, ... ).alias("v") ... ) >>> r.collect_column("v")[0].as_py() diff --git a/python/datafusion/user_defined.py b/python/datafusion/user_defined.py index 81a516af8..394c682ae 100644 --- a/python/datafusion/user_defined.py +++ b/python/datafusion/user_defined.py @@ -249,7 +249,8 @@ def udf(*args: Any, **kwargs: Any): # noqa: D417 input_fields (list[pa.Field | pa.DataType]): The data types or Fields of the arguments to ``func``. This list must be of the same length as the number of arguments. - return_field (_R): The field of the return value from the function. + return_field (pa.DataType | pa.Field): The field of the return value + from the function. volatility (Volatility | str): See `Volatility` for allowed values. name (Optional[str]): A descriptive name for the function. diff --git a/python/tests/test_catalog.py b/python/tests/test_catalog.py index c89da36bf..97a898825 100644 --- a/python/tests/test_catalog.py +++ b/python/tests/test_catalog.py @@ -123,6 +123,11 @@ def register_catalog( class CustomTableProviderFactory(dfn.catalog.TableProviderFactory): def create(self, cmd: dfn.expr.CreateExternalTable): assert cmd.name() == "test_table_factory" + assert cmd.locations() == ["foo"] + + with pytest.warns(DeprecationWarning, match=r"location\(\).+deprecated"): + assert cmd.location() == "foo" + return create_dataset() diff --git a/python/tests/test_context.py b/python/tests/test_context.py index 112a6fd7b..7d038c7a5 100644 --- a/python/tests/test_context.py +++ b/python/tests/test_context.py @@ -17,6 +17,7 @@ import datetime as dt import gzip import pathlib +import shutil import pyarrow as pa import pyarrow.dataset as ds @@ -116,6 +117,22 @@ def test_register_record_batches(ctx): assert result[0].column(1) == pa.array([-3, -3, -3]) +def test_register_record_batches_empty(ctx): + # A partition list with no record batches carries no schema, so this used to + # panic on unchecked `[0][0]` indexing. It should now raise a clear error. + with pytest.raises(ValueError, match="no record batches"): + ctx.register_record_batches("t", [[]]) + + # An empty outer partition list carries no schema either, and raises the same error. + with pytest.raises(ValueError, match="no record batches"): + ctx.register_record_batches("t", []) + + # The schema is still recovered from a later non-empty partition. + batch = pa.RecordBatch.from_arrays([pa.array([1, 2, 3])], names=["a"]) + ctx.register_record_batches("t2", [[], [batch]]) + assert ctx.sql("SELECT a FROM t2").collect()[0].column(0) == pa.array([1, 2, 3]) + + def test_create_dataframe_registers_unique_table_name(ctx): # create a RecordBatch and register it as memtable batch = pa.RecordBatch.from_arrays( @@ -783,16 +800,15 @@ def test_read_csv(ctx): csv_df.select(column("c1")).show() -def test_read_csv_list(ctx): - csv_df = ctx.read_csv(path=["testing/data/csv/aggregate_test_100.csv"]) +def test_read_csv_list(ctx, tmp_path): + source_path = pathlib.Path("testing/data/csv/aggregate_test_100.csv") + copied_path = tmp_path / source_path.name + shutil.copy(source_path, copied_path) + + csv_df = ctx.read_csv(path=[source_path]) expected = csv_df.count() * 2 - double_csv_df = ctx.read_csv( - path=[ - "testing/data/csv/aggregate_test_100.csv", - "testing/data/csv/aggregate_test_100.csv", - ] - ) + double_csv_df = ctx.read_csv(path=[source_path, copied_path]) actual = double_csv_df.count() double_csv_df.select(column("c1")).show() diff --git a/python/tests/test_expr.py b/python/tests/test_expr.py index 606c6a984..ef006bd91 100644 --- a/python/tests/test_expr.py +++ b/python/tests/test_expr.py @@ -115,6 +115,8 @@ def test_limit(test_ctx): plan = plan.to_variant() assert isinstance(plan, Limit) assert "Skip: None" in str(plan) + assert plan.skip() is None + assert plan.fetch().python_value().as_py() == 10 df = test_ctx.sql("select c1 from test LIMIT 10 OFFSET 5") plan = df.logical_plan() @@ -122,6 +124,8 @@ def test_limit(test_ctx): plan = plan.to_variant() assert isinstance(plan, Limit) assert "Skip: Some(Literal(Int64(5), None))" in str(plan) + assert plan.skip().python_value().as_py() == 5 + assert plan.fetch().python_value().as_py() == 10 def test_aggregate_query(test_ctx): @@ -500,6 +504,11 @@ def test_alias_with_metadata(df): pa.array([False, True, False, None], type=pa.bool_()), id="isnan", ), + pytest.param( + col("e").is_nan(), + pa.array([False, True, False, None], type=pa.bool_()), + id="is_nan", + ), pytest.param( functions.round(col("a").degrees(), lit(4)), pa.array([-42.9718, 28.6479, 0.0, None], type=pa.float64()), diff --git a/python/tests/test_functions.py b/python/tests/test_functions.py index 894d81fcf..fabe6dff2 100644 --- a/python/tests/test_functions.py +++ b/python/tests/test_functions.py @@ -145,7 +145,7 @@ def test_math_functions(): f.pow(col_v, literal(pa.scalar(4))), f.round(col_v), f.round(col_v, literal(pa.scalar(3))), - f.sqrt(col_v), + f.sqrt(f.abs(col_v)), f.signum(col_v), f.trunc(col_v), f.asinh(col_v), @@ -189,7 +189,7 @@ def test_math_functions(): np.testing.assert_array_almost_equal(result.column(16), np.power(values, 4)) np.testing.assert_array_almost_equal(result.column(17), np.round(values)) np.testing.assert_array_almost_equal(result.column(18), np.round(values, 3)) - np.testing.assert_array_almost_equal(result.column(19), np.sqrt(values)) + np.testing.assert_array_almost_equal(result.column(19), np.sqrt(np.abs(values))) np.testing.assert_array_almost_equal(result.column(20), np.sign(values)) np.testing.assert_array_almost_equal(result.column(21), np.trunc(values)) np.testing.assert_array_almost_equal(result.column(22), np.arcsinh(values)) @@ -215,6 +215,23 @@ def test_math_functions(): ) +def test_sqrt_rejects_negative_input(): + ctx = SessionContext() + df = ctx.from_pydict({"value": [-1.0]}) + + with pytest.raises(Exception, match="cannot take square root of a negative number"): + df.select(f.sqrt(column("value"))).collect() + + +def test_is_nan_alias(): + ctx = SessionContext() + df = ctx.from_pydict({"value": [1.0, np.nan, None]}) + + result = df.select(f.is_nan(column("value")).alias("is_nan")).to_pydict() + + assert result == {"is_nan": [False, True, None]} + + def py_indexof(arr, v): try: return arr.index(v) + 1 diff --git a/python/tests/test_io.py b/python/tests/test_io.py index 9f56f74d7..f0a9c3aa1 100644 --- a/python/tests/test_io.py +++ b/python/tests/test_io.py @@ -15,6 +15,7 @@ # specific language governing permissions and limitations # under the License. +import shutil from pathlib import Path import pyarrow as pa @@ -72,16 +73,15 @@ def test_read_csv(): csv_df.select(column("c1")).show() -def test_read_csv_list(): - csv_df = read_csv(path=["testing/data/csv/aggregate_test_100.csv"]) +def test_read_csv_list(tmp_path): + source_path = Path("testing/data/csv/aggregate_test_100.csv") + copied_path = tmp_path / source_path.name + shutil.copy(source_path, copied_path) + + csv_df = read_csv(path=[source_path]) expected = csv_df.count() * 2 - double_csv_df = read_csv( - path=[ - "testing/data/csv/aggregate_test_100.csv", - "testing/data/csv/aggregate_test_100.csv", - ] - ) + double_csv_df = read_csv(path=[source_path, copied_path]) actual = double_csv_df.count() double_csv_df.select(column("c1")).show() diff --git a/python/tests/test_lambda.py b/python/tests/test_lambda.py index 68be22a04..ce37546be 100644 --- a/python/tests/test_lambda.py +++ b/python/tests/test_lambda.py @@ -16,6 +16,8 @@ # under the License. """Tests for lambda expressions and higher-order array functions.""" +import pickle + import pytest from datafusion import SessionConfig, SessionContext, col, lit from datafusion import functions as f @@ -137,8 +139,9 @@ def test_sql_lambda_keyword_syntax(dialect): assert result.to_pylist() == [[2, 4, 6]] -def test_pickle_lambda_expr_not_supported(): - # v1 limitation: upstream proto serialization rejects lambda expressions. +def test_pickle_lambda_expr_round_trip(df): expr = f.array_transform(col("a"), lambda v: v * 2) - with pytest.raises(Exception, match="Lambda not implemented"): - expr.to_bytes() + decoded = pickle.loads(pickle.dumps(expr)) # noqa: S301 + + assert decoded.canonical_name() == expr.canonical_name() + assert _column(df, decoded, "r") == [[2, 4, 6], [8, 10]] diff --git a/python/tests/test_object_store_param.py b/python/tests/test_object_store_param.py new file mode 100644 index 000000000..13d51d792 --- /dev/null +++ b/python/tests/test_object_store_param.py @@ -0,0 +1,140 @@ +# 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. + +"""Tests for the object_store parameter on register/read file methods.""" + +import contextlib +from pathlib import Path +from unittest.mock import MagicMock, patch + +import pyarrow as pa +import pyarrow.parquet as pq +import pytest +from datafusion import SessionContext +from datafusion.object_store import LocalFileSystem + + +@pytest.fixture +def ctx(): + return SessionContext() + + +@pytest.mark.parametrize( + ("path", "scheme", "host"), + [ + ("s3://my-bucket/path/file.parquet", "s3://", "my-bucket"), + ("gs://my-gcs-bucket/data.parquet", "gs://", "my-gcs-bucket"), + ("az://my-container/data.parquet", "az://", "my-container"), + ("https://example.com/data.parquet", "https://", "example.com"), + ("file:///tmp/data.parquet", "file://", None), + ], +) +def test_register_object_store_for_url(ctx, path, scheme, host): + store = MagicMock() + + with patch.object(ctx, "register_object_store") as register: + ctx._register_object_store_for_path(path, store) + + register.assert_called_once_with(scheme, store, host=host) + + +@pytest.mark.parametrize( + ("path", "error"), + [ + ("/local/path/file.parquet", "Cannot determine object store URL"), + ("relative/path.parquet", "Cannot determine object store URL"), + (Path("/local/file.parquet"), "Cannot determine object store URL"), + ("C:\\Users\\data\\file.parquet", "must include a host or bucket"), + ("s3:///key.parquet", "must include a host or bucket"), + ], +) +def test_register_object_store_rejects_invalid_url(path, error, ctx): + with pytest.raises(ValueError, match=error): + ctx._register_object_store_for_path(path, MagicMock()) + + +@pytest.mark.parametrize( + ("method_name", "args", "path"), + [ + ("register_parquet", ("table",), "s3://bucket/data.parquet"), + ("read_parquet", (), "s3://bucket/data.parquet"), + ("register_csv", ("table",), "s3://bucket/data.csv"), + ("read_csv", (), "s3://bucket/data.csv"), + ("register_json", ("table",), "s3://bucket/data.json"), + ("read_json", (), "s3://bucket/data.json"), + ("register_avro", ("table",), "s3://bucket/data.avro"), + ("read_avro", (), "s3://bucket/data.avro"), + ("register_arrow", ("table",), "s3://bucket/data.arrow"), + ("read_arrow", (), "s3://bucket/data.arrow"), + ], +) +def test_file_methods_register_object_store(ctx, method_name, args, path): + store = MagicMock() + + # The remote file does not exist. Registration happens before DataFusion + # tries to inspect it, which is the behavior under test. + with ( + patch.object(ctx, "register_object_store") as register, + contextlib.suppress(Exception), + ): + getattr(ctx, method_name)(*args, path, object_store=store) + + register.assert_called_once_with("s3://", store, host="bucket") + + +def test_register_csv_uses_first_path_for_object_store(ctx): + store = MagicMock() + paths = ["s3://bucket/a.csv", "s3://bucket/b.csv"] + + with ( + patch.object(ctx, "register_object_store") as register, + contextlib.suppress(Exception), + ): + ctx.register_csv("table", paths, object_store=store) + + register.assert_called_once_with("s3://", store, host="bucket") + + +def test_object_store_none_does_not_register(ctx): + with ( + patch.object(ctx, "register_object_store") as register, + contextlib.suppress(Exception), + ): + ctx.register_parquet("table", "missing.parquet") + + register.assert_not_called() + + +def test_file_method_rejects_local_path_with_object_store(ctx): + with pytest.raises(ValueError, match="Cannot determine object store URL"): + ctx.register_parquet("table", "/local/file.parquet", object_store=MagicMock()) + + +@pytest.mark.parametrize("method_name", ["register_parquet", "read_parquet"]) +def test_parquet_methods_with_local_object_store(ctx, tmp_path, method_name): + table = pa.table({"value": [10, 20, 30]}) + parquet_path = tmp_path / "data.parquet" + pq.write_table(table, parquet_path) + + path = parquet_path.as_uri() + if method_name == "register_parquet": + ctx.register_parquet("test_table", path, object_store=LocalFileSystem()) + dataframe = ctx.sql("SELECT * FROM test_table") + else: + dataframe = ctx.read_parquet(path, object_store=LocalFileSystem()) + + assert dataframe.collect()[0].column("value").to_pylist() == [10, 20, 30] diff --git a/python/tests/test_spark_functions.py b/python/tests/test_spark_functions.py index aed46d7fa..39735a543 100644 --- a/python/tests/test_spark_functions.py +++ b/python/tests/test_spark_functions.py @@ -172,6 +172,10 @@ def test_array_and_size(df): def test_slice(df): + assert _val(df, spark.slice(col("a"), lit(2), lit(2))) == [2, 3] + + +def test_slice_spark_array(df): arr = spark.array(lit(1), lit(2), lit(3), lit(4)) assert _val(df, spark.slice(arr, lit(2), lit(2))) == [2, 3] diff --git a/python/tests/test_sql.py b/python/tests/test_sql.py index 6089c377d..57d22bd31 100644 --- a/python/tests/test_sql.py +++ b/python/tests/test_sql.py @@ -107,6 +107,7 @@ def test_register_csv(ctx, tmp_path): def test_register_csv_list(ctx, tmp_path): path = tmp_path / "test.csv" + second_path = tmp_path / "test2.csv" int_values = [1, 2, 3, 4] table = pa.Table.from_arrays( @@ -118,6 +119,7 @@ def test_register_csv_list(ctx, tmp_path): names=["int", "str", "float"], ) write_csv(table, path) + write_csv(table, second_path) ctx.register_csv("csv", path) csv_df = ctx.table("csv") @@ -126,7 +128,7 @@ def test_register_csv_list(ctx, tmp_path): "double_csv", path=[ path, - path, + second_path, ], ) diff --git a/rust-toolchain.toml b/rust-toolchain.toml new file mode 100644 index 000000000..5639a821f --- /dev/null +++ b/rust-toolchain.toml @@ -0,0 +1,23 @@ +# 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. + +# This file specifies the default version of Rust used +# to compile this workspace and run CI jobs. + +[toolchain] +channel = "1.97.0" +components = ["rustfmt", "clippy"]