This is an automated email from the ASF dual-hosted git repository.
tvalentyn pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/master by this push:
new 261f05f142b [Python] Fix LogElements crash on global window timestamps
(#40332)
261f05f142b is described below
commit 261f05f142b1ed1d7c901faa3b51fe8117c8b142
Author: Anish Mehta <[email protected]>
AuthorDate: Tue Oct 6 03:50:25 2026 +0530
[Python] Fix LogElements crash on global window timestamps (#40332)
---
sdks/python/apache_beam/transforms/util.py | 7 ++++++-
sdks/python/apache_beam/transforms/util_test.py | 8 ++++++++
2 files changed, 14 insertions(+), 1 deletion(-)
diff --git a/sdks/python/apache_beam/transforms/util.py
b/sdks/python/apache_beam/transforms/util.py
index 30f8dc01c45..e5fcf369342 100644
--- a/sdks/python/apache_beam/transforms/util.py
+++ b/sdks/python/apache_beam/transforms/util.py
@@ -1978,7 +1978,12 @@ class LogElements(PTransform):
def format_timestamp(self, timestamp):
if self.use_epoch_time:
return timestamp.seconds()
- return timestamp.to_rfc3339()
+ try:
+ return timestamp.to_rfc3339()
+ except OverflowError:
+ # MIN_TIMESTAMP and the global window bounds are outside the datetime
+ # range.
+ return str(timestamp)
def process(
self,
diff --git a/sdks/python/apache_beam/transforms/util_test.py
b/sdks/python/apache_beam/transforms/util_test.py
index 509f2f2c897..5a935731ca7 100644
--- a/sdks/python/apache_beam/transforms/util_test.py
+++ b/sdks/python/apache_beam/transforms/util_test.py
@@ -2481,6 +2481,14 @@ class LogElementsTest(unittest.TestCase):
| util.LogElements(prefix='prefix_'))
assert_that(result, equal_to(['a', 'b', 'c']))
+ def test_global_window_with_timestamp_and_window(self):
+ with TestPipeline() as p:
+ result = (
+ p
+ | beam.Create(['a'])
+ | util.LogElements(with_timestamp=True, with_window=True))
+ assert_that(result, equal_to(['a']))
+
@pytest.fixture(scope="function")
def _capture_logs(request, caplog):
with caplog.at_level(logging.INFO):