@@ -31,14 +31,8 @@ use datafusion_common::ScalarValue;
3131use datafusion_common:: test_util:: batches_to_sort_string;
3232use datafusion_execution:: object_store:: ObjectStoreUrl ;
3333
34- use datafusion_common:: config:: TableParquetOptions ;
3534use datafusion_datasource:: file_scan_config:: FileScanConfigBuilder ;
3635use datafusion_datasource:: source:: DataSourceExec ;
37- use datafusion_physical_expr:: expressions:: { Column , IsNotNullExpr } ;
38- use datafusion_physical_expr_common:: physical_expr:: PhysicalExpr ;
39- use datafusion_physical_plan:: ExecutionPlan ;
40- use datafusion_physical_plan:: filter:: FilterExec ;
41- use datafusion_physical_plan:: projection:: ProjectionExec ;
4236use insta:: assert_snapshot;
4337use object_store:: ObjectMeta ;
4438use parquet:: arrow:: ArrowWriter ;
@@ -152,9 +146,7 @@ async fn multi_parquet_coercion_projection() {
152146 " ) ;
153147}
154148
155- /// Writes a parquet file with mixed column encoding: `rle_cols` use RLE_DICTIONARY,
156- /// all other string columns use plain encoding. Returns the file so the caller
157- /// can keep it alive.
149+ /// Writes `batch` to a temp parquet file where `rle_cols` use RLE_DICTIONARY encoding and all other columns use plain.
158150fn store_mixed_encoding_parquet (
159151 batch : & RecordBatch ,
160152 rle_cols : & [ & str ] ,
@@ -178,9 +170,9 @@ fn store_mixed_encoding_parquet(
178170 ( meta, output)
179171}
180172
181- /// Logical planning: after `register_parquet` with the flag on, the table schema
182- /// already has RLE-encoded columns as `Dictionary(Int32, Utf8)` before any query
183- /// executes. Plain string columns remain `Utf8View` .
173+ /// Promotion happens at logical planning time : after `register_parquet` with the
174+ /// flag on, `ctx.table()` already shows RLE columns as `Dictionary(Int32, Utf8)`
175+ /// before any query executes .
184176#[ tokio:: test]
185177async fn rle_dictionary_logical_schema_promotion ( ) {
186178 let schema = Arc :: new ( Schema :: new ( vec ! [
@@ -195,7 +187,6 @@ async fn rle_dictionary_logical_schema_promotion() {
195187 ] ,
196188 )
197189 . unwrap ( ) ;
198- // env → RLE_DICTIONARY, version → plain
199190 let ( _meta, file) = store_mixed_encoding_parquet ( & batch, & [ "env" ] ) ;
200191
201192 let config = SessionConfig :: new ( ) . set (
@@ -211,9 +202,7 @@ async fn rle_dictionary_logical_schema_promotion() {
211202 . await
212203 . unwrap ( ) ;
213204
214- // Inspect the logical schema — no query execution.
215- let df = ctx. table ( "t" ) . await . unwrap ( ) ;
216- let logical_schema = df. schema ( ) ;
205+ let logical_schema = ctx. table ( "t" ) . await . unwrap ( ) . schema ( ) . clone ( ) ;
217206 let dict_type =
218207 DataType :: Dictionary ( Box :: new ( DataType :: Int32 ) , Box :: new ( DataType :: Utf8 ) ) ;
219208 assert_eq ! (
@@ -233,10 +222,10 @@ async fn rle_dictionary_logical_schema_promotion() {
233222}
234223
235224/// Five string columns, three RLE_DICTIONARY encoded and two plain.
236- /// With the flag on DataFusion should surface only the three RLE columns as
237- /// Dictionary(Int32, Utf8) and leave the plain columns as Utf8View.
225+ /// Flag on: only the three RLE columns become `Dictionary(Int32, Utf8)`; plain
226+ /// columns stay `Utf8View`. Flag off: all columns stay ` Utf8View` .
238227#[ tokio:: test]
239- async fn rle_dictionary_selective_promotion_only_rle_columns_become_dict ( ) {
228+ async fn rle_dictionary_selective_promotion ( ) {
240229 let schema = Arc :: new ( Schema :: new ( vec ! [
241230 Field :: new( "env" , DataType :: Utf8 , true ) ,
242231 Field :: new( "region" , DataType :: Utf8 , true ) ,
@@ -262,183 +251,47 @@ async fn rle_dictionary_selective_promotion_only_rle_columns_become_dict() {
262251 . unwrap ( ) ;
263252
264253 let ( _meta, file) = store_mixed_encoding_parquet ( & batch, & [ "env" , "region" , "tier" ] ) ;
265-
266- let config = SessionConfig :: new ( ) . set (
267- "datafusion.execution.parquet.enable_rle_to_dictionary" ,
268- & ScalarValue :: Boolean ( Some ( true ) ) ,
269- ) ;
270- let ctx = SessionContext :: new_with_config ( config) ;
271- ctx. register_parquet (
272- "t" ,
273- file. path ( ) . to_str ( ) . unwrap ( ) ,
274- ParquetReadOptions :: default ( ) ,
275- )
276- . await
277- . unwrap ( ) ;
278-
279- let result = ctx
280- . sql ( "SELECT env, region, tier, version, build FROM t LIMIT 1" )
281- . await
282- . unwrap ( )
283- . collect ( )
284- . await
285- . unwrap ( ) ;
286-
287- assert ! ( !result. is_empty( ) ) ;
288- let out_schema = result[ 0 ] . schema ( ) ;
289- let dict_type =
290- DataType :: Dictionary ( Box :: new ( DataType :: Int32 ) , Box :: new ( DataType :: Utf8 ) ) ;
291- assert_eq ! (
292- out_schema. field_with_name( "env" ) . unwrap( ) . data_type( ) ,
293- & dict_type
294- ) ;
295- assert_eq ! (
296- out_schema. field_with_name( "region" ) . unwrap( ) . data_type( ) ,
297- & dict_type
298- ) ;
299- assert_eq ! (
300- out_schema. field_with_name( "tier" ) . unwrap( ) . data_type( ) ,
301- & dict_type
302- ) ;
303- assert_eq ! (
304- out_schema. field_with_name( "version" ) . unwrap( ) . data_type( ) ,
305- & DataType :: Utf8View
306- ) ;
307- assert_eq ! (
308- out_schema. field_with_name( "build" ) . unwrap( ) . data_type( ) ,
309- & DataType :: Utf8View
310- ) ;
311- }
312-
313- fn store_rle_dict_parquet ( batch : & RecordBatch ) -> ( ObjectMeta , NamedTempFile ) {
314- let mut output = tempfile:: Builder :: new ( )
315- . suffix ( ".parquet" )
316- . tempfile ( )
317- . expect ( "creating temp file" ) ;
318- let props = WriterProperties :: builder ( )
319- . set_dictionary_enabled ( true )
320- . build ( ) ;
321- let mut writer = ArrowWriter :: try_new ( & mut output, batch. schema ( ) , Some ( props) )
322- . expect ( "creating writer" ) ;
323- writer. write ( batch) . expect ( "Writing batch" ) ;
324- writer. close ( ) . unwrap ( ) ;
325- let meta = local_unpartitioned_file ( & output) ;
326- ( meta, output)
327- }
328-
329- /// Runs scan → filter (IS NOT NULL on `env`) → project through a physical plan
330- /// and returns the collected batches. Exercises planning-time schema consistency:
331- /// FilterExec and ProjectionExec reference columns by the type the scan declares,
332- /// so a mismatch between the plan's declared type and the arrays the scan produces
333- /// will surface here.
334- async fn scan_filter_project_env (
335- meta : ObjectMeta ,
336- source : Arc < ParquetSource > ,
337- ) -> Vec < RecordBatch > {
338- let task_ctx = SessionContext :: new ( ) . task_ctx ( ) ;
339- let conf = FileScanConfigBuilder :: new ( ObjectStoreUrl :: local_filesystem ( ) , source)
340- . with_file_group ( vec ! [ meta. into( ) ] . into ( ) )
341- . build ( ) ;
342- let scan: Arc < dyn ExecutionPlan > = DataSourceExec :: from_data_source ( conf) ;
343-
344- let env_col = Arc :: new ( Column :: new ( "env" , 0 ) ) as Arc < dyn PhysicalExpr > ;
345- let not_null = Arc :: new ( IsNotNullExpr :: new ( Arc :: clone ( & env_col) ) ) ;
346- let filter: Arc < dyn ExecutionPlan > =
347- Arc :: new ( FilterExec :: try_new ( not_null, Arc :: clone ( & scan) ) . unwrap ( ) ) ;
348-
349- let projection_exprs = vec ! [ ( env_col, "env" . to_string( ) ) ] ;
350- let projection = Arc :: new ( ProjectionExec :: try_new ( projection_exprs, filter) . unwrap ( ) )
351- as Arc < dyn ExecutionPlan > ;
352-
353- collect ( projection, task_ctx) . await . unwrap ( )
354- }
355-
356- fn env_batch ( ) -> RecordBatch {
357- let schema = Arc :: new ( Schema :: new ( vec ! [ Field :: new( "env" , DataType :: Utf8 , true ) ] ) ) ;
358- let envs = Arc :: new ( StringArray :: from ( vec ! [
359- Some ( "prod" ) ,
360- Some ( "staging" ) ,
361- Some ( "prod" ) ,
362- Some ( "prod" ) ,
363- Some ( "canary" ) ,
364- Some ( "staging" ) ,
365- Some ( "prod" ) ,
366- Some ( "canary" ) ,
367- ] ) ) ;
368- RecordBatch :: try_new ( schema, vec ! [ envs as _] ) . unwrap ( )
369- }
370-
371- /// Direct ParquetSource (no infer_schema): flag off + Utf8 schema stays Utf8;
372- /// flag on + pre-promoted Dict schema (as infer_schema would produce) gives Dict.
373- /// FilterExec + ProjectionExec catch any schema mismatch at runtime.
374- #[ tokio:: test]
375- async fn rle_dictionary_direct_source ( ) {
376- let batch = env_batch ( ) ;
377- let utf8_schema = batch. schema ( ) ;
378- let ( meta, _file) = store_rle_dict_parquet ( & batch) ;
379- let dict_type =
380- DataType :: Dictionary ( Box :: new ( DataType :: Int32 ) , Box :: new ( DataType :: Utf8 ) ) ;
381-
382- let read = scan_filter_project_env (
383- meta. clone ( ) ,
384- Arc :: new ( ParquetSource :: new ( Arc :: clone ( & utf8_schema) ) ) ,
385- )
386- . await ;
387- assert_eq ! ( read[ 0 ] . schema( ) . field( 0 ) . data_type( ) , & DataType :: Utf8 ) ;
388-
389- let dict_schema = Arc :: new ( Schema :: new ( vec ! [ Field :: new(
390- "env" ,
391- dict_type. clone( ) ,
392- true ,
393- ) ] ) ) ;
394- let mut options = TableParquetOptions :: default ( ) ;
395- options. global . enable_rle_to_dictionary = true ;
396- let source =
397- Arc :: new ( ParquetSource :: new ( dict_schema) . with_table_parquet_options ( options) ) ;
398- let read = scan_filter_project_env ( meta, source) . await ;
399- assert_eq ! ( read[ 0 ] . schema( ) . field( 0 ) . data_type( ) , & dict_type) ;
400- }
401-
402- /// E2E scan → filter → aggregate: flag off gives Utf8View, flag on gives Dict.
403- #[ tokio:: test]
404- async fn rle_dictionary_e2e_aggregate ( ) {
405- let ( _meta, file) = store_rle_dict_parquet ( & env_batch ( ) ) ;
406254 let path = file. path ( ) . to_str ( ) . unwrap ( ) ;
407- let sql = "SELECT env, COUNT(*) AS cnt FROM envs WHERE env IS NOT NULL GROUP BY env " ;
255+ let sql = "SELECT env, region, tier, version, build FROM t LIMIT 1 " ;
408256 let dict_type =
409257 DataType :: Dictionary ( Box :: new ( DataType :: Int32 ) , Box :: new ( DataType :: Utf8 ) ) ;
410258
259+ // Flag off: all string columns stay Utf8View.
411260 let ctx = SessionContext :: new ( ) ;
412- ctx. register_parquet ( "envs " , path, ParquetReadOptions :: default ( ) )
261+ ctx. register_parquet ( "t " , path, ParquetReadOptions :: default ( ) )
413262 . await
414263 . unwrap ( ) ;
415264 let result = ctx. sql ( sql) . await . unwrap ( ) . collect ( ) . await . unwrap ( ) ;
416- assert_eq ! (
417- result[ 0 ]
418- . schema ( )
419- . field_with_name ( "env" )
420- . unwrap( )
421- . data_type ( ) ,
422- & DataType :: Utf8View ,
423- ) ;
265+ assert ! ( !result . is_empty ( ) ) ;
266+ let s = result[ 0 ] . schema ( ) ;
267+ for col in [ "env" , "region" , "tier" , "version" , "build" ] {
268+ assert_eq ! (
269+ s . field_with_name ( col ) . unwrap( ) . data_type ( ) ,
270+ & DataType :: Utf8View
271+ ) ;
272+ }
424273
274+ // Flag on: only RLE-encoded columns become Dict; plain columns stay Utf8View.
425275 let config = SessionConfig :: new ( ) . set (
426276 "datafusion.execution.parquet.enable_rle_to_dictionary" ,
427277 & ScalarValue :: Boolean ( Some ( true ) ) ,
428278 ) ;
429279 let ctx = SessionContext :: new_with_config ( config) ;
430- ctx. register_parquet ( "envs " , path, ParquetReadOptions :: default ( ) )
280+ ctx. register_parquet ( "t " , path, ParquetReadOptions :: default ( ) )
431281 . await
432282 . unwrap ( ) ;
433283 let result = ctx. sql ( sql) . await . unwrap ( ) . collect ( ) . await . unwrap ( ) ;
434- assert_eq ! (
435- result[ 0 ]
436- . schema( )
437- . field_with_name( "env" )
438- . unwrap( )
439- . data_type( ) ,
440- & dict_type,
441- ) ;
284+ assert ! ( !result. is_empty( ) ) ;
285+ let s = result[ 0 ] . schema ( ) ;
286+ for col in [ "env" , "region" , "tier" ] {
287+ assert_eq ! ( s. field_with_name( col) . unwrap( ) . data_type( ) , & dict_type) ;
288+ }
289+ for col in [ "version" , "build" ] {
290+ assert_eq ! (
291+ s. field_with_name( col) . unwrap( ) . data_type( ) ,
292+ & DataType :: Utf8View
293+ ) ;
294+ }
442295}
443296
444297/// Writes `batches` to a temporary parquet file
@@ -450,14 +303,10 @@ pub fn store_parquet(
450303 . into_iter ( )
451304 . map ( |batch| {
452305 let mut output = NamedTempFile :: new ( ) . expect ( "creating temp file" ) ;
453-
454- let builder = WriterProperties :: builder ( ) ;
455- let props = builder. build ( ) ;
456-
306+ let props = WriterProperties :: builder ( ) . build ( ) ;
457307 let mut writer =
458308 ArrowWriter :: try_new ( & mut output, batch. schema ( ) , Some ( props) )
459309 . expect ( "creating writer" ) ;
460-
461310 writer. write ( & batch) . expect ( "Writing batch" ) ;
462311 writer. close ( ) . unwrap ( ) ;
463312 output
0 commit comments