laskoviymishka commented on code in PR #3327:
URL: https://github.com/apache/iceberg-rust/pull/3327#discussion_r4232681735
##########
crates/iceberg/src/writer/partitioning/fanout_writer.rs:
##########
@@ -381,4 +382,86 @@ mod tests {
Ok(())
}
+
+ #[tokio::test]
+ async fn test_fanout_writer_signed_zero_partitions() -> Result<()> {
+ let temp_dir = TempDir::new()?;
+ let file_io = FileIO::new_with_fs();
+ let location_gen = DefaultLocationGenerator::with_data_location(
+ temp_dir.path().to_str().unwrap().to_string(),
+ );
+ let file_name_gen =
+ DefaultFileNameGenerator::new("test".to_string(), None,
DataFileFormat::Parquet);
+
+ let schema = Arc::new(
+ crate::spec::Schema::builder()
+ .with_schema_id(1)
+ .with_fields(vec![
+ NestedField::required(1, "id",
Type::Primitive(PrimitiveType::Int)).into(),
+ NestedField::required(2, "d",
Type::Primitive(PrimitiveType::Double)).into(),
+ ])
+ .build()?,
+ );
+ let partition_spec = PartitionSpec::builder(schema.clone())
+ .add_partition_field("d", "d", Transform::Identity)?
+ .build()?;
+
+ let partition_value_neg_zero =
Struct::from_iter([Some(Literal::double(-0.0))]);
+ let partition_key_neg_zero = PartitionKey::new(
+ partition_spec.clone(),
+ schema.clone(),
+ partition_value_neg_zero.clone(),
+ );
+
+ let partition_value_pos_zero =
Struct::from_iter([Some(Literal::double(0.0))]);
+ let partition_key_pos_zero = PartitionKey::new(
+ partition_spec.clone(),
+ schema.clone(),
+ partition_value_pos_zero.clone(),
+ );
+
+ let parquet_writer_builder =
+ ParquetWriterBuilder::new(WriterProperties::builder().build(),
schema.clone());
+ let rolling_writer_builder =
RollingFileWriterBuilder::new_with_default_file_size(
+ parquet_writer_builder,
+ file_io.clone(),
+ location_gen,
+ file_name_gen,
+ );
+ let data_file_writer_builder =
DataFileWriterBuilder::new(rolling_writer_builder);
+
+ let mut writer = FanoutWriter::new(data_file_writer_builder);
+
+ let arrow_schema = Arc::new(schema_to_arrow_schema(&schema)?);
+ let batch_neg_zero = RecordBatch::try_new(arrow_schema.clone(), vec![
+ Arc::new(Int32Array::from(vec![1])),
+ Arc::new(Float64Array::from(vec![-0.0])),
+ ])?;
+ let batch_pos_zero = RecordBatch::try_new(arrow_schema.clone(), vec![
+ Arc::new(Int32Array::from(vec![2])),
+ Arc::new(Float64Array::from(vec![0.0])),
+ ])?;
+
+ writer.write(partition_key_neg_zero, batch_neg_zero).await?;
+ writer.write(partition_key_pos_zero, batch_pos_zero).await?;
+
+ let data_files = writer.close().await?;
+
+ // -0.0 and 0.0 are different partition values, as in iceberg-java, so
each row
+ // goes to its own data file.
+ assert_eq!(data_files.len(), 2);
+ for data_file in &data_files {
+ assert_eq!(data_file.record_count, 1);
+ }
+ let partitions_found: std::collections::HashSet<Struct> = data_files
Review Comment:
Tiny follow-up, not blocking: `HashSet` is fully qualified here and just
below while the rest of the module imports its types up top. A `use
std::collections::HashSet;` in the test module would read cleaner.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]