This is an automated email from the ASF dual-hosted git repository.
github-merge-queue[bot] pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/datafusion.git
The following commit(s) were added to refs/heads/main by this push:
new f34a676302 chore: Update to arrow/parquet 59.1.0 (#23312)
f34a676302 is described below
commit f34a676302e2320526172705503a5ac8222804ea
Author: Andrew Lamb <[email protected]>
AuthorDate: Tue Jul 7 22:35:49 2026 -0400
chore: Update to arrow/parquet 59.1.0 (#23312)
(DRAFT until arrow is updated, I am using this PR to pre-test the
release)
## Which issue does this PR close?
- related to https://github.com/apache/arrow-rs/issues/9878
## Rationale for this change
Update to latest arrow
## What changes are included in this PR?
## Are these changes tested?
Yes by CI
## Are there any user-facing changes?
No API change (this is a minor update of Arrow)
---
Cargo.lock | 88 +++++++++++-----------
Cargo.toml | 18 ++---
datafusion-examples/examples/flight/server.rs | 4 +-
datafusion/common/src/scalar/mod.rs | 3 +-
datafusion/common/src/utils/mod.rs | 4 +-
.../src/min_max/min_max_struct.rs | 2 +-
datafusion/functions-nested/src/array_compact.rs | 4 +-
datafusion/functions-nested/src/arrays_zip.rs | 6 +-
datafusion/functions-nested/src/concat.rs | 10 +--
datafusion/functions-nested/src/extract.rs | 24 +++---
datafusion/functions-nested/src/make_array.rs | 4 +-
datafusion/functions-nested/src/map_extract.rs | 4 +-
datafusion/functions-nested/src/remove.rs | 12 +--
datafusion/functions-nested/src/replace.rs | 30 ++++----
datafusion/functions-nested/src/resize.rs | 13 ++--
datafusion/functions/src/core/getfield.rs | 8 +-
datafusion/physical-expr-common/src/utils.rs | 10 +--
datafusion/proto-common/src/to_proto/mod.rs | 4 +-
datafusion/spark/src/function/array/shuffle.rs | 8 +-
19 files changed, 131 insertions(+), 125 deletions(-)
diff --git a/Cargo.lock b/Cargo.lock
index db8c6e1252..5b43435ec0 100644
--- a/Cargo.lock
+++ b/Cargo.lock
@@ -112,7 +112,7 @@ version = "1.1.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "40c48f72fd53cd289104fc64099abca73db4166ad86ea0b4341abe65af83dadc"
dependencies = [
- "windows-sys 0.61.2",
+ "windows-sys 0.60.2",
]
[[package]]
@@ -123,7 +123,7 @@ checksum =
"291e6a250ff86cd4a820112fb8898808a366d8f9f58ce16d1f538353ad55747d"
dependencies = [
"anstyle",
"once_cell_polyfill",
- "windows-sys 0.61.2",
+ "windows-sys 0.60.2",
]
[[package]]
@@ -164,9 +164,9 @@ checksum =
"7c02d123df017efcdfbd739ef81735b36c5ba83ec3c59c80a9d7ecc718f92e50"
[[package]]
name = "arrow"
-version = "59.0.0"
+version = "59.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "ffaaa3e009861fd829d0a24dd6f115aa8e4634324bb092147d43baafe69ca4a7"
+checksum = "b952ca5a8046ad741b60f142d6eca4aeebcad615694202bc64c5341f23e32c5b"
dependencies = [
"arrow-arith",
"arrow-array",
@@ -187,9 +187,9 @@ dependencies = [
[[package]]
name = "arrow-arith"
-version = "59.0.0"
+version = "59.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "3ac95125e1d71c4a252b5a9c729aef111e80418f08aaa6dbabd1ba66918247fc"
+checksum = "64a13b8d3008c4e9063c597a08f46446fe3fd5789277127672d6c0bdbb43b1ff"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -201,9 +201,9 @@ dependencies = [
[[package]]
name = "arrow-array"
-version = "59.0.0"
+version = "59.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "0c60c79628e9a97cb90d7a0dc3e944f216a902f837d4ecabc14d524bddbbc137"
+checksum = "9486151b2f0785bafc6fa04fc5c99fcb4495455662e58787ea32eaaed33c4192"
dependencies = [
"ahash",
"arrow-buffer",
@@ -220,9 +220,9 @@ dependencies = [
[[package]]
name = "arrow-avro"
-version = "59.0.0"
+version = "59.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "2835d67df2b69bf5de251ee6d289f85650d41b8169dcee0f950ea88747812c32"
+checksum = "2e4f9b23a0d7b613acb59fa20bdbe0f80ffdae6411498378340b3915e45f5b84"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -244,9 +244,9 @@ dependencies = [
[[package]]
name = "arrow-buffer"
-version = "59.0.0"
+version = "59.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "6026f638c400e9878c1b1cc05c3cfd46fbf381285916ab408678701c1df46c1a"
+checksum = "c4776577a87794bfdf0b4e90e2ea12454fa7738ea2823c4be5b9d1851da7b434"
dependencies = [
"bytes",
"half",
@@ -256,9 +256,9 @@ dependencies = [
[[package]]
name = "arrow-cast"
-version = "59.0.0"
+version = "59.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "c82c236c3caf8df5664284f3f1fbe89938852163998c3fdbf37e84ac220445e9"
+checksum = "a9ad451ce4f98710828a455b96991b8f031deb2e67f5fcad6773f017e4a69c3a"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -278,9 +278,9 @@ dependencies = [
[[package]]
name = "arrow-csv"
-version = "59.0.0"
+version = "59.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "12714e5fb7954159af1e26d4e0d37108bcf1a2ad5ee5c5bf02a944d564d588b7"
+checksum = "8aa7bf96d6141a7bcca2eed57c7c9767d2a2175281857b8a7b68308992864784"
dependencies = [
"arrow-array",
"arrow-cast",
@@ -293,9 +293,9 @@ dependencies = [
[[package]]
name = "arrow-data"
-version = "59.0.0"
+version = "59.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "7bd568aa70c4ec5947027b0d5caee94877433b661a0bb9e8ddceeeb5f0c9b1ab"
+checksum = "b38fe43e2e8704360f1464e6e8cc4fc381ef02cc4fb0192afa8df1aaa0115c66"
dependencies = [
"arrow-buffer",
"arrow-schema",
@@ -306,9 +306,9 @@ dependencies = [
[[package]]
name = "arrow-flight"
-version = "59.0.0"
+version = "59.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "68365401e834743d708094927e2ca727a32d639fe900df04b936e07a36701b74"
+checksum = "42115e09dbb694b5955da998912121451c6910b338228cb80a5701370dba43ff"
dependencies = [
"arrow-arith",
"arrow-array",
@@ -334,9 +334,9 @@ dependencies = [
[[package]]
name = "arrow-ipc"
-version = "59.0.0"
+version = "59.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "e57ee4d470eab1a021bc4b63fa2b2c15d572892bf227b0a982d3b755a6c662b5"
+checksum = "29dac499fcbc6ba74ee0324057821d381929a48526a3966bd9dffb44aa06d98c"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -350,9 +350,9 @@ dependencies = [
[[package]]
name = "arrow-json"
-version = "59.0.0"
+version = "59.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "38f47e0e7a284e1f3707a780dc8cd5451b1614e9e398ea2d9ca03c7a2fe9a9ed"
+checksum = "0fe05e916ddc50f4c7a363cd69c0ef5894fcee063517e9a0b8582f0c56746af6"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -375,9 +375,9 @@ dependencies = [
[[package]]
name = "arrow-ord"
-version = "59.0.0"
+version = "59.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "a79cf73ad2eba8686ec2aa9bbf8671208e509025f166afc040cedbd94ffe4983"
+checksum = "0e13dbdc2a9c053c10c7baa6e30faee04a180aa7ce88e471835850ce37abd20b"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -388,9 +388,9 @@ dependencies = [
[[package]]
name = "arrow-row"
-version = "59.0.0"
+version = "59.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "cea0f7d8ed6182f14952761e2c0f989852d5aa334fcbc49f73a9f2247c25b879"
+checksum = "4d5a1f8c733d15260b305683472ee8ad89c62cbd706703ca873b90d051b41592"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -401,9 +401,9 @@ dependencies = [
[[package]]
name = "arrow-schema"
-version = "59.0.0"
+version = "59.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "80b3e786a0dd9103acd583a6fb486dbf2f3268466cc0bd571dcf34cef231c1f1"
+checksum = "d9e4969dc350d571766247143ab36a5187d095d3d3690970408bc630d47c69e5"
dependencies = [
"bitflags",
"serde",
@@ -413,9 +413,9 @@ dependencies = [
[[package]]
name = "arrow-select"
-version = "59.0.0"
+version = "59.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "067a67e0361f6c31f4a7248759f36ca4ca71b187a941ed4d49da1c7d3d4db624"
+checksum = "402770dba90865359d98d1ef92ef16e23d75c0cca9c2c880c8a05468b7743bf9"
dependencies = [
"ahash",
"arrow-array",
@@ -427,9 +427,9 @@ dependencies = [
[[package]]
name = "arrow-string"
-version = "59.0.0"
+version = "59.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "99bc95847f3ff62a2b03d6f8ce2e3e78f01362060549a2a311898dd442f6256d"
+checksum = "a2b0afbb8b9016700938291123df30838b89decc3213dba00852021988b170d3"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -2761,7 +2761,7 @@ dependencies = [
"libc",
"option-ext",
"redox_users",
- "windows-sys 0.61.2",
+ "windows-sys 0.60.2",
]
[[package]]
@@ -2900,7 +2900,7 @@ source =
"registry+https://github.com/rust-lang/crates.io-index"
checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb"
dependencies = [
"libc",
- "windows-sys 0.61.2",
+ "windows-sys 0.52.0",
]
[[package]]
@@ -4172,7 +4172,7 @@ version = "0.50.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "7957b9740744892f114936ab4a57b3f487491bbeafaf8083688b16841a4240e5"
dependencies = [
- "windows-sys 0.61.2",
+ "windows-sys 0.60.2",
]
[[package]]
@@ -4446,9 +4446,9 @@ dependencies = [
[[package]]
name = "parquet"
-version = "59.0.0"
+version = "59.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "970dff83e97d953c827ae8176f6bf4e9f77bf62daacc01ec5df348ec5eacd913"
+checksum = "5302d4da74d6596a1f11f9928767995b53bca657cbeea1e4e8c5074f8a1157dd"
dependencies = [
"ahash",
"arrow-array",
@@ -5330,7 +5330,7 @@ dependencies = [
"errno",
"libc",
"linux-raw-sys",
- "windows-sys 0.61.2",
+ "windows-sys 0.52.0",
]
[[package]]
@@ -5792,7 +5792,7 @@ source =
"registry+https://github.com/rust-lang/crates.io-index"
checksum = "3a766e1110788c36f4fa1c2b71b387a7815aa65f88ce0229841826633d93723e"
dependencies = [
"libc",
- "windows-sys 0.61.2",
+ "windows-sys 0.60.2",
]
[[package]]
@@ -5892,7 +5892,7 @@ dependencies = [
"cfg-if",
"libc",
"psm",
- "windows-sys 0.61.2",
+ "windows-sys 0.60.2",
]
[[package]]
@@ -6061,7 +6061,7 @@ dependencies = [
"getrandom 0.4.2",
"once_cell",
"rustix",
- "windows-sys 0.61.2",
+ "windows-sys 0.52.0",
]
[[package]]
@@ -6980,7 +6980,7 @@ version = "0.1.11"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22"
dependencies = [
- "windows-sys 0.61.2",
+ "windows-sys 0.52.0",
]
[[package]]
diff --git a/Cargo.toml b/Cargo.toml
index 24a4c7a5a8..0bfaad9a68 100644
--- a/Cargo.toml
+++ b/Cargo.toml
@@ -89,30 +89,30 @@ version = "54.0.0"
#
# See for more details: https://github.com/rust-lang/cargo/issues/11329
apache-avro = { version = "0.21", default-features = false }
-arrow = { version = "59.0.0", features = [
+arrow = { version = "59.1.0", features = [
"prettyprint",
"chrono-tz",
] }
-arrow-avro = { version = "59.0.0", default-features = false, features = [
+arrow-avro = { version = "59.1.0", default-features = false, features = [
"deflate",
"snappy",
"zstd",
"bzip2",
"xz",
] }
-arrow-buffer = { version = "59.0.0", default-features = false }
-arrow-data = { version = "59.0.0", default-features = false }
-arrow-flight = { version = "59.0.0", features = [
+arrow-buffer = { version = "59.1.0", default-features = false }
+arrow-data = { version = "59.1.0", default-features = false }
+arrow-flight = { version = "59.1.0", features = [
"flight-sql-experimental",
] }
# Both codecs are required here to make sure that code paths like
# file-spilling have access to all compression codecs.
-arrow-ipc = { version = "59.0.0", default-features = false, features = [
+arrow-ipc = { version = "59.1.0", default-features = false, features = [
"lz4",
"zstd",
] }
-arrow-ord = { version = "59.0.0", default-features = false }
-arrow-schema = { version = "59.0.0", default-features = false }
+arrow-ord = { version = "59.1.0", default-features = false }
+arrow-schema = { version = "59.1.0", default-features = false }
async-trait = "0.1.89"
bigdecimal = "0.4.8"
bytes = "1.11"
@@ -178,7 +178,7 @@ memchr = "2.8.1"
num-traits = { version = "0.2" }
object_store = { version = "0.13.2", default-features = false }
parking_lot = "0.12"
-parquet = { version = "59.0.0", default-features = false, features = [
+parquet = { version = "59.1.0", default-features = false, features = [
"arrow",
"async",
"object_store",
diff --git a/datafusion-examples/examples/flight/server.rs
b/datafusion-examples/examples/flight/server.rs
index b73c81dd7d..ac8908d7c8 100644
--- a/datafusion-examples/examples/flight/server.rs
+++ b/datafusion-examples/examples/flight/server.rs
@@ -19,7 +19,7 @@
use std::sync::Arc;
-use arrow::ipc::writer::{CompressionContext, DictionaryTracker,
IpcDataGenerator};
+use arrow::ipc::writer::{DictionaryTracker, IpcDataGenerator, IpcWriteContext};
use arrow_flight::{
Action, ActionType, Criteria, Empty, FlightData, FlightDescriptor,
FlightInfo,
HandshakeRequest, HandshakeResponse, PutResult, SchemaResult, Ticket,
@@ -112,7 +112,7 @@ impl FlightService for FlightServiceImpl {
// add an initial FlightData message that sends schema
let options = arrow::ipc::writer::IpcWriteOptions::default();
- let mut compression_context = CompressionContext::default();
+ let mut compression_context = IpcWriteContext::default();
let schema_flight_data = SchemaAsIpc::new(&schema, &options);
let mut flights = vec![FlightData::from(schema_flight_data)];
diff --git a/datafusion/common/src/scalar/mod.rs
b/datafusion/common/src/scalar/mod.rs
index bba7f77b89..3d2e8d08ac 100644
--- a/datafusion/common/src/scalar/mod.rs
+++ b/datafusion/common/src/scalar/mod.rs
@@ -5265,7 +5265,8 @@ impl ScalarValue {
/// as necessary.
pub fn copy_array_data(src_data: &ArrayData) -> ArrayData {
let mut copy = MutableArrayData::new(vec![&src_data], true,
src_data.len());
- copy.extend(0, 0, src_data.len());
+ copy.try_extend(0, 0, src_data.len())
+ .expect("copy_array_data failed due to offset overflow");
copy.freeze()
}
diff --git a/datafusion/common/src/utils/mod.rs
b/datafusion/common/src/utils/mod.rs
index 041cb1aa0b..30a6c45dfc 100644
--- a/datafusion/common/src/utils/mod.rs
+++ b/datafusion/common/src/utils/mod.rs
@@ -1288,11 +1288,11 @@ fn truncate_list_nulls<O: OffsetSizeTrait>(
let (valid_or_empty, _nulls) = valid_or_empty.into_parts();
for (start, end) in valid_or_empty.set_slices() {
- mutable_array_data.extend(
+ mutable_array_data.try_extend(
0,
offsets[start].as_usize(),
offsets[end].as_usize(),
- );
+ )?;
}
let lengths = std::iter::zip(offsets.lengths(), nulls)
diff --git a/datafusion/functions-aggregate/src/min_max/min_max_struct.rs
b/datafusion/functions-aggregate/src/min_max/min_max_struct.rs
index 10580ac18d..15df0f1d44 100644
--- a/datafusion/functions-aggregate/src/min_max/min_max_struct.rs
+++ b/datafusion/functions-aggregate/src/min_max/min_max_struct.rs
@@ -118,7 +118,7 @@ impl GroupsAccumulator for MinMaxStructAccumulator {
let mut copy = MutableArrayData::new(min_maxes_refs, true,
min_maxes_data.len());
for (i, item) in min_maxes_data.iter().enumerate() {
- copy.extend(i, 0, item.len());
+ copy.try_extend(i, 0, item.len())?;
}
let result = copy.freeze();
assert_eq!(&self.inner.data_type, result.data_type());
diff --git a/datafusion/functions-nested/src/array_compact.rs
b/datafusion/functions-nested/src/array_compact.rs
index adea9efef2..4222d6264b 100644
--- a/datafusion/functions-nested/src/array_compact.rs
+++ b/datafusion/functions-nested/src/array_compact.rs
@@ -173,7 +173,7 @@ fn compact_list<O: OffsetSizeTrait>(
if values_nulls.is_null(i) {
// Null breaks the current batch — flush it
if let Some(bs) = batch_start {
- mutable.extend(0, bs, i);
+ mutable.try_extend(0, bs, i)?;
batch_start = None;
}
} else if batch_start.is_none() {
@@ -182,7 +182,7 @@ fn compact_list<O: OffsetSizeTrait>(
}
// Flush any remaining batch after the loop
if let Some(bs) = batch_start {
- mutable.extend(0, bs, end);
+ mutable.try_extend(0, bs, end)?;
}
offsets.push(offsets[row_index] + O::usize_as(kept));
diff --git a/datafusion/functions-nested/src/arrays_zip.rs
b/datafusion/functions-nested/src/arrays_zip.rs
index 76b1b589f4..c574821707 100644
--- a/datafusion/functions-nested/src/arrays_zip.rs
+++ b/datafusion/functions-nested/src/arrays_zip.rs
@@ -280,15 +280,15 @@ fn arrays_zip_inner(args: &[ArrayRef]) ->
Result<ArrayRef> {
let end = v.offsets[row_idx + 1];
let len = end - start;
let builder = builders[col_idx].as_mut().unwrap();
- builder.extend(0, start, end);
+ builder.try_extend(0, start, end)?;
if len < max_len {
- builder.extend_nulls(max_len - len);
+ builder.try_extend_nulls(max_len - len)?;
}
}
_ => {
// Null list entry or None (Null-typed) arg — all nulls.
if let Some(builder) = builders[col_idx].as_mut() {
- builder.extend_nulls(max_len);
+ builder.try_extend_nulls(max_len)?;
}
}
}
diff --git a/datafusion/functions-nested/src/concat.rs
b/datafusion/functions-nested/src/concat.rs
index 8d06140889..5dc437b3c2 100644
--- a/datafusion/functions-nested/src/concat.rs
+++ b/datafusion/functions-nested/src/concat.rs
@@ -432,7 +432,7 @@ fn concat_internal<O: OffsetSizeTrait>(args: &[ArrayRef])
-> Result<ArrayRef> {
let start = list_array.offsets()[row_idx].to_usize().unwrap();
let end = list_array.offsets()[row_idx + 1].to_usize().unwrap();
if start < end {
- mutable.extend(arr_idx, start, end);
+ mutable.try_extend(arr_idx, start, end)?;
}
}
offsets.push(O::usize_as(mutable.len()));
@@ -553,11 +553,11 @@ where
let start = offset_window[0].to_usize().unwrap();
let end = offset_window[1].to_usize().unwrap();
if is_append {
- mutable.extend(values_index, start, end);
- mutable.extend(element_index, row_index, row_index + 1);
+ mutable.try_extend(values_index, start, end)?;
+ mutable.try_extend(element_index, row_index, row_index + 1)?;
} else {
- mutable.extend(element_index, row_index, row_index + 1);
- mutable.extend(values_index, start, end);
+ mutable.try_extend(element_index, row_index, row_index + 1)?;
+ mutable.try_extend(values_index, start, end)?;
}
offsets.push(offsets[row_index] + O::usize_as(end - start + 1));
}
diff --git a/datafusion/functions-nested/src/extract.rs
b/datafusion/functions-nested/src/extract.rs
index 202a76bd0b..900b408bff 100644
--- a/datafusion/functions-nested/src/extract.rs
+++ b/datafusion/functions-nested/src/extract.rs
@@ -258,7 +258,7 @@ where
// array or index is null
if array.is_null(row_index) || indexes.is_null(row_index) {
- mutable.extend_nulls(1);
+ mutable.try_extend_nulls(1)?;
continue;
}
@@ -266,10 +266,10 @@ where
if let Some(index) = index {
let start = start.as_usize() + index.as_usize();
- mutable.extend(0, start, start + 1_usize);
+ mutable.try_extend(0, start, start + 1_usize)?;
} else {
// Index out of bounds
- mutable.extend_nulls(1);
+ mutable.try_extend_nulls(1)?;
}
}
@@ -639,7 +639,7 @@ where
let len = end - start;
if nulls.as_ref().is_some_and(|n| n.is_null(row_index)) {
- mutable.extend_nulls(1);
+ mutable.try_extend_nulls(1)?;
offsets.push(offsets[row_index] + O::usize_as(1));
continue;
}
@@ -665,14 +665,14 @@ where
} => {
let start_index = (start + rel_start).to_usize().unwrap();
let end_index = (start + rel_start +
slice_len).to_usize().unwrap();
- mutable.extend(0, start_index, end_index);
+ mutable.try_extend(0, start_index, end_index)?;
offsets.push(offsets[row_index] + slice_len);
}
SlicePlan::Indices(indices) => {
let count = indices.len();
for rel_index in indices {
let absolute_index = (start +
rel_index).to_usize().unwrap();
- mutable.extend(0, absolute_index, absolute_index + 1);
+ mutable.try_extend(0, absolute_index, absolute_index + 1)?;
}
offsets.push(offsets[row_index] + O::usize_as(count));
}
@@ -754,7 +754,7 @@ where
} => {
let start_index = (start + rel_start).to_usize().unwrap();
let end_index = (start + rel_start +
slice_len).to_usize().unwrap();
- mutable.extend(0, start_index, end_index);
+ mutable.try_extend(0, start_index, end_index)?;
offsets.push(current_offset);
sizes.push(slice_len);
current_offset += slice_len;
@@ -763,7 +763,7 @@ where
let count = indices.len();
for rel_index in indices {
let absolute_index = (start +
rel_index).to_usize().unwrap();
- mutable.extend(0, absolute_index, absolute_index + 1);
+ mutable.try_extend(0, absolute_index, absolute_index + 1)?;
}
let length = O::usize_as(count);
offsets.push(current_offset);
@@ -1065,7 +1065,7 @@ where
// array is null
if array.is_null(row_index) {
- mutable.extend_nulls(1);
+ mutable.try_extend_nulls(1)?;
continue;
}
@@ -1077,16 +1077,16 @@ where
row_nulls_buffer.valid_indices().next()
{
let index = start.as_usize() + first_non_null_index;
- mutable.extend(0, index, index + 1)
+ mutable.try_extend(0, index, index + 1)?;
} else {
// all the elements in the array are null
- mutable.extend_nulls(1);
+ mutable.try_extend_nulls(1)?;
}
}
None => {
// no nulls are present in the array so take the first element
let index = start.as_usize();
- mutable.extend(0, index, index + 1);
+ mutable.try_extend(0, index, index + 1)?;
}
}
}
diff --git a/datafusion/functions-nested/src/make_array.rs
b/datafusion/functions-nested/src/make_array.rs
index 32af5df2c6..6f083ab700 100644
--- a/datafusion/functions-nested/src/make_array.rs
+++ b/datafusion/functions-nested/src/make_array.rs
@@ -224,9 +224,9 @@ pub fn array_array<O: OffsetSizeTrait>(
&& !arg.is_null(row_idx)
&& arg.is_valid(row_idx)
{
- mutable.extend(arr_idx, row_idx, row_idx + 1);
+ mutable.try_extend(arr_idx, row_idx, row_idx + 1)?;
} else {
- mutable.extend_nulls(1);
+ mutable.try_extend_nulls(1)?;
}
}
offsets.push(O::usize_as(mutable.len()));
diff --git a/datafusion/functions-nested/src/map_extract.rs
b/datafusion/functions-nested/src/map_extract.rs
index aab0d013a4..69c5088fc9 100644
--- a/datafusion/functions-nested/src/map_extract.rs
+++ b/datafusion/functions-nested/src/map_extract.rs
@@ -161,10 +161,10 @@ fn general_map_extract_inner(
match value_index {
Some(index) => {
- mutable.extend(0, start + index, start + index + 1);
+ mutable.try_extend(0, start + index, start + index + 1)?;
}
None => {
- mutable.extend_nulls(1);
+ mutable.try_extend_nulls(1)?;
}
}
offsets.push(offsets[row_index] + 1);
diff --git a/datafusion/functions-nested/src/remove.rs
b/datafusion/functions-nested/src/remove.rs
index 111147659a..491e823c9a 100644
--- a/datafusion/functions-nested/src/remove.rs
+++ b/datafusion/functions-nested/src/remove.rs
@@ -508,7 +508,7 @@ fn general_remove<OffsetSize: OffsetSizeTrait>(
// Fast path: no elements to remove, copy entire row
if num_to_remove == 0 {
- mutable.extend(0, start, end);
+ mutable.try_extend(0, start, end)?;
offsets.push(offsets[row_index] + OffsetSize::usize_as(end -
start));
valid.append_non_null();
continue;
@@ -524,7 +524,7 @@ fn general_remove<OffsetSize: OffsetSizeTrait>(
if keep == Some(false) && removed < max_removals {
// Flush pending batch before skipping this element
if let Some(bs) = pending_batch_to_retain {
- mutable.extend(0, start + bs, start + i);
+ mutable.try_extend(0, start + bs, start + i)?;
copied += i - bs;
pending_batch_to_retain = None;
}
@@ -536,7 +536,7 @@ fn general_remove<OffsetSize: OffsetSizeTrait>(
// Flush remaining batch
if let Some(bs) = pending_batch_to_retain {
- mutable.extend(0, start + bs, start + eq_array.len());
+ mutable.try_extend(0, start + bs, start + eq_array.len())?;
copied += eq_array.len() - bs;
}
@@ -610,7 +610,7 @@ fn general_remove_with_scalar<OffsetSize: OffsetSizeTrait>(
let num_to_remove = row_remove_bits.count_set_bits();
if num_to_remove == 0 {
- mutable.extend(0, start, end);
+ mutable.try_extend(0, start, end)?;
offsets.push(offsets[row_index] + OffsetSize::usize_as(row_len));
continue;
}
@@ -626,7 +626,7 @@ fn general_remove_with_scalar<OffsetSize: OffsetSizeTrait>(
for remove_pos in row_remove_bits.set_indices() {
let abs_pos = start + remove_pos;
if abs_pos > prev_end {
- mutable.extend(0, prev_end, abs_pos);
+ mutable.try_extend(0, prev_end, abs_pos)?;
copied += abs_pos - prev_end;
}
prev_end = abs_pos + 1;
@@ -637,7 +637,7 @@ fn general_remove_with_scalar<OffsetSize: OffsetSizeTrait>(
}
// Copy the remaining tail after the last removal
if prev_end < end {
- mutable.extend(0, prev_end, end);
+ mutable.try_extend(0, prev_end, end)?;
copied += end - prev_end;
}
diff --git a/datafusion/functions-nested/src/replace.rs
b/datafusion/functions-nested/src/replace.rs
index 28808a05db..71d6f57815 100644
--- a/datafusion/functions-nested/src/replace.rs
+++ b/datafusion/functions-nested/src/replace.rs
@@ -441,11 +441,11 @@ fn general_replace<O: OffsetSizeTrait>(
// All elements are false, no need to replace, just copy original data
if n <= 0 || !eq_array.has_true() {
- mutable.extend(
+ mutable.try_extend(
original_idx.to_usize().unwrap(),
start.to_usize().unwrap(),
end.to_usize().unwrap(),
- );
+ )?;
offsets.push(offsets[row_index] + (end - start));
valid.append_non_null();
continue;
@@ -457,21 +457,25 @@ fn general_replace<O: OffsetSizeTrait>(
if to_replace == Some(true) && counter < n {
// Flush any pending retain run before emitting the
replacement.
if let Some(rs) = pending_retain.take() {
- mutable.extend(
+ mutable.try_extend(
original_idx.to_usize().unwrap(),
(start + rs).to_usize().unwrap(),
(start + i).to_usize().unwrap(),
- );
+ )?;
}
- mutable.extend(replace_idx.to_usize().unwrap(), row_index,
row_index + 1);
+ mutable.try_extend(
+ replace_idx.to_usize().unwrap(),
+ row_index,
+ row_index + 1,
+ )?;
counter += 1;
if counter == n {
// copy original data for any matches past n
- mutable.extend(
+ mutable.try_extend(
original_idx.to_usize().unwrap(),
(start + i).to_usize().unwrap() + 1,
end.to_usize().unwrap(),
- );
+ )?;
break;
}
} else if pending_retain.is_none() {
@@ -484,11 +488,11 @@ fn general_replace<O: OffsetSizeTrait>(
if counter < n
&& let Some(rs) = pending_retain
{
- mutable.extend(
+ mutable.try_extend(
original_idx.to_usize().unwrap(),
(start + rs).to_usize().unwrap(),
end.to_usize().unwrap(),
- );
+ )?;
}
offsets.push(offsets[row_index] + (end - start));
@@ -564,7 +568,7 @@ fn general_replace_with_scalar<O: OffsetSizeTrait>(
.take(max_replacements as usize)
.peekable();
if match_positions.peek().is_none() {
- mutable.extend(0, start, end);
+ mutable.try_extend(0, start, end)?;
offsets.push(offsets[row_index] + O::usize_as(row_len));
continue;
}
@@ -576,16 +580,16 @@ fn general_replace_with_scalar<O: OffsetSizeTrait>(
for match_pos in match_positions {
// Retain elements before this match.
if match_pos > prev_end {
- mutable.extend(0, start + prev_end, start + match_pos);
+ mutable.try_extend(0, start + prev_end, start + match_pos)?;
}
// Emit the replacement element.
- mutable.extend(1, 0, 1);
+ mutable.try_extend(1, 0, 1)?;
prev_end = match_pos + 1;
}
// Copy remaining elements after the last replacement.
if prev_end < row_len {
- mutable.extend(0, start + prev_end, end);
+ mutable.try_extend(0, start + prev_end, end)?;
}
offsets.push(offsets[row_index] + O::usize_as(row_len));
diff --git a/datafusion/functions-nested/src/resize.rs
b/datafusion/functions-nested/src/resize.rs
index e4fd8421fe..832ddbdc0a 100644
--- a/datafusion/functions-nested/src/resize.rs
+++ b/datafusion/functions-nested/src/resize.rs
@@ -264,7 +264,7 @@ fn general_list_resize<O: OffsetSizeTrait + TryInto<i64>>(
&original_data,
&default_value_data,
output_values_len,
- |mutable, _, extra_count| mutable.extend(1, 0, extra_count),
+ |mutable, _, extra_count| Ok(mutable.try_extend(1, 0,
extra_count)?),
)
} else {
// Slow path: rows may need different fill values, so append from the
@@ -286,8 +286,9 @@ fn general_list_resize<O: OffsetSizeTrait + TryInto<i64>>(
output_values_len,
|mutable, row_index, extra_count| {
for _ in 0..extra_count {
- mutable.extend(1, row_index, row_index + 1);
+ mutable.try_extend(1, row_index, row_index + 1)?;
}
+ Ok(())
},
)
}
@@ -304,7 +305,7 @@ fn build_resized_list<O, F>(
) -> Result<ArrayRef>
where
O: OffsetSizeTrait + TryInto<i64>,
- F: FnMut(&mut MutableArrayData, usize, usize),
+ F: FnMut(&mut MutableArrayData, usize, usize) -> Result<()>,
{
let capacity = Capacities::Array(output_values_len);
let mut offsets = vec![O::usize_as(0)];
@@ -331,11 +332,11 @@ where
if start + count > offset_window[1] {
let extra_count = (start + count -
offset_window[1]).to_usize().unwrap();
let end = offset_window[1];
- mutable.extend(0, start.to_usize().unwrap(),
end.to_usize().unwrap());
- append_fill_values(&mut mutable, row_index, extra_count);
+ mutable.try_extend(0, start.to_usize().unwrap(),
end.to_usize().unwrap())?;
+ append_fill_values(&mut mutable, row_index, extra_count)?;
} else {
let end = start + count;
- mutable.extend(0, start.to_usize().unwrap(),
end.to_usize().unwrap());
+ mutable.try_extend(0, start.to_usize().unwrap(),
end.to_usize().unwrap())?;
};
offsets.push(offsets[row_index] + count);
}
diff --git a/datafusion/functions/src/core/getfield.rs
b/datafusion/functions/src/core/getfield.rs
index 93a4cddef4..70fc8bb0ea 100644
--- a/datafusion/functions/src/core/getfield.rs
+++ b/datafusion/functions/src/core/getfield.rs
@@ -140,11 +140,11 @@ fn process_map_array(
.find(|(_, t)| t.unwrap());
if maybe_matched.is_none() {
- mutable.extend_nulls(1);
+ mutable.try_extend_nulls(1)?;
continue;
}
let (match_offset, _) = maybe_matched.unwrap();
- mutable.extend(0, start + match_offset, start + match_offset + 1);
+ mutable.try_extend(0, start + match_offset, start + match_offset + 1)?;
}
let data = mutable.freeze();
@@ -177,14 +177,14 @@ fn process_map_with_nested_key(
let mut found_match = false;
for i in start..end {
if comparator(i, 0).is_eq() {
- mutable.extend(0, i, i + 1);
+ mutable.try_extend(0, i, i + 1)?;
found_match = true;
break;
}
}
if !found_match {
- mutable.extend_nulls(1);
+ mutable.try_extend_nulls(1)?;
}
}
diff --git a/datafusion/physical-expr-common/src/utils.rs
b/datafusion/physical-expr-common/src/utils.rs
index 117da23df2..5dadcdcabb 100644
--- a/datafusion/physical-expr-common/src/utils.rs
+++ b/datafusion/physical-expr-common/src/utils.rs
@@ -370,21 +370,21 @@ fn scatter_fallback(
let mut true_pos = 0;
let mask_array = BooleanArray::new(mask.clone(), None);
- SlicesIterator::new(&mask_array).for_each(|(start, end)| {
+ for (start, end) in SlicesIterator::new(&mask_array) {
// the gap needs to be filled with nulls
if start > filled {
- mutable.extend_nulls(start - filled);
+ mutable.try_extend_nulls(start - filled)?;
}
// fill with truthy values
let len = end - start;
- mutable.extend(0, true_pos, true_pos + len);
+ mutable.try_extend(0, true_pos, true_pos + len)?;
true_pos += len;
filled = end;
- });
+ }
// the remaining part is falsy
if filled < output_len {
- mutable.extend_nulls(output_len - filled);
+ mutable.try_extend_nulls(output_len - filled)?;
}
let data = mutable.freeze();
diff --git a/datafusion/proto-common/src/to_proto/mod.rs
b/datafusion/proto-common/src/to_proto/mod.rs
index a6fa13ca74..d2e1ca50c8 100644
--- a/datafusion/proto-common/src/to_proto/mod.rs
+++ b/datafusion/proto-common/src/to_proto/mod.rs
@@ -29,7 +29,7 @@ use arrow::datatypes::{
SchemaRef, TimeUnit, UnionMode,
};
use arrow::ipc::writer::{
- CompressionContext, DictionaryTracker, IpcDataGenerator, IpcWriteOptions,
+ DictionaryTracker, IpcDataGenerator, IpcWriteContext, IpcWriteOptions,
};
use datafusion_common::parsers::CsvQuoteStyle;
use datafusion_common::{
@@ -1112,7 +1112,7 @@ fn encode_scalar_nested_value(
&mut dict_tracker,
&write_options,
);
- let mut compression_context = CompressionContext::default();
+ let mut compression_context = IpcWriteContext::default();
let (encoded_dictionaries, encoded_message) = ipc_gen
.encode(
&batch,
diff --git a/datafusion/spark/src/function/array/shuffle.rs
b/datafusion/spark/src/function/array/shuffle.rs
index 031dd17177..2673c9155f 100644
--- a/datafusion/spark/src/function/array/shuffle.rs
+++ b/datafusion/spark/src/function/array/shuffle.rs
@@ -185,7 +185,7 @@ fn general_array_shuffle<O: OffsetSizeTrait>(
if array.is_null(row_index) {
nulls.push(false);
offsets.push(offsets[row_index] + O::one());
- mutable.extend(0, 0, 1);
+ mutable.try_extend(0, 0, 1)?;
continue;
}
nulls.push(true);
@@ -200,7 +200,7 @@ fn general_array_shuffle<O: OffsetSizeTrait>(
// Add shuffled elements
for &index in &indices {
- mutable.extend(0, index, index + 1);
+ mutable.try_extend(0, index, index + 1)?;
}
offsets.push(offsets[row_index] + O::usize_as(length));
@@ -239,7 +239,7 @@ fn fixed_size_array_shuffle(
// skip the null value
if array.is_null(row_index) {
nulls.push(false);
- mutable.extend(0, 0, value_length);
+ mutable.try_extend(0, 0, value_length)?;
continue;
}
nulls.push(true);
@@ -253,7 +253,7 @@ fn fixed_size_array_shuffle(
// Add shuffled elements
for &index in &indices {
- mutable.extend(0, index, index + 1);
+ mutable.try_extend(0, index, index + 1)?;
}
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]