Skip to content

Commit dd0464e

Browse files
authored
Merge branch 'main' into alamb/update_arrow_58
2 parents bf17d4a + b9328b9 commit dd0464e

17 files changed

Lines changed: 823 additions & 126 deletions

File tree

Cargo.lock

Lines changed: 4 additions & 4 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

Cargo.toml

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -183,7 +183,7 @@ regex = "1.12"
183183
rstest = "0.26.1"
184184
serde_json = "1"
185185
sha2 = "^0.10.9"
186-
sqlparser = { version = "0.60.0", default-features = false, features = ["std", "visitor"] }
186+
sqlparser = { version = "0.61.0", default-features = false, features = ["std", "visitor"] }
187187
strum = "0.27.2"
188188
strum_macros = "0.27.2"
189189
tempfile = "3"

datafusion-cli/tests/cli_integration.rs

Lines changed: 40 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -44,7 +44,7 @@ fn make_settings() -> Settings {
4444
settings
4545
}
4646

47-
async fn setup_minio_container() -> ContainerAsync<minio::MinIO> {
47+
async fn setup_minio_container() -> Result<ContainerAsync<minio::MinIO>, String> {
4848
const MINIO_ROOT_USER: &str = "TEST-DataFusionLogin";
4949
const MINIO_ROOT_PASSWORD: &str = "TEST-DataFusionPassword";
5050

@@ -99,27 +99,23 @@ async fn setup_minio_container() -> ContainerAsync<minio::MinIO> {
9999
let stdout = container.stdout_to_vec().await.unwrap_or_default();
100100
let stderr = container.stderr_to_vec().await.unwrap_or_default();
101101

102-
panic!(
102+
return Err(format!(
103103
"Failed to execute command: {}\nError: {}\nStdout: {:?}\nStderr: {:?}",
104104
cmd_ref,
105105
e,
106106
String::from_utf8_lossy(&stdout),
107107
String::from_utf8_lossy(&stderr)
108-
);
108+
));
109109
}
110110
}
111111

112-
container
112+
Ok(container)
113113
}
114114

115-
Err(TestcontainersError::Client(e)) => {
116-
panic!(
117-
"Failed to start MinIO container. Ensure Docker is running and accessible: {e}"
118-
);
119-
}
120-
Err(e) => {
121-
panic!("Failed to start MinIO container: {e}");
122-
}
115+
Err(TestcontainersError::Client(e)) => Err(format!(
116+
"Failed to start MinIO container. Ensure Docker is running and accessible: {e}"
117+
)),
118+
Err(e) => Err(format!("Failed to start MinIO container: {e}")),
123119
}
124120
}
125121

@@ -253,7 +249,14 @@ async fn test_cli() {
253249
return;
254250
}
255251

256-
let container = setup_minio_container().await;
252+
let container = match setup_minio_container().await {
253+
Ok(c) => c,
254+
Err(e) if e.contains("toomanyrequests") => {
255+
eprintln!("Skipping test: Docker pull rate limit reached: {e}");
256+
return;
257+
}
258+
e @ Err(_) => e.unwrap(),
259+
};
257260

