|
26 | 26 | //! select * from data limit 10; |
27 | 27 | //! ``` |
28 | 28 |
|
29 | | -use arrow::array::{ArrayRef, Int32Array, StringArray}; |
| 29 | +use arrow::array::{ArrayRef, Int32Array, Int64Array, StringArray}; |
30 | 30 | use arrow::compute::concat_batches; |
31 | 31 | use arrow::error::ArrowError; |
32 | 32 | use arrow::record_batch::RecordBatch; |
| 33 | +use datafusion::execution::SessionStateBuilder; |
33 | 34 | use datafusion::physical_plan::metrics::{MetricValue, MetricsSet}; |
34 | 35 | use datafusion::physical_plan::{collect, displayable}; |
35 | 36 | use datafusion::prelude::{ |
@@ -809,3 +810,91 @@ async fn pushed_down_predicate_reports_the_original_error() { |
809 | 810 | "expected the original cast error, got {root:?}" |
810 | 811 | ); |
811 | 812 | } |
| 813 | + |
| 814 | +/// A MIN/MAX dynamic filter must survive a file that lacks the aggregated |
| 815 | +/// column. That file's Partial aggregate evaluates to a typed null |
| 816 | +/// (`Int64(NULL)`), which must not replace the real MIN published later. |
| 817 | +/// |
| 818 | +/// With one partition the files are read in order `01_missing`, `02_high`, |
| 819 | +/// `03_low`, so the bad interleaving happens on every run. Without the fix, |
| 820 | +/// the MIN bound stays null, the filter becomes `latency_ms > 204`, and |
| 821 | +/// `03_low` is pruned, returning MIN = 200. |
| 822 | +#[tokio::test] |
| 823 | +async fn aggregate_dynamic_filter_ignores_typed_null_bound() { |
| 824 | + let tempdir = TempDir::new_in(Path::new(".")).unwrap(); |
| 825 | + let dir = tempdir.path(); |
| 826 | + |
| 827 | + let write = |name: &str, col_name: &str, array: ArrayRef| { |
| 828 | + let batch = RecordBatch::try_from_iter([(col_name, array)]).unwrap(); |
| 829 | + let file = File::create(dir.join(name)).unwrap(); |
| 830 | + let mut writer = ArrowWriter::try_new(file, batch.schema(), None).unwrap(); |
| 831 | + writer.write(&batch).unwrap(); |
| 832 | + writer.close().unwrap(); |
| 833 | + }; |
| 834 | + write( |
| 835 | + "01_missing.parquet", |
| 836 | + "host", |
| 837 | + Arc::new(StringArray::from(vec!["h1"; 5])), |
| 838 | + ); |
| 839 | + write( |
| 840 | + "02_high.parquet", |
| 841 | + "latency_ms", |
| 842 | + Arc::new(Int64Array::from(vec![200, 201, 202, 203, 204])), |
| 843 | + ); |
| 844 | + write( |
| 845 | + "03_low.parquet", |
| 846 | + "latency_ms", |
| 847 | + Arc::new(Int64Array::from(vec![100, 101, 102, 103, 104])), |
| 848 | + ); |
| 849 | + |
| 850 | + // One partition keeps the file order fixed. Drop the rule that would |
| 851 | + // otherwise fold Partial/Final into a Single aggregate, which does not |
| 852 | + // create a dynamic filter. |
| 853 | + let default_state = SessionStateBuilder::new().with_default_features().build(); |
| 854 | + let rules = default_state |
| 855 | + .physical_optimizers() |
| 856 | + .iter() |
| 857 | + .filter(|rule| rule.name() != "CombinePartialFinalAggregate") |
| 858 | + .cloned() |
| 859 | + .collect(); |
| 860 | + let config = SessionConfig::new().with_target_partitions(1); |
| 861 | + let state = SessionStateBuilder::new() |
| 862 | + .with_config(config) |
| 863 | + .with_default_features() |
| 864 | + .with_physical_optimizer_rules(rules) |
| 865 | + .build(); |
| 866 | + let ctx = SessionContext::new_with_state(state); |
| 867 | + |
| 868 | + ctx.sql(&format!( |
| 869 | + "CREATE EXTERNAL TABLE t (latency_ms BIGINT, host VARCHAR) \ |
| 870 | + STORED AS PARQUET LOCATION '{}/'", |
| 871 | + dir.display() |
| 872 | + )) |
| 873 | + .await |
| 874 | + .unwrap(); |
| 875 | + |
| 876 | + let df = ctx |
| 877 | + .sql("SELECT min(latency_ms), max(latency_ms) FROM t") |
| 878 | + .await |
| 879 | + .unwrap(); |
| 880 | + let plan = df.create_physical_plan().await.unwrap(); |
| 881 | + let plan_str = displayable(plan.as_ref()).indent(false).to_string(); |
| 882 | + assert!( |
| 883 | + plan_str.contains("mode=Partial") && plan_str.contains("DynamicFilter"), |
| 884 | + "expected a Partial aggregate with a dynamic filter:\n{plan_str}" |
| 885 | + ); |
| 886 | + |
| 887 | + let batches = collect(plan, ctx.task_ctx()).await.unwrap(); |
| 888 | + let batch = concat_batches(&batches[0].schema(), &batches).unwrap(); |
| 889 | + let min = batch |
| 890 | + .column(0) |
| 891 | + .as_any() |
| 892 | + .downcast_ref::<Int64Array>() |
| 893 | + .unwrap(); |
| 894 | + let max = batch |
| 895 | + .column(1) |
| 896 | + .as_any() |
| 897 | + .downcast_ref::<Int64Array>() |
| 898 | + .unwrap(); |
| 899 | + assert_eq!((min.value(0), max.value(0)), (100, 204), "\n{plan_str}"); |
| 900 | +} |
0 commit comments