diff --git a/LICENSE-3rdparty.csv b/LICENSE-3rdparty.csv index bc99a4da800..0748c22d08c 100644 --- a/LICENSE-3rdparty.csv +++ b/LICENSE-3rdparty.csv @@ -24,6 +24,16 @@ anyhow,https://github.com/dtolnay/anyhow,MIT OR Apache-2.0,David Tolnay , Manish Goregaokar , Simonas Kazlauskas , Brian L. Troutwine , Corey Farwell " arc-swap,https://github.com/vorner/arc-swap,MIT OR Apache-2.0,Michal 'vorner' Vaner arrayvec,https://github.com/bluss/arrayvec,MIT OR Apache-2.0,bluss +arrow-array,https://github.com/apache/arrow-rs,Apache-2.0 AND MIT,Apache Arrow +arrow-buffer,https://github.com/apache/arrow-rs,Apache-2.0,Apache Arrow +arrow-cast,https://github.com/apache/arrow-rs,Apache-2.0,Apache Arrow +arrow-cmp,https://github.com/apache/arrow-rs,Apache-2.0,Apache Arrow +arrow-data,https://github.com/apache/arrow-rs,Apache-2.0,Apache Arrow +arrow-ipc,https://github.com/apache/arrow-rs,Apache-2.0,Apache Arrow +arrow-json,https://github.com/apache/arrow-rs,Apache-2.0,Apache Arrow +arrow-ord,https://github.com/apache/arrow-rs,Apache-2.0,Apache Arrow +arrow-schema,https://github.com/apache/arrow-rs,Apache-2.0,Apache Arrow +arrow-select,https://github.com/apache/arrow-rs,Apache-2.0,Apache Arrow ascii-canvas,https://github.com/lalrpop/ascii-canvas,Apache-2.0 OR MIT,Niko Matsakis assert-json-diff,https://github.com/davidpdrsn/assert-json-diff,MIT,David Pedersen async-channel,https://github.com/smol-rs/async-channel,Apache-2.0 OR MIT,Stjepan Glavina @@ -261,6 +271,7 @@ fiat-crypto,https://github.com/mit-plv/fiat-crypto,MIT OR Apache-2.0 OR BSD-1-Cl find-msvc-tools,https://github.com/rust-lang/cc-rs,MIT OR Apache-2.0,The find-msvc-tools Authors findshlibs,https://github.com/gimli-rs/findshlibs,MIT OR Apache-2.0,The findshlibs Authors fixedbitset,https://github.com/petgraph/fixedbitset,MIT OR Apache-2.0,bluss +flatbuffers,https://github.com/google/flatbuffers,Apache-2.0,"Robert Winslow , FlatBuffers Maintainers" flate2,https://github.com/rust-lang/flate2-rs,MIT OR Apache-2.0,"Alex Crichton , Josh Triplett " float-cmp,https://github.com/mikedilger/float-cmp,MIT,Mike Dilger fluent-uri,https://github.com/yescallop/fluent-uri-rs,MIT,Scallop Ye @@ -386,6 +397,12 @@ lambda_runtime_api_client,https://github.com/aws/aws-lambda-rust-runtime,Apache- lazy_static,https://github.com/rust-lang-nursery/lazy-static.rs,MIT OR Apache-2.0,Marvin Löbel left-right,https://github.com/jonhoo/left-right,MIT OR Apache-2.0,Jon Gjengset levenshtein_automata,https://github.com/tantivy-search/levenshtein-automata,MIT,Paul Masurel +lexical-core,https://github.com/Alexhuszagh/rust-lexical,MIT OR Apache-2.0,Alex Huszagh +lexical-parse-float,https://github.com/Alexhuszagh/rust-lexical,MIT OR Apache-2.0,Alex Huszagh +lexical-parse-integer,https://github.com/Alexhuszagh/rust-lexical,MIT OR Apache-2.0,Alex Huszagh +lexical-util,https://github.com/Alexhuszagh/rust-lexical,MIT OR Apache-2.0,Alex Huszagh +lexical-write-float,https://github.com/Alexhuszagh/rust-lexical,MIT OR Apache-2.0,Alex Huszagh +lexical-write-integer,https://github.com/Alexhuszagh/rust-lexical,MIT OR Apache-2.0,Alex Huszagh libc,https://github.com/rust-lang/libc,MIT OR Apache-2.0,The Rust Project Developers libm,https://github.com/rust-lang/compiler-builtins,MIT,"Alex Crichton , Amanieu d'Antras , Jorge Aparicio , Trevor Gross " libredox,https://gitlab.redox-os.org/redox-os/libredox,MIT,4lDO2 <4lDO2@protonmail.com> @@ -442,6 +459,7 @@ normalize-line-endings,https://github.com/derekdreery/normalize-line-endings,Apa nu-ansi-term,https://github.com/nushell/nu-ansi-term,MIT,"ogham@bsago.me, Ryan Scheel (Havvy) , Josh Triplett , The Nushell Project Developers" num,https://github.com/rust-num/num,MIT OR Apache-2.0,The Rust Project Developers num-bigint,https://github.com/rust-num/num-bigint,MIT OR Apache-2.0,The Rust Project Developers +num-bigint,https://github.com/rust-num/num-bigint,MIT OR Apache-2.0,The num-bigint Authors num-bigint-dig,https://github.com/dignifiedquire/num-bigint,MIT OR Apache-2.0,"dignifiedquire , The Rust Project Developers" num-cmp,https://github.com/lifthrasiir/num-cmp,MIT OR Apache-2.0,Kang Seonghoon num-complex,https://github.com/rust-num/num-complex,MIT OR Apache-2.0,The Rust Project Developers @@ -500,6 +518,7 @@ papergrid,https://github.com/zhiburt/tabled,MIT,Maxim Zhiburt , The Rust Project Developers" parking_lot,https://github.com/Amanieu/parking_lot,MIT OR Apache-2.0,Amanieu d'Antras parking_lot_core,https://github.com/Amanieu/parking_lot,MIT OR Apache-2.0,Amanieu d'Antras +parquet,https://github.com/apache/arrow-rs,Apache-2.0,Apache Arrow parse-size,https://github.com/kennytm/parse-size,MIT,kennytm paste,https://github.com/dtolnay/paste,MIT OR Apache-2.0,David Tolnay pbkdf2,https://github.com/RustCrypto/password-hashes/tree/master/pbkdf2,MIT OR Apache-2.0,RustCrypto Developers @@ -667,6 +686,7 @@ security-framework-sys,https://github.com/kornelski/rust-security-framework,MIT seize,https://github.com/ibraheemdev/seize,MIT,Ibraheem Ahmed semver,https://github.com/dtolnay/semver,MIT OR Apache-2.0,David Tolnay separator,https://github.com/saghm/rust-separator,MIT,Saghm Rossi +seq-macro,https://github.com/dtolnay/seq-macro,MIT OR Apache-2.0,David Tolnay serde,https://github.com/serde-rs/serde,MIT OR Apache-2.0,"Erick Tryzelaar , David Tolnay " serde-tuple-vec-map,https://github.com/daboross/serde-tuple-vec-map,MIT,David Ross serde-value,https://github.com/arcnmx/serde-value,MIT,arcnmx diff --git a/docs/internals/parquet-bulk-load.md b/docs/internals/parquet-bulk-load.md new file mode 100644 index 00000000000..27f3a63d430 --- /dev/null +++ b/docs/internals/parquet-bulk-load.md @@ -0,0 +1,38 @@ +# Parquet bulk load (cli only) + +Index a parquet file without merges. + + + +``` +quickwit tool local-ingest --index idx --input-path file:///data/x.parquet --num-pipelines N +``` + +## How it works + +``` +CLI + ├─ plan = ParquetLoadPlan::try_new(file) // parses the footer once + ├─ source_loader: File → ParquetSourceFactory { plan } + ├─ IndexingService::new(...).with_source_loader(source_loader) + └─ SpawnPipeline + DetachIndexingPipeline × N → ParquetSource { plan } +``` + +- Sources pull row groups from the plan, decode them on the +blocking runtime, and emit JSON documents (`arrow_json`). +- Only the CLI replaces the source loader, so servers can't run a Parquet +source. + +## Failure handling + +No checkpoint. On failure, the CLI stops actors, attempts to clear the +index, and reports cleanup errors alongside the load error. +Rollback is best effort: detached uploads can outlive actor shutdown after an +early pipeline failure. Process termination is not covered. + +## Benchmark settings + +Size splits with `split_num_docs_target`; use `no_merge` to keep them +unchanged when running a server afterwards. +`--num-pipelines` defaults to half the CPUs: 16 pipelines was the fastest on an +m8gd.8xlarge (32 vCPUs) with 5M-document splits. \ No newline at end of file diff --git a/docs/reference/cli.md b/docs/reference/cli.md index e4cbaa5e1e3..22ad95cf013 100644 --- a/docs/reference/cli.md +++ b/docs/reference/cli.md @@ -761,7 +761,7 @@ Performs utility operations. Requires a node config. ### tool local-ingest -Indexes NDJSON documents locally. +Indexes NDJSON or Parquet documents locally. `quickwit tool local-ingest [args]` *Synopsis* @@ -771,6 +771,8 @@ quickwit tool local-ingest --index [--input-path ] [--input-format ] + [--num-pipelines ] + [--batch-num-rows ] [--overwrite] [--transform-script ] [--keep-cache] @@ -782,10 +784,14 @@ quickwit tool local-ingest |-----------------|-------------|--------:| | `--index` | ID of the target index | | | `--input-path` | Location of the input file. | | -| `--input-format` | Format of the input data. | `json` | +| `--input-format` | Format of the input documents: `json` or `plain`. Parquet rows are indexed as `json` documents. | `json` | +| `--num-pipelines` | Number of indexing pipelines running in parallel. Values greater than 1 require a `.parquet` input file. Defaults to half the number of CPUs for a `.parquet` input file, 1 otherwise. | | +| `--batch-num-rows` | Number of Parquet rows decoded at once. Only valid with a `.parquet` input file. | | | `--overwrite` | Overwrites pre-existing index. | | | `--transform-script` | VRL program to transform docs before ingesting. | | | `--keep-cache` | Does not clear local cache directory upon completion. | | + +An `--input-path` ending with `.parquet` is bulk loaded by several pipelines without merging (benchmark only, requires the `parquet` Cargo feature). The index must be empty (`--overwrite` clears it); cleanup on failure is best effort. See `docs/internals/parquet-bulk-load.md`. ### tool extract-split Downloads and extracts a split to a directory. diff --git a/quickwit/Cargo.lock b/quickwit/Cargo.lock index 05da65b9ec2..0b17b3d0613 100644 --- a/quickwit/Cargo.lock +++ b/quickwit/Cargo.lock @@ -67,6 +67,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5a15f179cd60c4584b8a8c596927aadc462e27f2ca70c04e0071964a73ba7a75" dependencies = [ "cfg-if", + "const-random", "getrandom 0.3.4", "once_cell", "serde", @@ -183,7 +184,7 @@ version = "1.1.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "40c48f72fd53cd289104fc64099abca73db4166ad86ea0b4341abe65af83dadc" dependencies = [ - "windows-sys 0.60.2", + "windows-sys 0.61.2", ] [[package]] @@ -194,7 +195,7 @@ checksum = "291e6a250ff86cd4a820112fb8898808a366d8f9f58ce16d1f538353ad55747d" dependencies = [ "anstyle", "once_cell_polyfill", - "windows-sys 0.60.2", + "windows-sys 0.61.2", ] [[package]] @@ -224,6 +225,156 @@ version = "0.7.8" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d3fb67a6e08acf24fdeccbac2cb6ac4305825bd1f117462e0e6f2f193345ad56" +[[package]] +name = "arrow-array" +version = "60.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "51dff6e4b9f158864a0aeb6a131ad53858ecbae7a2fc8989307b9b3c6b2f122e" +dependencies = [ + "ahash", + "arrow-buffer", + "arrow-data", + "arrow-schema", + "chrono", + "chrono-tz", + "half", + "hashbrown 0.17.1", + "num-complex", + "num-integer", + "num-traits", +] + +[[package]] +name = "arrow-buffer" +version = "60.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c7235b863533e9db3ab86905b4521251de11226275b436ff3cb2d0010558d1c5" +dependencies = [ + "bytes", + "half", + "num-bigint 0.5.1", + "num-traits", +] + +[[package]] +name = "arrow-cast" +version = "60.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d52cf840f0c34aeeb7e8b05aba72d66878e6c001827f423331f0b260209e7b7e" +dependencies = [ + "arrow-array", + "arrow-buffer", + "arrow-data", + "arrow-ord", + "arrow-schema", + "arrow-select", + "atoi 3.1.0", + "base64 0.23.0", + "chrono", + "half", + "lexical-core", + "num-traits", + "ryu", +] + +[[package]] +name = "arrow-cmp" +version = "60.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "73a0e9adb28d95739fabcdae637261a2c14e618caf58ba1b88857120f9ee2687" +dependencies = [ + "arrow-array", + "arrow-buffer", + "arrow-schema", +] + +[[package]] +name = "arrow-data" +version = "60.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5aa06d5e22786c4966ffc4d2f171d0cb65d95f6bdea108a14e8c335e9faa437f" +dependencies = [ + "arrow-buffer", + "arrow-schema", + "half", + "num-integer", + "num-traits", +] + +[[package]] +name = "arrow-ipc" +version = "60.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "86c523472d22e31f984df7a44fc36e6323d3c9b05ffe53d267ee4c3aa7c128d2" +dependencies = [ + "arrow-array", + "arrow-buffer", + "arrow-data", + "arrow-schema", + "arrow-select", + "flatbuffers", +] + +[[package]] +name = "arrow-json" +version = "60.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "342fc6600a93a48038d853431803c941436b7b51f756763e76cad9784b7e3718" +dependencies = [ + "arrow-array", + "arrow-buffer", + "arrow-cast", + "arrow-ord", + "arrow-schema", + "arrow-select", + "chrono", + "half", + "indexmap 2.14.0", + "itoa", + "lexical-core", + "memchr", + "num-traits", + "ryu", + "serde_core", + "serde_json", + "simdutf8", +] + +[[package]] +name = "arrow-ord" +version = "60.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5a7ac2c6a9cdc782fa9548d3763461eaec8c93347e297fddbaa44807e70a3698" +dependencies = [ + "arrow-array", + "arrow-buffer", + "arrow-cmp", + "arrow-data", + "arrow-schema", + "arrow-select", +] + +[[package]] +name = "arrow-schema" +version = "60.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f85a9a58cd03560526b0066487574aacb9008e413504fa0ee23d61141e96c338" + +[[package]] +name = "arrow-select" +version = "60.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ba3482b7d03bf4a3bf0b195f08a7f2b071595ccd54897e35b975ac63d3104977" +dependencies = [ + "ahash", + "arrow-array", + "arrow-buffer", + "arrow-cmp", + "arrow-data", + "arrow-schema", + "num-traits", +] + [[package]] name = "ascii-canvas" version = "4.0.0" @@ -409,6 +560,15 @@ dependencies = [ "num-traits", ] +[[package]] +name = "atoi" +version = "3.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e7a8bbe9949e43a1edaa043038c68703b04774156afdfb62ba2cef5bf93d67be" +dependencies = [ + "num-traits", +] + [[package]] name = "atomic-waker" version = "1.1.2" @@ -1528,7 +1688,7 @@ dependencies = [ "data-encoding", "half", "nom 7.1.3", - "num-bigint", + "num-bigint 0.4.8", "num-rational", "num-traits", "separator", @@ -1803,7 +1963,7 @@ version = "3.1.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "faf9468729b8cbcea668e36183cb69d317348c2e08e994829fb56ebfdfbaac34" dependencies = [ - "windows-sys 0.48.0", + "windows-sys 0.61.2", ] [[package]] @@ -3029,7 +3189,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb" dependencies = [ "libc", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -3176,6 +3336,16 @@ version = "0.5.7" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1d674e81391d1e1ab681a28d99df07927c6d4aa5b027d7da16ba32d1d21ecd99" +[[package]] +name = "flatbuffers" +version = "25.12.19" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "35f6839d7b3b98adde531effaf34f0c2badc6f4735d26fe74709d8e513a96ef3" +dependencies = [ + "bitflags 2.13.0", + "rustc_version", +] + [[package]] name = "flate2" version = "1.1.9" @@ -3742,6 +3912,7 @@ checksum = "6ea2d84b969582b4b1864a92dc5d27cd2b77b622a8d79306834f1be5ba20d84b" dependencies = [ "cfg-if", "crunchy", + "num-traits", "zerocopy", ] @@ -4139,7 +4310,7 @@ dependencies = [ "libc", "percent-encoding", "pin-project-lite", - "socket2 0.5.10", + "socket2 0.6.4", "tokio", "tower-service", "tracing", @@ -4431,7 +4602,7 @@ checksum = "3640c1c38b8e4e43584d8df18be5fc6b0aa314ce6ebf51b53313d4306cca8e46" dependencies = [ "hermit-abi", "libc", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -4792,6 +4963,63 @@ version = "0.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0c2cdeb66e45e9f36bfad5bbdb4d2384e70936afbee843c6f6543f0c551ebb25" +[[package]] +name = "lexical-core" +version = "1.0.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7d8d125a277f807e55a77304455eb7b1cb52f2b18c143b60e766c120bd64a594" +dependencies = [ + "lexical-parse-float", + "lexical-parse-integer", + "lexical-util", + "lexical-write-float", + "lexical-write-integer", +] + +[[package]] +name = "lexical-parse-float" +version = "1.0.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "52a9f232fbd6f550bc0137dcb5f99ab674071ac2d690ac69704593cb4abbea56" +dependencies = [ + "lexical-parse-integer", + "lexical-util", +] + +[[package]] +name = "lexical-parse-integer" +version = "1.0.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9a7a039f8fb9c19c996cd7b2fcce303c1b2874fe1aca544edc85c4a5f8489b34" +dependencies = [ + "lexical-util", +] + +[[package]] +name = "lexical-util" +version = "1.0.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2604dd126bb14f13fb5d1bd6a66155079cb9fa655b37f875b3a742c705dbed17" + +[[package]] +name = "lexical-write-float" +version = "1.0.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "50c438c87c013188d415fbabbb1dceb44249ab81664efbd31b14ae55dabb6361" +dependencies = [ + "lexical-util", + "lexical-write-integer", +] + +[[package]] +name = "lexical-write-integer" +version = "1.0.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "409851a618475d2d5796377cad353802345cba92c867d9fbcde9cf4eac4e14df" +dependencies = [ + "lexical-util", +] + [[package]] name = "libc" version = "0.2.186" @@ -4947,6 +5175,9 @@ name = "lz4_flex" version = "0.14.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ecbdfe44b1bd960b68170b417450a628c43f7cf56bb3c5317e61cb230ee7f226" +dependencies = [ + "twox-hash", +] [[package]] name = "mach2" @@ -5325,7 +5556,7 @@ version = "0.50.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7957b9740744892f114936ab4a57b3f487491bbeafaf8083688b16841a4240e5" dependencies = [ - "windows-sys 0.60.2", + "windows-sys 0.61.2", ] [[package]] @@ -5334,7 +5565,7 @@ version = "0.4.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "35bd024e8b2ff75562e5f34e7f4905839deb4b22955ef5e73d2fea1b9813cb23" dependencies = [ - "num-bigint", + "num-bigint 0.4.8", "num-complex", "num-integer", "num-iter", @@ -5352,6 +5583,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-bigint-dig" version = "0.8.6" @@ -5424,7 +5665,7 @@ version = "0.4.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f83d14da390562dca69fc84082e73e548e1ad308d24accdedd2720017cb37824" dependencies = [ - "num-bigint", + "num-bigint 0.4.8", "num-integer", "num-traits", ] @@ -5506,7 +5747,7 @@ version = "5.0.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "51e219e79014df21a225b1860a479e2dcd7cbd9130f4defd4bd0e191ea31d67d" dependencies = [ - "base64 0.21.7", + "base64 0.22.1", "chrono", "getrandom 0.2.17", "http 1.4.2", @@ -5956,7 +6197,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7d8fae84b431384b68627d0f9b3b1245fcf9f46f6c0e3dc902e9dce64edd1967" dependencies = [ "libc", - "windows-sys 0.48.0", + "windows-sys 0.61.2", ] [[package]] @@ -6083,6 +6324,34 @@ dependencies = [ "windows-link", ] +[[package]] +name = "parquet" +version = "60.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8af83d2940bc0510f9aef86d865f56fdc6095f87ab115ac885a80b7c5226d3ba" +dependencies = [ + "ahash", + "arrow-array", + "arrow-buffer", + "arrow-data", + "arrow-ipc", + "arrow-schema", + "arrow-select", + "base64 0.23.0", + "bytes", + "chrono", + "half", + "hashbrown 0.17.1", + "lz4_flex 0.14.0", + "num-bigint 0.5.1", + "num-integer", + "num-traits", + "seq-macro", + "snap", + "twox-hash", + "zstd 0.14.0", +] + [[package]] name = "parse-size" version = "1.1.0" @@ -7143,6 +7412,7 @@ name = "quickwit-cli" version = "0.9.1" dependencies = [ "anyhow", + "async-trait", "backtrace", "bytesize", "chrono", @@ -7486,10 +7756,14 @@ version = "0.9.1" dependencies = [ "anyhow", "arc-swap", + "arrow-array", + "arrow-json", + "arrow-schema", "async-compression", "async-trait", "aws-sdk-kinesis", "aws-sdk-sqs", + "base64 0.23.0", "bytes", "bytesize", "criterion", @@ -7506,6 +7780,7 @@ dependencies = [ "mockall", "oneshot 0.2.1", "openssl", + "parquet", "percent-encoding", "proptest", "prost 0.14.4", @@ -8177,7 +8452,7 @@ dependencies = [ "quinn-udp", "rustc-hash", "rustls 0.23.41", - "socket2 0.5.10", + "socket2 0.6.4", "thiserror 2.0.18", "tokio", "tracing", @@ -8216,9 +8491,9 @@ dependencies = [ "cfg_aliases", "libc", "once_cell", - "socket2 0.5.10", + "socket2 0.6.4", "tracing", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -8990,7 +9265,7 @@ dependencies = [ "errno", "libc", "linux-raw-sys 0.4.15", - "windows-sys 0.52.0", + "windows-sys 0.59.0", ] [[package]] @@ -9003,7 +9278,7 @@ dependencies = [ "errno", "libc", "linux-raw-sys 0.12.1", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -9083,7 +9358,7 @@ dependencies = [ "security-framework", "security-framework-sys", "webpki-root-certs", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -9324,7 +9599,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5b55fb86dfd3a2f5f76ea78310a88f96c4ea21a3031f8d212443d56123fd0521" dependencies = [ "libc", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -9343,6 +9618,12 @@ version = "0.4.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f97841a747eef040fcd2e7b3b9a220a7205926e60488e673d9e4926d27772ce5" +[[package]] +name = "seq-macro" +version = "0.3.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1bc711410fbe7399f390ca1c3b60ad0f53f80e95c5eb935e52268a0e2cd49acc" + [[package]] name = "serde" version = "1.0.228" @@ -9731,7 +10012,7 @@ version = "0.6.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0d585997b0ac10be3c5ee635f1bab02d512760d14b7c468801ac8a01d9ae5f1d" dependencies = [ - "num-bigint", + "num-bigint 0.4.8", "num-traits", "thiserror 2.0.18", "time", @@ -9854,7 +10135,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "52d1cfed4120b4d927bf7c0f86d2087a4a7d6027c906d9f9d525a80573b9be51" dependencies = [ "libc", - "windows-sys 0.60.2", + "windows-sys 0.61.2", ] [[package]] @@ -9988,7 +10269,7 @@ version = "0.8.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "aa003f0038df784eb8fecbbac13affe3da23b45194bd57dba231c8f48199c526" dependencies = [ - "atoi", + "atoi 2.0.0", "base64 0.22.1", "bitflags 2.13.0", "byteorder", @@ -10031,7 +10312,7 @@ version = "0.8.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "db58fcd5a53cf07c184b154801ff91347e4c30d17a3562a635ff028ad5deda46" dependencies = [ - "atoi", + "atoi 2.0.0", "base64 0.22.1", "bitflags 2.13.0", "byteorder", @@ -10069,7 +10350,7 @@ version = "0.8.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2d12fe70b2c1b4401038055f90f151b78208de1f9f89a7dbfd41587a10c3eea" dependencies = [ - "atoi", + "atoi 2.0.0", "flume 0.11.1", "futures-channel", "futures-core", @@ -10471,10 +10752,10 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "32497e9a4c7b38532efcdebeef879707aa9f794296a4f0244f6f69e9bc8574bd" dependencies = [ "fastrand 2.4.1", - "getrandom 0.3.4", + "getrandom 0.4.3", "once_cell", "rustix 1.1.4", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -10483,7 +10764,7 @@ version = "1.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d8c27177b12a6399ffc08b98f76f7c9a1f4fe9fc967c784c5a071fa8d93cf7e1" dependencies = [ - "windows-sys 0.60.2", + "windows-sys 0.61.2", ] [[package]] @@ -11891,7 +12172,7 @@ version = "0.1.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22" dependencies = [ - "windows-sys 0.48.0", + "windows-sys 0.61.2", ] [[package]] diff --git a/quickwit/Cargo.toml b/quickwit/Cargo.toml index 99dc769561e..58f454e6f6b 100644 --- a/quickwit/Cargo.toml +++ b/quickwit/Cargo.toml @@ -54,6 +54,9 @@ license = "Apache-2.0" ahash = "0.8" anyhow = "1" arc-swap = "1.9" +arrow-array = { version = "60", default-features = false } +arrow-json = { version = "60", default-features = false } +arrow-schema = { version = "60", default-features = false } assert-json-diff = "2" async-compression = { version = "0.4", features = ["tokio", "gzip"] } async-speed-limit = "0.4" @@ -156,6 +159,7 @@ opentelemetry-otlp = { version = "0.32", default-features = false, features = [ "trace", ] } ouroboros = "0.18" +parquet = { version = "60", default-features = false } percent-encoding = "2.3" pin-project = "1.1" pnet = { version = "0.35", features = ["std"] } diff --git a/quickwit/quickwit-cli/Cargo.toml b/quickwit/quickwit-cli/Cargo.toml index 5257d6f5c85..546040f053d 100644 --- a/quickwit/quickwit-cli/Cargo.toml +++ b/quickwit/quickwit-cli/Cargo.toml @@ -71,12 +71,14 @@ quickwit-transport = { workspace = true } procfs = { workspace = true } [dev-dependencies] +async-trait = { workspace = true } predicates = { workspace = true } reqwest = { workspace = true } quickwit-actors = { workspace = true, features = ["testsuite"] } quickwit-common = { workspace = true, features = ["testsuite"] } quickwit-config = { workspace = true, features = ["testsuite"] } +quickwit-indexing = { workspace = true, features = ["testsuite"] } quickwit-metastore = { workspace = true, features = ["testsuite"] } quickwit-proto = { workspace = true, features = ["testsuite"] } quickwit-storage = { workspace = true, features = ["testsuite"] } @@ -93,6 +95,8 @@ jemalloc-profiled = [ ci-test = [] pprof = ["quickwit-serve/pprof"] openssl-support = ["openssl-probe"] +# Enables `quickwit tool local-ingest` on `.parquet` files (benchmark-only bulk load). +parquet = ["quickwit-indexing/parquet"] # Requires to enable tokio unstable via RUSTFLAGS="--cfg tokio_unstable" tokio-console = [ "dep:console-subscriber", diff --git a/quickwit/quickwit-cli/src/main.rs b/quickwit/quickwit-cli/src/main.rs index 0af72990277..37e105e7789 100644 --- a/quickwit/quickwit-cli/src/main.rs +++ b/quickwit/quickwit-cli/src/main.rs @@ -448,7 +448,10 @@ mod tests { overwrite, vrl_script: Some(vrl_script), clear_cache, + num_pipelines, + batch_num_rows_opt: None, })) if &index_id == "wikipedia" + && num_pipelines.get() == 1 && config_uri == Uri::from_str("file:///config.yaml").unwrap() && vrl_script == ".message = downcase(string!(.message))" && overwrite @@ -457,6 +460,82 @@ mod tests { )); } + #[test] + fn test_parse_local_ingest_parquet_args() { + let app = build_cli().no_binary_name(true); + let matches = app + .try_get_matches_from([ + "tool", + "local-ingest", + "--index", + "wikipedia", + "--config", + "/config.yaml", + "--input-path", + "/data/wikipedia.parquet", + "--num-pipelines", + "8", + "--batch-num-rows", + "4096", + ]) + .unwrap(); + let command = CliCommand::parse_cli_args(matches).unwrap(); + let CliCommand::Tool(ToolCliCommand::LocalIngest(args)) = command else { + panic!("expected local ingest command"); + }; + assert_eq!(args.input_format, SourceInputFormat::Json); + assert_eq!(args.num_pipelines.get(), 8); + assert_eq!(args.batch_num_rows_opt.unwrap().get(), 4096); + } + + #[test] + fn test_parse_local_ingest_args_rejects_multiple_json_pipelines() { + for extra_args in [["--num-pipelines", "2"], ["--batch-num-rows", "1000"]] { + let app = build_cli().no_binary_name(true); + let matches = app + .try_get_matches_from( + [ + "tool", + "local-ingest", + "--index", + "wikipedia", + "--config", + "/config.yaml", + ] + .into_iter() + .chain(extra_args), + ) + .unwrap(); + let error = CliCommand::parse_cli_args(matches).unwrap_err(); + assert!( + error + .to_string() + .contains("requires a `.parquet` input file") + ); + } + } + + #[test] + fn test_parse_local_ingest_args_rejects_plain_parquet() { + let app = build_cli().no_binary_name(true); + let matches = app + .try_get_matches_from([ + "tool", + "local-ingest", + "--index", + "wikipedia", + "--config", + "/config.yaml", + "--input-path", + "/data/wikipedia.parquet", + "--input-format", + "plain", + ]) + .unwrap(); + let error = CliCommand::parse_cli_args(matches).unwrap_err(); + assert!(error.to_string().contains("indexed as `json` documents")); + } + #[test] fn test_parse_search_args() -> anyhow::Result<()> { let app = build_cli().no_binary_name(true); diff --git a/quickwit/quickwit-cli/src/tool.rs b/quickwit/quickwit-cli/src/tool.rs index dfff0b937c2..696b02f8316 100644 --- a/quickwit/quickwit-cli/src/tool.rs +++ b/quickwit/quickwit-cli/src/tool.rs @@ -41,6 +41,7 @@ use quickwit_indexing::docs_clustering::Fingerprinter; use quickwit_indexing::models::{ DetachIndexingPipeline, DetachMergePipeline, IndexingStatistics, SpawnPipeline, }; +use quickwit_indexing::source::{SourceLoader, quickwit_supported_sources}; use quickwit_indexing::{IndexingPipeline, IndexingSplitCache}; use quickwit_ingest::IngesterPool; use quickwit_metastore::IndexMetadataResponseExt; @@ -53,7 +54,7 @@ use quickwit_search::{SearchResponseRest, single_node_search}; use quickwit_serve::{ BodyFormat, SearchRequestQueryString, SortBy, search_request_from_api_request, }; -use quickwit_storage::{BundleStorage, Storage}; +use quickwit_storage::{BundleStorage, Storage, StorageResolver}; use quickwit_transport::ChannelFactory; use thousands::Separable; use tracing::debug; @@ -64,6 +65,11 @@ use crate::{ run_index_checklist, start_actor_runtimes, }; +#[cfg(feature = "parquet")] +mod parquet_ingest; +#[cfg(feature = "parquet")] +mod parquet_report; + pub fn build_tool_command() -> Command { Command::new("tool") .about("Performs utility operations. Requires a node config.") @@ -71,17 +77,21 @@ pub fn build_tool_command() -> Command { .subcommand( Command::new("local-ingest") .display_order(10) - .about("Indexes NDJSON documents locally.") - .long_about("Local ingest indexes locally NDJSON documents from a file or from stdin and uploads splits on the configured storage.") + .about("Indexes NDJSON or Parquet documents locally.") + .long_about("Local ingest indexes locally NDJSON documents from a file or from stdin, or rows of a local `.parquet` file (requires the `parquet` feature), and uploads splits on the configured storage.") .args(&[ arg!(--index "ID of the target index") .display_order(1) .required(true), arg!(--"input-path" "Location of the input file.") .required(false), - arg!(--"input-format" "Format of the input data.") + arg!(--"input-format" "Format of the input documents: `json` or `plain`. Parquet rows are indexed as `json` documents.") .default_value("json") .required(false), + arg!(--"num-pipelines" "Number of indexing pipelines running in parallel. Values greater than 1 require a `.parquet` input file. Defaults to half the number of CPUs for a `.parquet` input file, 1 otherwise.") + .required(false), + arg!(--"batch-num-rows" "Number of Parquet rows decoded at once. Only valid with a `.parquet` input file.") + .required(false), arg!(--overwrite "Overwrites pre-existing index.") .required(false), arg!(--"transform-script"