258261
let settings = make_settings();
259262
let _bound = settings.bind_to_scope();
@@ -286,7 +289,14 @@ async fn test_aws_options() {
286289
let settings = make_settings();
287290
let _bound = settings.bind_to_scope();
288291

289-
let container = setup_minio_container().await;
292+
let container = match setup_minio_container().await {
293+
Ok(c) => c,
294+
Err(e) if e.contains("toomanyrequests") => {
295+
eprintln!("Skipping test: Docker pull rate limit reached: {e}");
296+
return;
297+
}
298+
e @ Err(_) => e.unwrap(),
299+
};
290300
let port = container.get_host_port_ipv4(9000).await.unwrap();
291301

292302
let input = format!(
@@ -377,7 +387,14 @@ async fn test_s3_url_fallback() {
377387
return;
378388
}
379389

380-
let container = setup_minio_container().await;
390+
let container = match setup_minio_container().await {
391+
Ok(c) => c,
392+
Err(e) if e.contains("toomanyrequests") => {
393+
eprintln!("Skipping test: Docker pull rate limit reached: {e}");
394+
return;
395+
}
396+
e @ Err(_) => e.unwrap(),
397+
};
381398

382399
let mut settings = make_settings();
383400
settings.set_snapshot_suffix("s3_url_fallback");
@@ -407,8 +424,14 @@ async fn test_object_store_profiling() {
407424
return;
408425
}
409426

410-
let container = setup_minio_container().await;
411-
427+
let container = match setup_minio_container().await {
428+
Ok(c) => c,
429+
Err(e) if e.contains("toomanyrequests") => {
430+
eprintln!("Skipping test: Docker pull rate limit reached: {e}");
431+
return;
432+
}
433+
e @ Err(_) => e.unwrap(),
434+
};
412435
let mut settings = make_settings();
413436

414437
// as the object store profiling contains timestamps and durations, we must

datafusion/core/benches/aggregate_query_sql.rs

Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -251,6 +251,39 @@ fn criterion_benchmark(c: &mut Criterion) {
251251
)
252252
})
253253
});
254+
255+
c.bench_function("array_agg_query_group_by_few_groups", |b| {
256+
b.iter(|| {
257+
query(
258+
ctx.clone(),
259+
&rt,
260+
"SELECT u64_narrow, array_agg(f64) \
261+
FROM t GROUP BY u64_narrow",
262+
)
263+
})
264+
});
265+
266+
c.bench_function("array_agg_query_group_by_mid_groups", |b| {
267+
b.iter(|| {
268+
query(
269+
ctx.clone(),
270+
&rt,
271+
"SELECT u64_mid, array_agg(f64) \
272+
FROM t GROUP BY u64_mid",
273+
)
274+
})
275+
});
276+
277+
c.bench_function("array_agg_query_group_by_many_groups", |b| {
278+
b.iter(|| {
279+
query(
280+
ctx.clone(),
281+
&rt,
282+
"SELECT u64_wide, array_agg(f64) \
283+
FROM t GROUP BY u64_wide",
284+
)
285+
})
286+
});
254287
}
255288

256289
criterion_group!(benches, criterion_benchmark);

datafusion/core/benches/data_utils/mod.rs

Lines changed: 47 additions & 54 deletions
Original file line numberDiff line numberDiff line change
@@ -45,7 +45,7 @@ pub fn create_table_provider(
4545
) -> Result<Arc<MemTable>> {
4646
let schema = Arc::new(create_schema());
4747
let partitions =
48-
create_record_batches(schema.clone(), array_len, partitions_len, batch_size);
48+
create_record_batches(&schema, array_len, partitions_len, batch_size);
4949
// declare a table in memory. In spark API, this corresponds to createDataFrame(...).
5050
MemTable::try_new(schema, partitions).map(Arc::new)
5151
}
@@ -56,21 +56,19 @@ pub fn create_schema() -> Schema {
5656
Field::new("utf8", DataType::Utf8, false),
5757
Field::new("f32", DataType::Float32, false),
5858
Field::new("f64", DataType::Float64, true),
59-
// This field will contain integers randomly selected from a large
60-
// range of values, i.e. [0, u64::MAX], such that there are none (or
61-
// very few) repeated values.
62-
Field::new("u64_wide", DataType::UInt64, true),
63-
// This field will contain integers randomly selected from a narrow
64-
// range of values such that there are a few distinct values, but they
65-
// are repeated often.
59+
// Integers randomly selected from a wide range of values, i.e. [0,
60+
// u64::MAX], such that there are ~no repeated values.
61+
Field::new("u64_wide", DataType::UInt64, false),
62+
// Integers randomly selected from a mid-range of values [0, 1000),
63+
// providing ~1000 distinct groups.
64+
Field::new("u64_mid", DataType::UInt64, false),
65+
// Integers randomly selected from a narrow range of values such that
66+
// there are a few distinct values, but they are repeated often.
6667
Field::new("u64_narrow", DataType::UInt64, false),
6768
])
6869
}
6970

