diff --git a/CHANGES.md b/CHANGES.md index 788ffd08154f..dbe5a0544556 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -89,6 +89,7 @@ * (Java) BigQueryIO now treats a 404 when deleting a temporary table or dataset as success, so a replayed work item whose earlier attempt already deleted it no longer retries forever ([#24997](https://github.com/apache/beam/issues/24997)). * (Java) IcebergIO now writes rows containing `EnumerationType` (proto enum) fields as strings, instead of throwing `Unsupported Beam logical type Enum` ([#40299](https://github.com/apache/beam/issues/40299)). * (Python) Fixed stateful DoFns with side inputs sometimes taking the timer key coder from a side input instead of the main input, which could make the worker fail to decode timer keys with `Unknown type tag` ([#40374](https://github.com/apache/beam/issues/40374)). +* (Python) `Duration` built from float seconds now rounds to the nearest microsecond instead of truncating, which could lose a microsecond ([#40263](https://github.com/apache/beam/issues/40263)). * Fixed X (Java/Python) ([#X](https://github.com/apache/beam/issues/X)). ## Security Fixes diff --git a/sdks/python/apache_beam/utils/timestamp.py b/sdks/python/apache_beam/utils/timestamp.py index c09cc93439c3..81d8430d3af7 100644 --- a/sdks/python/apache_beam/utils/timestamp.py +++ b/sdks/python/apache_beam/utils/timestamp.py @@ -497,7 +497,10 @@ def __init__( self, seconds: Union[int, float] = 0, micros: Union[int, float] = 0) -> None: - self.micros = int(seconds * 1000000) + int(micros) + # Round rather than truncate, as the Timestamp constructor does: the float + # multiplication can land just below the exact integer (for example + # 2.000002 * 1e6 == 2000001.9999999998) and int() would drop a microsecond. + self.micros = round(seconds * 1000000) + int(micros) @staticmethod def of(seconds: DurationTypes) -> 'Duration': diff --git a/sdks/python/apache_beam/utils/timestamp_test.py b/sdks/python/apache_beam/utils/timestamp_test.py index ff37ff84b5de..ebdd42d37029 100644 --- a/sdks/python/apache_beam/utils/timestamp_test.py +++ b/sdks/python/apache_beam/utils/timestamp_test.py @@ -411,6 +411,16 @@ def test_of(self): with self.assertRaises(TypeError): Duration.of(Timestamp(10)) + def test_constructor_float_rounds_to_nearest(self): + # Same float-to-micros conversion as Timestamp: truncating dropped a + # microsecond whenever seconds * 1e6 landed just below the integer. + self.assertEqual(Duration(2.000002).micros, 2000002) + self.assertEqual(Duration(-2.000002).micros, -2000002) + self.assertEqual(Duration.of(1.000001).micros, 1000001) + self.assertEqual(Duration(1.5).micros, 1500000) + for micros in (1056803002554, 544368929943, 33956373746): + self.assertEqual(Duration(micros / 1000000).micros, micros) + def test_precision(self): self.assertEqual(Duration(10000000) % 0.1, 0) self.assertEqual(Duration(10000000) % 0.05, 0)