Skip to content

Commit d824985

Browse files
authored
Merge branch 'main' into optimize_mem_hash_join_arrow
2 parents d814704 + 206ac6b commit d824985

62 files changed

Lines changed: 3011 additions & 1066 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

benchmarks/queries/clickbench/README.md

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -228,6 +228,22 @@ Results look like
228228
Elapsed 30.195 seconds.
229229
```
230230

231+
232+
### Q9-Q12: FIRST_VALUE Aggregation Performance
233+
234+
These queries test the performance of the `FIRST_VALUE` aggregation function with different data types and grouping cardinalities.
235+
236+
| Query | `FIRST_VALUE` Column | Column Type | Group By Column | Group By Type | Number of Groups |
237+
|-------|----------------------|-------------|-----------------|---------------|------------------|
238+
| Q9 | `URL` | `Utf8` | `UserID` | `Int64` | 17,630,976 |
239+
| Q10 | `URL` | `Utf8` | `OS` | `Int16` | 91 |
240+
| Q11 | `WatchID` | `Int64` | `UserID` | `Int64` | 17,630,976 |
241+
| Q12 | `WatchID` | `Int64` | `OS` | `Int16` | 91 |
242+
243+
244+
245+
246+
231247
## Data Notes
232248

233249
Here are some interesting statistics about the data used in the queries
Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,8 @@
1+
-- Must set for ClickBench hits_partitioned dataset. See https://github.com/apache/datafusion/issues/16591
2+
-- set datafusion.execution.parquet.binary_as_string = true
3+
4+
SELECT MAX(len) FROM (
5+
SELECT LENGTH(FIRST_VALUE("URL" ORDER BY "EventTime")) as len
6+
FROM hits
7+
GROUP BY "OS"
8+
);
Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,8 @@
1+
-- Must set for ClickBench hits_partitioned dataset. See https://github.com/apache/datafusion/issues/16591
2+
-- set datafusion.execution.parquet.binary_as_string = true
3+
4+
SELECT MAX(fv) FROM (
5+
SELECT FIRST_VALUE("WatchID" ORDER BY "EventTime") as fv
6+
FROM hits
7+
GROUP BY "UserID"
8+
);
Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,8 @@
1+
-- Must set for ClickBench hits_partitioned dataset. See https://github.com/apache/datafusion/issues/16591
2+
-- set datafusion.execution.parquet.binary_as_string = true
3+
4+
SELECT MAX(fv) FROM (
5+
SELECT FIRST_VALUE("WatchID" ORDER BY "EventTime") as fv
6+
FROM hits
7+
GROUP BY "OS"
8+
);
Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,4 @@
1+
-- Must set for ClickBench hits_partitioned dataset. See https://github.com/apache/datafusion/issues/16591
2+
-- set datafusion.execution.parquet.binary_as_string = true
3+
4+
SELECT "RegionID", "UserAgent", "OS", AVG(to_timestamp("ResponseEndTiming")-to_timestamp("ResponseStartTiming")) as avg_response_time, AVG(to_timestamp("ResponseEndTiming")-to_timestamp("ConnectTiming")) as avg_latency FROM hits GROUP BY "RegionID", "UserAgent", "OS" ORDER BY avg_latency DESC limit 10;
Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,8 @@
1+
-- Must set for ClickBench hits_partitioned dataset. See https://github.com/apache/datafusion/issues/16591
2+
-- set datafusion.execution.parquet.binary_as_string = true
3+
4+
SELECT MAX(len) FROM (
5+
SELECT LENGTH(FIRST_VALUE("URL" ORDER BY "EventTime")) as len
6+
FROM hits
7+
GROUP BY "UserID"
8+
);

datafusion-examples/examples/data_io/parquet_index.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -462,7 +462,7 @@ impl PruningStatistics for ParquetMetadataIndex {
462462
}
463463

464464
/// return the row counts for each file
465-
fn row_counts(&self, _column: &Column) -> Option<ArrayRef> {
465+
fn row_counts(&self) -> Option<ArrayRef> {
466466
Some(self.row_counts_ref().clone())
467467
}
468468

datafusion-examples/examples/query_planning/pruning.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -174,7 +174,7 @@ impl PruningStatistics for MyCatalog {
174174
None
175175
}
176176

177-
fn row_counts(&self, _column: &Column) -> Option<ArrayRef> {
177+
fn row_counts(&self) -> Option<ArrayRef> {
178178
// In this example, we know nothing about the number of rows in each file
179179
None
180180
}

datafusion/common/src/config.rs

Lines changed: 211 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -687,6 +687,137 @@ config_namespace! {
687687
}
688688
}
689689

690+
/// Options for content-defined chunking (CDC) when writing parquet files.
691+
/// See [`ParquetOptions::use_content_defined_chunking`].
692+
///
693+
/// Can be enabled with default options by setting
694+
/// `use_content_defined_chunking` to `true`, or configured with sub-fields
695+
/// like `use_content_defined_chunking.min_chunk_size`.
696+
#[derive(Debug, Clone, PartialEq)]
697+
pub struct CdcOptions {
698+
/// Minimum chunk size in bytes. The rolling hash will not trigger a split
699+
/// until this many bytes have been accumulated. Default is 256 KiB.
700+
pub min_chunk_size: usize,
701+
702+
/// Maximum chunk size in bytes. A split is forced when the accumulated
703+
/// size exceeds this value. Default is 1 MiB.
704+
pub max_chunk_size: usize,
705+
706+
/// Normalization level. Increasing this improves deduplication ratio
707+
/// but increases fragmentation. Recommended range is [-3, 3], default is 0.
708+
pub norm_level: i32,
709+
}
710+
711+
// Note: `CdcOptions` intentionally does NOT implement `Default` so that the
712+
// blanket `impl<F: ConfigField + Default> ConfigField for Option<F>` does not
713+
// apply. This allows the specific `impl ConfigField for Option<CdcOptions>`
714+
// below to handle "true"/"false" for enabling/disabling CDC.
715+
// Use `CdcOptions::default()` (the inherent method) instead of `Default::default()`.
716+
impl CdcOptions {
717+
/// Returns a new `CdcOptions` with default values.
718+
#[expect(clippy::should_implement_trait)]
719+
pub fn default() -> Self {
720+
Self {
721+
min_chunk_size: 256 * 1024,
722+
max_chunk_size: 1024 * 1024,
723+
norm_level: 0,
724+
}
725+
}
726+
}
727+
728+
impl ConfigField for CdcOptions {
729+
fn set(&mut self, key: &str, value: &str) -> Result<()> {
730+
let (key, rem) = key.split_once('.').unwrap_or((key, ""));
731+
match key {
732+
"min_chunk_size" => self.min_chunk_size.set(rem, value),
733+
"max_chunk_size" => self.max_chunk_size.set(rem, value),
734+
"norm_level" => self.norm_level.set(rem, value),
735+
_ => _config_err!("Config value \"{}\" not found on CdcOptions", key),
736+
}
737+
}
738+
739+
fn visit<V: Visit>(&self, v: &mut V, key_prefix: &str, _description: &'static str) {
740+
let key = format!("{key_prefix}.min_chunk_size");
741+
self.min_chunk_size.visit(v, &key, "Minimum chunk size in bytes. The rolling hash will not trigger a split until this many bytes have been accumulated. Default is 256 KiB.");
742+
let key = format!("{key_prefix}.max_chunk_size");
743+
self.max_chunk_size.visit(v, &key, "Maximum chunk size in bytes. A split is forced when the accumulated size exceeds this value. Default is 1 MiB.");
744+
let key = format!("{key_prefix}.norm_level");
745+
self.norm_level.visit(v, &key, "Normalization level. Increasing this improves deduplication ratio but increases fragmentation. Recommended range is [-3, 3], default is 0.");
746+
}
747+
748+
fn reset(&mut self, key: &str) -> Result<()> {
749+
let (key, rem) = key.split_once('.').unwrap_or((key, ""));
750+
match key {
751+
"min_chunk_size" => {
752+
if rem.is_empty() {
753+
self.min_chunk_size = CdcOptions::default().min_chunk_size;
754+
Ok(())
755+
} else {
756+
self.min_chunk_size.reset(rem)
757+
}
758+
}
759+
"max_chunk_size" => {
760+
if rem.is_empty() {
761+
self.max_chunk_size = CdcOptions::default().max_chunk_size;
762+
Ok(())
763+
} else {
764+
self.max_chunk_size.reset(rem)
765+
}
766+
}
767+
"norm_level" => {
768+
if rem.is_empty() {
769+
self.norm_level = CdcOptions::default().norm_level;
770+
Ok(())
771+
} else {
772+
self.norm_level.reset(rem)
773+
}
774+
}
775+
_ => _config_err!("Config value \"{}\" not found on CdcOptions", key),
776+
}
777+
}
778+
}
779+
780+
/// `ConfigField` for `Option<CdcOptions>` — allows setting the option to
781+
/// `"true"` (enable with defaults) or `"false"` (disable), in addition to
782+
/// setting individual sub-fields like `min_chunk_size`.
783+
impl ConfigField for Option<CdcOptions> {
784+
fn visit<V: Visit>(&self, v: &mut V, key: &str, description: &'static str) {
785+
match self {
786+
Some(s) => s.visit(v, key, description),
787+
None => v.none(key, description),
788+
}
789+
}
790+
791+
fn set(&mut self, key: &str, value: &str) -> Result<()> {
792+
if key.is_empty() {
793+
match value.to_ascii_lowercase().as_str() {
794+
"true" => {
795+
*self = Some(CdcOptions::default());
796+
Ok(())
797+
}
798+
"false" => {
799+
*self = None;
800+
Ok(())
801+
}
802+
_ => _config_err!(
803+
"Expected 'true' or 'false' for use_content_defined_chunking, got '{value}'"
804+
),
805+
}
806+
} else {
807+
self.get_or_insert_with(CdcOptions::default).set(key, value)
808+
}
809+
}
810+
811+
fn reset(&mut self, key: &str) -> Result<()> {
812+
if key.is_empty() {
813+
*self = None;
814+
Ok(())
815+
} else {
816+
self.get_or_insert_with(CdcOptions::default).reset(key)
817+
}
818+
}
819+
}
820+
690821
config_namespace! {
691822
/// Options for reading and writing parquet files
692823
///
@@ -872,6 +1003,12 @@ config_namespace! {
8721003
/// writing out already in-memory data, such as from a cached
8731004
/// data frame.
8741005
pub maximum_buffered_record_batches_per_stream: usize, default = 2
1006+
1007+
/// (writing) EXPERIMENTAL: Enable content-defined chunking (CDC) when writing
1008+
/// parquet files. When `Some`, CDC is enabled with the given options; when `None`
1009+
/// (the default), CDC is disabled. When CDC is enabled, parallel writing is
1010+
/// automatically disabled since the chunker state must persist across row groups.
1011+
pub use_content_defined_chunking: Option<CdcOptions>, default = None
8751012
}
8761013
}
8771014

@@ -1826,6 +1963,7 @@ config_field!(usize);
18261963
config_field!(f64);
18271964
config_field!(u64);
18281965
config_field!(u32);
1966+
config_field!(i32);
18291967

18301968
impl ConfigField for u8 {
18311969
fn visit<V: Visit>(&self, v: &mut V, key: &str, description: &'static str) {
@@ -3579,4 +3717,77 @@ mod tests {
35793717
"Invalid or Unsupported Configuration: Invalid parquet writer version: 3.0. Expected one of: 1.0, 2.0"
35803718
);
35813719
}
3720+
3721+
#[cfg(feature = "parquet")]
3722+
#[test]
3723+
fn set_cdc_option_with_boolean_true() {
3724+
use crate::config::ConfigOptions;
3725+
3726+
let mut config = ConfigOptions::default();
3727+
assert!(
3728+
config
3729+
.execution
3730+
.parquet
3731+
.use_content_defined_chunking
3732+
.is_none()
3733+
);
3734+
3735+
// Setting to "true" should enable CDC with default options
3736+
config
3737+
.set(
3738+
"datafusion.execution.parquet.use_content_defined_chunking",
3739+
"true",
3740+
)
3741+
.unwrap();
3742+
let cdc = config
3743+
.execution
3744+
.parquet
3745+
.use_content_defined_chunking
3746+
.as_ref()
3747+
.expect("CDC should be enabled");
3748+
assert_eq!(cdc.min_chunk_size, 256 * 1024);
3749+
assert_eq!(cdc.max_chunk_size, 1024 * 1024);
3750+
assert_eq!(cdc.norm_level, 0);
3751+
3752+
// Setting to "false" should disable CDC
3753+
config
3754+
.set(
3755+
"datafusion.execution.parquet.use_content_defined_chunking",
3756+
"false",
3757+
)
3758+
.unwrap();
3759+
assert!(
3760+
config
3761+
.execution
3762+
.parquet
3763+
.use_content_defined_chunking
3764+
.is_none()
3765+
);
3766+
}
3767+
3768+
#[cfg(feature = "parquet")]
3769+
#[test]
3770+
fn set_cdc_option_with_subfields() {
3771+
use crate::config::ConfigOptions;
3772+
3773+
let mut config = ConfigOptions::default();
3774+
3775+
// Setting sub-fields should also enable CDC
3776+
config
3777+
.set(
3778+
"datafusion.execution.parquet.use_content_defined_chunking.min_chunk_size",
3779+
"1024",
3780+
)
3781+
.unwrap();
3782+
let cdc = config
3783+
.execution
3784+
.parquet
3785+
.use_content_defined_chunking
3786+
.as_ref()
3787+
.expect("CDC should be enabled");
3788+
assert_eq!(cdc.min_chunk_size, 1024);
3789+
// Other fields should be defaults
3790+
assert_eq!(cdc.max_chunk_size, 1024 * 1024);
3791+
assert_eq!(cdc.norm_level, 0);
3792+
}
35823793
}

0 commit comments

Comments
 (0)