@@ -4,6 +4,7 @@ use std::sync::Arc;
44
55use delta_kernel:: arrow:: array:: {
66 Array , BinaryArray , BooleanArray , Int32Array , RecordBatch , StringArray , StructArray ,
7+ TimestampMicrosecondArray ,
78} ;
89use delta_kernel:: arrow:: compute:: filter_record_batch;
910use delta_kernel:: engine:: arrow_data:: ArrowEngineData ;
@@ -391,3 +392,194 @@ async fn empty_string_partition_pruning(#[values(false, true)] native_checkpoint
391392 "empty-bytes file must be pruned under p_bin = X'6f74686572'"
392393 ) ;
393394}
395+
396+ async fn write_timezone_partition_table ( table_path : & std:: path:: Path , timestamp : & str ) -> Url {
397+ std:: fs:: create_dir_all ( table_path) . unwrap ( ) ;
398+ let url = Url :: from_directory_path ( table_path) . unwrap ( ) ;
399+ let table_root = url. to_string ( ) ;
400+ let store: Arc < DynObjectStore > = Arc :: new ( LocalFileSystem :: new ( ) ) ;
401+ let schema_string = serde_json:: json!( {
402+ "type" : "struct" ,
403+ "fields" : [
404+ { "name" : "p_ts" , "type" : "timestamp" , "nullable" : true , "metadata" : { } } ,
405+ { "name" : "p_ntz" , "type" : "timestamp_ntz" , "nullable" : true , "metadata" : { } } ,
406+ { "name" : "p_int" , "type" : "integer" , "nullable" : true , "metadata" : { } } ,
407+ { "name" : "value" , "type" : "integer" , "nullable" : true , "metadata" : { } } ,
408+ ] ,
409+ } )
410+ . to_string ( ) ;
411+ let protocol = r#"{"protocol":{"minReaderVersion":3,"minWriterVersion":7,"readerFeatures":["timestampNtz"],"writerFeatures":["timestampNtz"]}}"# ;
412+ let metadata = serde_json:: json!( {
413+ "metaData" : {
414+ "id" : "00000000-0000-0000-0000-000000000000" ,
415+ "format" : { "provider" : "parquet" , "options" : { } } ,
416+ "schemaString" : schema_string,
417+ "partitionColumns" : [ "p_ts" , "p_ntz" , "p_int" ] ,
418+ "configuration" : { "delta.checkpoint.writeStatsAsStruct" : "true" } ,
419+ "createdTime" : 1700000000000_i64 ,
420+ } ,
421+ } )
422+ . to_string ( ) ;
423+ add_commit (
424+ & table_root,
425+ store. as_ref ( ) ,
426+ 0 ,
427+ format ! ( "{protocol}\n {metadata}" ) ,
428+ )
429+ . await
430+ . unwrap ( ) ;
431+ let add = serde_json:: json!( {
432+ "add" : {
433+ "path" : "part.parquet" ,
434+ "partitionValues" : {
435+ "p_ts" : timestamp,
436+ "p_ntz" : "2024-01-15 12:30:45.123456" ,
437+ "p_int" : "7" ,
438+ } ,
439+ "size" : 100 ,
440+ "modificationTime" : 1700000000000_i64 ,
441+ "dataChange" : true ,
442+ "stats" : "{\" numRecords\" :1}" ,
443+ } ,
444+ } ) ;
445+ add_commit ( & table_root, store. as_ref ( ) , 1 , add. to_string ( ) )
446+ . await
447+ . unwrap ( ) ;
448+ url
449+ }
450+
451+ #[ rstest]
452+ #[ case:: utc( None , "2024-01-15 12:30:45.123456" , "2024-01-15T12:30:45.123456Z" ) ]
453+ #[ case:: los_angeles_winter(
454+ Some ( "America/Los_Angeles" ) ,
455+ "2024-01-15 12:30:45.123456" ,
456+ "2024-01-15T20:30:45.123456Z"
457+ ) ]
458+ #[ case:: los_angeles_summer(
459+ Some ( "America/Los_Angeles" ) ,
460+ "2024-06-15 08:00:00.500500" ,
461+ "2024-06-15T15:00:00.500500Z"
462+ ) ]
463+ #[ case:: fixed_offset_seconds(
464+ Some ( "+12:45:30" ) ,
465+ "2024-01-15 12:30:45.123456" ,
466+ "2024-01-14T23:45:15.123456Z"
467+ ) ]
468+ #[ case:: explicit_offset_wins(
469+ Some ( "America/Los_Angeles" ) ,
470+ "2024-01-15 12:30:45+02:00" ,
471+ "2024-01-15T10:30:45Z"
472+ ) ]
473+ #[ case:: dst_overlap_uses_earlier_instant(
474+ Some ( "America/Los_Angeles" ) ,
475+ "2024-11-03 01:30:00" ,
476+ "2024-11-03T08:30:00Z"
477+ ) ]
478+ #[ case:: dst_gap_uses_pre_transition_offset(
479+ Some ( "America/Los_Angeles" ) ,
480+ "2024-03-10 02:30:00" ,
481+ "2024-03-10T10:30:00Z"
482+ ) ]
483+ #[ tokio:: test( flavor = "multi_thread" , worker_threads = 2 ) ]
484+ async fn timezone_aware_parsed_partition_values_across_json_and_native_checkpoint (
485+ #[ case] timestamp_timezone : Option < & str > ,
486+ #[ case] raw_timestamp : & str ,
487+ #[ case] expected_timestamp : & str ,
488+ #[ values( false , true ) ] native_checkpoint : bool ,
489+ ) {
490+ let temp_dir = tempfile:: tempdir ( ) . unwrap ( ) ;
491+ let url = write_timezone_partition_table ( temp_dir. path ( ) , raw_timestamp) . await ;
492+ let engine = create_default_engine_mt_executor ( & url) . unwrap ( ) ;
493+ if native_checkpoint {
494+ let snapshot = Snapshot :: builder_for ( url. clone ( ) )
495+ . build ( engine. as_ref ( ) )
496+ . unwrap ( ) ;
497+ snapshot. checkpoint ( engine. as_ref ( ) , None ) . unwrap ( ) ;
498+ }
499+
500+ let snapshot = Snapshot :: builder_for ( url) . build ( engine. as_ref ( ) ) . unwrap ( ) ;
501+ let partition_values = timestamp_timezone
502+ . map_or_else ( PartitionValuesOptions :: with_struct, |timezone| {
503+ PartitionValuesOptions :: with_struct ( ) . with_timestamp_timezone ( timezone)
504+ } ) ;
505+ let scan = snapshot
506+ . scan_builder ( )
507+ . with_partition_values ( partition_values)
508+ . build ( )
509+ . unwrap ( ) ;
510+ let mut rows = 0 ;
511+ for metadata in scan. scan_metadata ( engine. as_ref ( ) ) . unwrap ( ) {
512+ let ( data, selection) = metadata. unwrap ( ) . scan_files . into_parts ( ) ;
513+ let batch: RecordBatch = ArrowEngineData :: try_from_engine_data ( data) . unwrap ( ) . into ( ) ;
514+ let batch = filter_record_batch ( & batch, & BooleanArray :: from ( selection) ) . unwrap ( ) ;
515+ if batch. num_rows ( ) == 0 {
516+ continue ;
517+ }
518+ let parsed = get_column ! ( batch, "partitionValues_parsed" , StructArray ) ;
519+ let timestamp = parsed
520+ . column_by_name ( "p_ts" )
521+ . unwrap ( )
522+ . as_any ( )
523+ . downcast_ref :: < TimestampMicrosecondArray > ( )
524+ . unwrap ( ) ;
525+ let timestamp_ntz = parsed
526+ . column_by_name ( "p_ntz" )
527+ . unwrap ( )
528+ . as_any ( )
529+ . downcast_ref :: < TimestampMicrosecondArray > ( )
530+ . unwrap ( ) ;
531+ let integer = parsed
532+ . column_by_name ( "p_int" )
533+ . unwrap ( )
534+ . as_any ( )
535+ . downcast_ref :: < Int32Array > ( )
536+ . unwrap ( ) ;
537+ assert_eq ! (
538+ timestamp. value( 0 ) ,
539+ chrono:: DateTime :: parse_from_rfc3339( expected_timestamp)
540+ . unwrap( )
541+ . timestamp_micros( )
542+ ) ;
543+ assert_eq ! (
544+ timestamp_ntz. value( 0 ) ,
545+ chrono:: DateTime :: parse_from_rfc3339( "2024-01-15T12:30:45.123456Z" )
546+ . unwrap( )
547+ . timestamp_micros( ) ,
548+ "reader timezone must not affect TIMESTAMP_NTZ"
549+ ) ;
550+ assert_eq ! ( integer. value( 0 ) , 7 ) ;
551+ rows += batch. num_rows ( ) ;
552+ }
553+ assert_eq ! ( rows, 1 ) ;
554+ }
555+
556+ #[ rstest]
557+ #[ case:: invalid_zone( "2024-01-15 12:30:45" , "Not/AZone" , "Not/AZone" ) ]
558+ #[ case:: invalid_timestamp( "not a timestamp" , "America/Los_Angeles" , "not a timestamp" ) ]
559+ #[ tokio:: test( flavor = "multi_thread" , worker_threads = 2 ) ]
560+ async fn timezone_aware_parsed_partition_values_report_invalid_input (
561+ #[ case] raw_timestamp : & str ,
562+ #[ case] timestamp_timezone : & str ,
563+ #[ case] expected_error : & str ,
564+ ) {
565+ let temp_dir = tempfile:: tempdir ( ) . unwrap ( ) ;
566+ let url = write_timezone_partition_table ( temp_dir. path ( ) , raw_timestamp) . await ;
567+ let engine = create_default_engine_mt_executor ( & url) . unwrap ( ) ;
568+ let snapshot = Snapshot :: builder_for ( url) . build ( engine. as_ref ( ) ) . unwrap ( ) ;
569+ let scan = snapshot
570+ . scan_builder ( )
571+ . with_partition_values (
572+ PartitionValuesOptions :: with_struct ( ) . with_timestamp_timezone ( timestamp_timezone) ,
573+ )
574+ . build ( )
575+ . unwrap ( ) ;
576+ let result = scan
577+ . scan_metadata ( engine. as_ref ( ) )
578+ . unwrap ( )
579+ . next ( )
580+ . expect ( "one metadata result" ) ;
581+ let Err ( error) = result else {
582+ panic ! ( "invalid timestamp input must fail map parsing" ) ;
583+ } ;
584+ assert ! ( error. to_string( ) . contains( expected_error) , "{error}" ) ;
585+ }
0 commit comments