70-
fn create_data(size: usize, null_density: f64) -> Vec<Option<f64>> {
71-
// use random numbers to avoid spurious compiler optimizations wrt to branching
72-
let mut rng = StdRng::seed_from_u64(42);
73-
71+
fn create_data(rng: &mut StdRng, size: usize, null_density: f64) -> Vec<Option<f64>> {
7472
(0..size)
7573
.map(|_| {
7674
if rng.random::<f64>() > null_density {
@@ -82,56 +80,43 @@ fn create_data(size: usize, null_density: f64) -> Vec<Option<f64>> {
8280
.collect()
8381
}
8482

85-
fn create_integer_data(
86-
rng: &mut StdRng,
87-
size: usize,
88-
value_density: f64,
89-
) -> Vec<Option<u64>> {
90-
(0..size)
91-
.map(|_| {
92-
if rng.random::<f64>() > value_density {
93-
None
94-
} else {
95-
Some(rng.random::<u64>())
96-
}
97-
})
98-
.collect()
99-
}
100-
10183
fn create_record_batch(
10284
schema: SchemaRef,
10385
rng: &mut StdRng,
10486
batch_size: usize,
105-
i: usize,
87+
batch_index: usize,
10688
) -> RecordBatch {
107-
// the 4 here is the number of different keys.
108-
// a higher number increase sparseness
109-
let vs = [0, 1, 2, 3];
110-
let keys: Vec<String> = (0..batch_size)
111-
.map(
112-
// use random numbers to avoid spurious compiler optimizations wrt to branching
113-
|_| format!("hi{:?}", vs.choose(rng)),
114-
)
115-
.collect();
116-
let keys: Vec<&str> = keys.iter().map(|e| &**e).collect();
89+
// Randomly choose from 4 distinct key values; a higher number increases sparseness.
90+
let key_suffixes = [0, 1, 2, 3];
91+
let keys = StringArray::from_iter_values(
92+
(0..batch_size).map(|_| format!("hi{}", key_suffixes.choose(rng).unwrap())),
93+
);
11794

118-
let values = create_data(batch_size, 0.5);
95+
let values = create_data(rng, batch_size, 0.5);
11996

12097
// Integer values between [0, u64::MAX].
121-
let integer_values_wide = create_integer_data(rng, batch_size, 9.0);
98+
let integer_values_wide = (0..batch_size)
99+
.map(|_| rng.random::<u64>())
100+
.collect::<Vec<_>>();
122101

123-
// Integer values between [0, 9].
102+
// Integer values between [0, 1000).
103+
let integer_values_mid = (0..batch_size)
104+
.map(|_| rng.random_range(0..1000))
105+
.collect::<Vec<_>>();
106+
107+
// Integer values between [0, 10).
124108
let integer_values_narrow = (0..batch_size)
125-
.map(|_| rng.random_range(0_u64..10))
109+
.map(|_| rng.random_range(0..10))
126110
.collect::<Vec<_>>();
127111

128112
RecordBatch::try_new(
129113
schema,
130114
vec![
131-
Arc::new(StringArray::from(keys)),
132-
Arc::new(Float32Array::from(vec![i as f32; batch_size])),
115+
Arc::new(keys),
116+
Arc::new(Float32Array::from(vec![batch_index as f32; batch_size])),
133117
Arc::new(Float64Array::from(values)),
134118
Arc::new(UInt64Array::from(integer_values_wide)),
119+
Arc::new(UInt64Array::from(integer_values_mid)),
135120
Arc::new(UInt64Array::from(integer_values_narrow)),
136121
],
137122
)
@@ -140,21 +125,29 @@ fn create_record_batch(
140125

141126
/// Create record batches of `partitions_len` partitions and `batch_size` for each batch,
142127
/// with a total number of `array_len` records
143-
#[expect(clippy::needless_pass_by_value)]
144128
pub fn create_record_batches(
145-
schema: SchemaRef,
129+
schema: &SchemaRef,
146130
array_len: usize,
147131
partitions_len: usize,
148132
batch_size: usize,
149133
) -> Vec<Vec<RecordBatch>> {
150134
let mut rng = StdRng::seed_from_u64(42);
151-
(0..partitions_len)
152-
.map(|_| {
153-
(0..array_len / batch_size / partitions_len)
154-
.map(|i| create_record_batch(schema.clone(), &mut rng, batch_size, i))
155-
.collect::<Vec<_>>()
156-
})
157-
.collect::<Vec<_>>()
135+
let mut partitions = Vec::with_capacity(partitions_len);
136+
let batches_per_partition = array_len / batch_size / partitions_len;
137+
138+
for _ in 0..partitions_len {
139+
let mut batches = Vec::with_capacity(batches_per_partition);
140+
for batch_index in 0..batches_per_partition {
141+
batches.push(create_record_batch(
142+
schema.clone(),
143+
&mut rng,
144+
batch_size,
145+
batch_index,
146+
));
147+
}
148+
partitions.push(batches);
149+
}
150+
partitions
158151
}
159152

160153
/// An enum that wraps either a regular StringBuilder or a GenericByteViewBuilder

0 commit comments

Comments
 (0)