This is an automated email from the ASF dual-hosted git repository.

ahmedabu98 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 f5feab8c598 Enhance Python Timestamp to be precision-variable up to 
nanos, and map it to Timestamp logical type (#39537)
f5feab8c598 is described below

commit f5feab8c598a4cbd03f1bd4cdad8f540c977a335
Author: Ahmed Abualsaud <[email protected]>
AuthorDate: Thu Jul 30 14:08:26 2026 -0700

    Enhance Python Timestamp to be precision-variable up to nanos, and map it 
to Timestamp logical type (#39537)
    
    * integrate logical timestamp type with python
    
    * trigger xlang ITs
    
    * Changes
    
    * Changes.md
    
    ---------
    
    Co-authored-by: Ahmed Abualsaud <[email protected]>
---
 .../beam_PostCommit_Python_Xlang_Gcp_Direct.json   |   2 +-
 .../beam_PostCommit_Python_Xlang_IO_Direct.json    |   2 +-
 CHANGES.md                                         |   3 +
 .../beam/model/fnexecution/v1/standard_coders.yaml |  31 ++
 .../org/apache/beam/model/pipeline/v1/schema.proto |  13 +
 .../apache/beam/sdk/schemas/SchemaTranslation.java |   2 +
 .../beam/sdk/schemas/logicaltypes/Timestamp.java   |   9 +-
 .../beam/sdk/schemas/SchemaTranslationTest.java    |   7 +
 sdks/python/apache_beam/coders/row_coder_test.py   |  54 ++++
 sdks/python/apache_beam/portability/common_urns.py |   1 +
 sdks/python/apache_beam/typehints/schemas.py       |  94 +++++-
 sdks/python/apache_beam/typehints/schemas_test.py  |  94 ++++++
 sdks/python/apache_beam/utils/timestamp.py         | 327 +++++++++++++++++----
 sdks/python/apache_beam/utils/timestamp_test.py    | 187 +++++++++++-
 14 files changed, 753 insertions(+), 73 deletions(-)

diff --git a/.github/trigger_files/beam_PostCommit_Python_Xlang_Gcp_Direct.json 
b/.github/trigger_files/beam_PostCommit_Python_Xlang_Gcp_Direct.json
index c5309eebb07..e3d6056a5de 100644
--- a/.github/trigger_files/beam_PostCommit_Python_Xlang_Gcp_Direct.json
+++ b/.github/trigger_files/beam_PostCommit_Python_Xlang_Gcp_Direct.json
@@ -1,4 +1,4 @@
 {
   "comment": "Modify this file in a trivial way to cause this test suite to 
run",
-  "modification": 17
+  "modification": 1
 }
diff --git a/.github/trigger_files/beam_PostCommit_Python_Xlang_IO_Direct.json 
b/.github/trigger_files/beam_PostCommit_Python_Xlang_IO_Direct.json
index d6a91b7e2e8..e3d6056a5de 100644
--- a/.github/trigger_files/beam_PostCommit_Python_Xlang_IO_Direct.json
+++ b/.github/trigger_files/beam_PostCommit_Python_Xlang_IO_Direct.json
@@ -1,4 +1,4 @@
 {
   "comment": "Modify this file in a trivial way to cause this test suite to 
run",
-  "modification": 7
+  "modification": 1
 }
diff --git a/CHANGES.md b/CHANGES.md
index c5da08bd1a5..75a362d1580 100644
--- a/CHANGES.md
+++ b/CHANGES.md
@@ -73,6 +73,9 @@
 * (Java) Supported acknowledge mode for JmsIO 
([#39253](https://github.com/apache/beam/issues/39253)).
 * (Python) Staged files directory is now automatically added to `sys.path` on 
the Python SDK worker at startup. This makes Python files provided via the 
'--files_to_stage' pipeline option importable in the pipeline code and makes it 
easier to initialize Python SDK harness at startup via the `--beam_plugins` 
pipeline option. For more information, see the [Staging Individual 
Files](https://beam.apache.org/documentation/sdks/python-pipeline-dependencies/#staging-files)
 section of the depend [...]
 * (Python) Added `equal_to_approx`, an `assert_that` matcher that compares 
numeric pipeline outputs with a configurable tolerance 
([#18028](https://github.com/apache/beam/issues/18028)).
+* (Python) `Timestamp` now supports variable subsecond precision, up to 
nanoseconds. The portable
+  `beam:logical_type:timestamp:v1` logical type now maps to Python's 
`Timestamp` ([#39344](https://github.com/apache/beam/issues/39344)).
+* X feature added (Java/Python) 
([#X](https://github.com/apache/beam/issues/X)).
 
 ## Breaking Changes
 
diff --git 
a/model/fn-execution/src/main/resources/org/apache/beam/model/fnexecution/v1/standard_coders.yaml
 
b/model/fn-execution/src/main/resources/org/apache/beam/model/fnexecution/v1/standard_coders.yaml
index 8ce9d275154..eb799968007 100644
--- 
a/model/fn-execution/src/main/resources/org/apache/beam/model/fnexecution/v1/standard_coders.yaml
+++ 
b/model/fn-execution/src/main/resources/org/apache/beam/model/fnexecution/v1/standard_coders.yaml
@@ -626,3 +626,34 @@ examples:
   "\u0001\u0003\u0061\u0062\u0063" : "abc"
   "\u0001\u000a\u006d\u006f\u0072\u0065\u0020\u0062\u0079\u0074\u0065\u0073" : 
"more bytes"
   "\u0000" : null
+
+---
+
+coder:
+  urn: "beam:coder:row:v1"
+  # f_timestamp: logical(timestamp:v1, precision=3, representation 
row<seconds:int64, subseconds:int16>)
+  payload: 
"\x0a\x88\x01\x0a\x0bf_timestamp\x1ay:w\x0a\x1ebeam:logical_type:timestamp:v1\x1aK2I\x0aG\x0a\x0d\x0a\x07seconds\x1a\x02\x10\x04\x0a\x10\x0a\x0asubseconds\x1a\x02\x10\x02\x12$5f2fbf6c-8b8e-4a30-a9be-e0f65f0f79c3\"\x02\x10\x03*\x04\x0a\x02\x18\x03\x12$9c7c1b6c-2a5e-4c53-9b32-91b71c1f8ac3"
+examples:
+  "\x01\x00\x02\x00\xb6\x95\xd5\xf9\x05\x03\xe7": {f_timestamp: {seconds: 
1597328054, subseconds: 999}}
+  "\x01\x00\x02\x00\xfe\xff\xff\xff\xff\xff\xff\xff\xff\x01\x01\xf4": 
{f_timestamp: {seconds: -2, subseconds: 500}}
+  "\x01\x00\x02\x00\xb6\x95\xd5\xf9\x05\x00\x00": {f_timestamp: {seconds: 
1597328054, subseconds: 0}}
+
+---
+
+coder:
+  urn: "beam:coder:row:v1"
+  # f_timestamp: logical(timestamp:v1, precision=6, representation 
row<seconds:int64, subseconds:int32>)
+  payload: 
"\x0a\x88\x01\x0a\x0bf_timestamp\x1ay:w\x0a\x1ebeam:logical_type:timestamp:v1\x1aK2I\x0aG\x0a\x0d\x0a\x07seconds\x1a\x02\x10\x04\x0a\x10\x0a\x0asubseconds\x1a\x02\x10\x03\x12$b8b3faf8-63f8-4b40-9f4d-3c1e83d4e07b\"\x02\x10\x03*\x04\x0a\x02\x18\x06\x12$e0a72e10-b214-4fd2-a4dc-fca0e7f7d924"
+examples:
+  "\x01\x00\x02\x00\xb6\x95\xd5\xf9\x05\xc0\xc4\x07": {f_timestamp: {seconds: 
1597328054, subseconds: 123456}}
+  "\x01\x00\x02\x00\xfe\xff\xff\xff\xff\xff\xff\xff\xff\x01\xa0\xc2\x1e": 
{f_timestamp: {seconds: -2, subseconds: 500000}}
+
+---
+
+coder:
+  urn: "beam:coder:row:v1"
+  # f_timestamp: logical(timestamp:v1, precision=9, representation 
row<seconds:int64, subseconds:int32>)
+  payload: 
"\x0a\x88\x01\x0a\x0bf_timestamp\x1ay:w\x0a\x1ebeam:logical_type:timestamp:v1\x1aK2I\x0aG\x0a\x0d\x0a\x07seconds\x1a\x02\x10\x04\x0a\x10\x0a\x0asubseconds\x1a\x02\x10\x03\x12$3d6f7b81-2f13-4a41-8be3-4d8cd8f30f42\"\x02\x10\x03*\x04\x0a\x02\x18\x09\x12$cf03a3c1-4a3a-4f77-8e2b-6e5b19c19f27"
+examples:
+  "\x01\x00\x02\x00\xb6\x95\xd5\xf9\x05\x95\x9a\xef:": {f_timestamp: {seconds: 
1597328054, subseconds: 123456789}}
+  
"\x01\x00\x02\x00\xfe\xff\xff\xff\xff\xff\xff\xff\xff\x01\xff\x93\xeb\xdc\x03": 
{f_timestamp: {seconds: -2, subseconds: 999999999}}
diff --git 
a/model/pipeline/src/main/proto/org/apache/beam/model/pipeline/v1/schema.proto 
b/model/pipeline/src/main/proto/org/apache/beam/model/pipeline/v1/schema.proto
index d514be101be..c114c35f2d9 100644
--- 
a/model/pipeline/src/main/proto/org/apache/beam/model/pipeline/v1/schema.proto
+++ 
b/model/pipeline/src/main/proto/org/apache/beam/model/pipeline/v1/schema.proto
@@ -234,6 +234,19 @@ message LogicalTypes {
     //     since the epoch.
     DATE = 8 [(org.apache.beam.model.pipeline.v1.beam_urn) =
       "beam:logical_type:date:v1"];
+
+    // A URN for Timestamp type
+    //   - Representation type: ROW<seconds: INT64, subseconds: INT16> for
+    //     precision < 5, otherwise ROW<seconds: INT64, subseconds: INT32>
+    //   - Argument type: INT32.
+    //     A timestamp without a timezone, represented as seconds since the
+    //     epoch plus a non-negative fraction of a second. The argument is
+    //     the precision: the number of decimal digits of the subseconds
+    //     field (e.g. 3 for milliseconds, 6 for micros, 9 for nanos).
+    //     subseconds stays non-negative even for timestamps before the epoch,
+    //     e.g. -1.5s at precision 6 is {seconds: -2, subseconds: 500000}.
+    TIMESTAMP = 9 [(org.apache.beam.model.pipeline.v1.beam_urn) =
+      "beam:logical_type:timestamp:v1"];
   }
 }
 
diff --git 
a/sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/SchemaTranslation.java
 
b/sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/SchemaTranslation.java
index d527b294923..8122704e643 100644
--- 
a/sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/SchemaTranslation.java
+++ 
b/sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/SchemaTranslation.java
@@ -50,6 +50,7 @@ import org.apache.beam.sdk.schemas.logicaltypes.FixedString;
 import org.apache.beam.sdk.schemas.logicaltypes.MicrosInstant;
 import org.apache.beam.sdk.schemas.logicaltypes.PythonCallable;
 import org.apache.beam.sdk.schemas.logicaltypes.SchemaLogicalType;
+import org.apache.beam.sdk.schemas.logicaltypes.Timestamp;
 import org.apache.beam.sdk.schemas.logicaltypes.UnknownLogicalType;
 import org.apache.beam.sdk.schemas.logicaltypes.VariableBytes;
 import org.apache.beam.sdk.schemas.logicaltypes.VariableString;
@@ -115,6 +116,7 @@ public class SchemaTranslation {
           .put(FixedString.IDENTIFIER, FixedString.class)
           .put(VariableString.IDENTIFIER, VariableString.class)
           .put(Date.IDENTIFIER, Date.class)
+          .put(Timestamp.IDENTIFIER, Timestamp.class)
           .build();
 
   public static SchemaApi.Schema schemaToProto(Schema schema, boolean 
serializeLogicalType) {
diff --git 
a/sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/logicaltypes/Timestamp.java
 
b/sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/logicaltypes/Timestamp.java
index 87e47f5961e..51fa9008b74 100644
--- 
a/sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/logicaltypes/Timestamp.java
+++ 
b/sdks/java/core/src/main/java/org/apache/beam/sdk/schemas/logicaltypes/Timestamp.java
@@ -22,6 +22,8 @@ import static 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Pr
 import static 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkState;
 
 import java.time.Instant;
+import org.apache.beam.model.pipeline.v1.RunnerApi;
+import org.apache.beam.model.pipeline.v1.SchemaApi;
 import org.apache.beam.sdk.schemas.Schema;
 import org.apache.beam.sdk.values.Row;
 import org.checkerframework.checker.nullness.qual.NonNull;
@@ -56,7 +58,12 @@ import org.checkerframework.checker.nullness.qual.NonNull;
  * }</pre>
  */
 public class Timestamp implements Schema.LogicalType<Instant, Row> {
-  public static final String IDENTIFIER = "beam:logical_type:timestamp:v1";
+  public static final String IDENTIFIER =
+      SchemaApi.LogicalTypes.Enum.TIMESTAMP
+          .getValueDescriptor()
+          .getOptions()
+          .getExtension(RunnerApi.beamUrn);
+
   static final int MIN_PRECISION = 0;
   static final int MAX_PRECISION = 9;
 
diff --git 
a/sdks/java/core/src/test/java/org/apache/beam/sdk/schemas/SchemaTranslationTest.java
 
b/sdks/java/core/src/test/java/org/apache/beam/sdk/schemas/SchemaTranslationTest.java
index 8a382c4ad1b..925b1e02b78 100644
--- 
a/sdks/java/core/src/test/java/org/apache/beam/sdk/schemas/SchemaTranslationTest.java
+++ 
b/sdks/java/core/src/test/java/org/apache/beam/sdk/schemas/SchemaTranslationTest.java
@@ -51,6 +51,7 @@ import org.apache.beam.sdk.schemas.logicaltypes.NanosInstant;
 import org.apache.beam.sdk.schemas.logicaltypes.PythonCallable;
 import org.apache.beam.sdk.schemas.logicaltypes.SchemaLogicalType;
 import org.apache.beam.sdk.schemas.logicaltypes.SqlTypes;
+import org.apache.beam.sdk.schemas.logicaltypes.Timestamp;
 import org.apache.beam.sdk.schemas.logicaltypes.UnknownLogicalType;
 import org.apache.beam.sdk.schemas.logicaltypes.VariableBytes;
 import org.apache.beam.sdk.schemas.logicaltypes.VariableString;
@@ -460,6 +461,12 @@ public class SchemaTranslationTest {
           .add(FieldType.logicalType(FixedString.of(10)))
           .add(FieldType.logicalType(VariableString.of(10)))
           .add(FieldType.logicalType(FixedPrecisionNumeric.of(10)))
+          .add(FieldType.logicalType(Timestamp.MILLIS))
+          .add(FieldType.logicalType(Timestamp.MICROS))
+          .add(FieldType.logicalType(Timestamp.NANOS))
+          // INT16/INT32 subseconds representation boundary.
+          .add(FieldType.logicalType(Timestamp.of(4)))
+          .add(FieldType.logicalType(Timestamp.of(5)))
           .add(FieldType.logicalType(new PortableNullArgLogicalType()))
           .add(FieldType.logicalType(new NullArgumentLogicalType()))
           .build();
diff --git a/sdks/python/apache_beam/coders/row_coder_test.py 
b/sdks/python/apache_beam/coders/row_coder_test.py
index 396cdc9b3f6..28170721c28 100644
--- a/sdks/python/apache_beam/coders/row_coder_test.py
+++ b/sdks/python/apache_beam/coders/row_coder_test.py
@@ -203,6 +203,60 @@ class RowCoderTest(unittest.TestCase):
     for test_case in self.PEOPLE:
       self.assertEqual(test_case, coder.decode(coder.encode(test_case)))
 
+  def test_create_row_coder_from_schema_with_parameterized_timestamp(self):
+    # Simulates schemas produced by the Java SDK's parameterized Timestamp
+    # logical type. Subseconds field uses INT16 for
+    # precision < 5 and INT32 otherwise.
+    for precision, subseconds_type in [(3, schema_pb2.INT16),
+                                       (4, schema_pb2.INT16),
+                                       (5, schema_pb2.INT32),
+                                       (6, schema_pb2.INT32),
+                                       (9, schema_pb2.INT32)]:
+      schema = schema_pb2.Schema(
+          id="timestamp_precision_%d" % precision,
+          fields=[
+              schema_pb2.Field(
+                  name="event_time",
+                  type=schema_pb2.FieldType(
+                      logical_type=schema_pb2.LogicalType(
+                          urn="beam:logical_type:timestamp:v1",
+                          representation=schema_pb2.FieldType(
+                              row_type=schema_pb2.RowType(
+                                  schema=schema_pb2.Schema(
+                                      id="timestamp_repr_%d" % precision,
+                                      fields=[
+                                          schema_pb2.Field(
+                                              name="seconds",
+                                              type=schema_pb2.FieldType(
+                                                  atomic_type=schema_pb2.INT64)
+                                          ),
+                                          schema_pb2.Field(
+                                              name="subseconds",
+                                              type=schema_pb2.FieldType(
+                                                  
atomic_type=subseconds_type)),
+                                      ]))),
+                          argument_type=schema_pb2.FieldType(
+                              atomic_type=schema_pb2.INT32),
+                          argument=schema_pb2.FieldValue(
+                              atomic_value=schema_pb2.AtomicTypeValue(
+                                  int32=precision))))),
+          ])
+      coder = RowCoder(schema)
+      row = named_tuple_from_schema(schema)
+
+      subseconds = 10**precision - 1
+      for value in [
+          Timestamp(seconds=1500000000,
+                    subseconds=subseconds,
+                    precision=precision),
+          Timestamp(seconds=-2, subseconds=subseconds, precision=precision),
+          Timestamp(seconds=1500000000, precision=precision),
+      ]:
+        test_case = row(event_time=value)
+        decoded = coder.decode(coder.encode(test_case))
+        self.assertEqual(test_case, decoded)
+        self.assertEqual(decoded.event_time.precision(), precision)
+
   def test_row_coder_negative_varint(self):
     schema = schema_pb2.Schema(
         id="negative",
diff --git a/sdks/python/apache_beam/portability/common_urns.py 
b/sdks/python/apache_beam/portability/common_urns.py
index 7777f63ffbe..913fb282f86 100644
--- a/sdks/python/apache_beam/portability/common_urns.py
+++ b/sdks/python/apache_beam/portability/common_urns.py
@@ -93,3 +93,4 @@ var_bytes = LogicalTypes.Enum.VAR_BYTES
 fixed_char = LogicalTypes.Enum.FIXED_CHAR
 var_char = LogicalTypes.Enum.VAR_CHAR
 date = LogicalTypes.Enum.DATE
+timestamp = LogicalTypes.Enum.TIMESTAMP
diff --git a/sdks/python/apache_beam/typehints/schemas.py 
b/sdks/python/apache_beam/typehints/schemas.py
index ca33e1eefcb..084ecc93581 100644
--- a/sdks/python/apache_beam/typehints/schemas.py
+++ b/sdks/python/apache_beam/typehints/schemas.py
@@ -977,6 +977,92 @@ MicrosInstantRepresentation = NamedTuple(
     'MicrosInstantRepresentation', [('seconds', np.int64),
                                     ('micros', np.int64)])
 
+ParameterizedTimestampRepresentation = NamedTuple(
+    'ParameterizedTimestampRepresentation', [('seconds', np.int64),
+                                             ('subseconds', np.int32)])
+
+# The subseconds field is INT16 for precision < 5
+_TIMESTAMP_SHORT_PRECISION_LIMIT = 5
+ParameterizedTimestampShortRepresentation = NamedTuple(
+    'ParameterizedTimestampShortRepresentation', [('seconds', np.int64),
+                                                  ('subseconds', np.int16)])
+
+
+@LogicalType._register_internal
+class ParameterizedTimestamp(LogicalType[Timestamp,
+                                         ParameterizedTimestampRepresentation,
+                                         np.int32]):
+  """Timestamp logical type parameterized by subsecond precision.
+
+  The argument is the precision: the number of decimal digits used to
+  represent the fraction of a second, e.g. 3 for milliseconds, 6 for
+  microseconds, 9 for nanoseconds.
+
+  Values are represented as a row of ``seconds`` (INT64, floored seconds
+  since the epoch) and ``subseconds`` (units of 10**-precision seconds,
+  always in ``[0, 10**precision)``.
+  ``subseconds`` is an INT16 field for precision < 5 and an INT32
+  field otherwise.
+
+  Note: Timestamp originating from Python to xlang still defaults to
+  MicrosInstant for backwards compatibility. To override the mapping of
+  Timestamp to this logical type, re-register using
+  :func:`~LogicalType.register_logical_type(ParameterizedTimestamp)`.
+  """
+  def __init__(self, precision: int = Timestamp.MICROS_PRECISION) -> None:
+    # The argument arrives as np.int32 when decoded from a schema proto.
+    precision = int(precision)
+    if not 0 <= precision <= Timestamp.NANOS_PRECISION:
+      raise ValueError(
+          'Timestamp precision must be between 0 and %d (inclusive), '
+          'but was %d.' % (Timestamp.NANOS_PRECISION, precision))
+    self._precision = precision
+
+  @classmethod
+  def urn(cls):
+    return common_urns.timestamp.urn
+
+  def representation_type(self) -> type:  # type: ignore[override]
+    # Unlike other logical types, the representation depends on the
+    # argument, so this is an instance method rather than a classmethod.
+    if self._precision < _TIMESTAMP_SHORT_PRECISION_LIMIT:
+      return ParameterizedTimestampShortRepresentation
+    return ParameterizedTimestampRepresentation
+
+  @classmethod
+  def language_type(cls):
+    return Timestamp
+
+  def to_representation_type(self, value: Timestamp):
+    # Verify that the value can be represented exactly at this type's precision
+    if value.precision() != self._precision:
+      value = value.to_precision(self._precision)
+    return self.representation_type()(value.seconds(), value.subseconds())
+
+  def to_language_type(self, value) -> Timestamp:
+    subseconds = int(value.subseconds)
+    # Match Java's toInputType: out-of-range subseconds indicate data
+    # corruption or a precision mismatch.
+    if not 0 <= subseconds < 10**self._precision:
+      raise ValueError(
+          'Invalid subseconds %d for Timestamp with precision %d.' %
+          (subseconds, self._precision))
+    return Timestamp(
+        seconds=int(value.seconds),
+        subseconds=subseconds,
+        precision=self._precision)
+
+  @classmethod
+  def argument_type(cls):
+    return np.int32
+
+  def argument(self):
+    return self._precision
+
+  @classmethod
+  def _from_typing(cls, typ):
+    return cls()
+
 
 @LogicalType._register_internal
 class MillisInstant(NoArgumentLogicalType[Timestamp, np.int64]):
@@ -1019,9 +1105,11 @@ class MillisInstant(NoArgumentLogicalType[Timestamp, 
np.int64]):
     return Timestamp(micros=millis * 1000)
 
 
-# Make sure MicrosInstant is registered after MillisInstant so that it
-# overwrites the mapping of Timestamp language type representation choice and
-# thus does not lose microsecond precision inside python sdk.
+# Make sure MicrosInstant is registered after MillisInstant and
+# ParameterizedTimestamp so that it overwrites the mapping of Timestamp
+# language type representation choice: plain Timestamp typehints keep their
+# historical micros_instant encoding and do not lose microsecond precision
+# inside the python sdk.
 @LogicalType._register_internal
 class MicrosInstant(NoArgumentLogicalType[Timestamp,
                                           MicrosInstantRepresentation]):
diff --git a/sdks/python/apache_beam/typehints/schemas_test.py 
b/sdks/python/apache_beam/typehints/schemas_test.py
index a8db60c9260..c2c21a7ce39 100644
--- a/sdks/python/apache_beam/typehints/schemas_test.py
+++ b/sdks/python/apache_beam/typehints/schemas_test.py
@@ -770,6 +770,100 @@ class SchemaTest(unittest.TestCase):
                                                      ('unique', Any)])
 
 
+class ParameterizedTimestampTest(unittest.TestCase):
+  def test_urn(self):
+    self.assertEqual(
+        schemas.ParameterizedTimestamp.urn(), 'beam:logical_type:timestamp:v1')
+
+  def test_timestamp_typehint_still_maps_to_micros_instant(self):
+    # Existing pipelines rely on plain Timestamp fields being encoded as
+    # micros_instant; registering ParameterizedTimestamp must not change
+    # the default.
+    field_type = typing_to_runner_api(Timestamp)
+    self.assertEqual(
+        field_type.logical_type.urn, common_urns.micros_instant.urn)
+
+  def test_from_runner_api_reconstructs_precision(self):
+    logical_type_proto = schema_pb2.LogicalType(
+        urn=common_urns.timestamp.urn,
+        representation=typing_to_runner_api(
+            schemas.ParameterizedTimestampRepresentation),
+        argument_type=schema_pb2.FieldType(atomic_type=schema_pb2.INT32),
+        argument=schema_pb2.FieldValue(
+            atomic_value=schema_pb2.AtomicTypeValue(int32=9)))
+    logical_type = schemas.LogicalType.from_runner_api(logical_type_proto)
+    self.assertIsInstance(logical_type, schemas.ParameterizedTimestamp)
+    self.assertEqual(logical_type.argument(), 9)
+    value = logical_type.to_language_type(
+        schemas.ParameterizedTimestampRepresentation(
+            seconds=np.int64(1234), subseconds=np.int32(123456789)))
+    self.assertEqual(
+        value, Timestamp(seconds=1234, subseconds=123456789, precision=9))
+    self.assertEqual(value.precision(), 9)
+
+  def test_representation_type_matches_java(self):
+    # The Java SDK uses an INT16 subseconds field for precision < 5 and
+    # INT32 otherwise; the wire formats differ, so we must match.
+    for precision in range(0, 5):
+      self.assertEqual(
+          schemas.ParameterizedTimestamp(precision).representation_type(),
+          schemas.ParameterizedTimestampShortRepresentation)
+    for precision in range(5, 10):
+      self.assertEqual(
+          schemas.ParameterizedTimestamp(precision).representation_type(),
+          schemas.ParameterizedTimestampRepresentation)
+
+  def test_precision_validation(self):
+    with self.assertRaises(ValueError):
+      schemas.ParameterizedTimestamp(10)
+    with self.assertRaises(ValueError):
+      schemas.ParameterizedTimestamp(-1)
+
+  def test_value_round_trip(self):
+    for precision, subseconds in [(3, 500), (6, 500000), (9, 123456789)]:
+      logical_type = schemas.ParameterizedTimestamp(precision)
+      value = Timestamp(
+          seconds=1234, subseconds=subseconds, precision=precision)
+      representation = logical_type.to_representation_type(value)
+      self.assertEqual(representation.seconds, 1234)
+      self.assertEqual(representation.subseconds, subseconds)
+      self.assertEqual(logical_type.to_language_type(representation), value)
+
+  def test_negative_timestamps_use_floored_seconds(self):
+    # -1.5s is represented as {seconds: -2, subseconds: 500000} at
+    # microsecond precision, matching the Java SDK and java.time.Instant.
+    logical_type = schemas.ParameterizedTimestamp(6)
+    representation = logical_type.to_representation_type(Timestamp(-1.5))
+    self.assertEqual(representation.seconds, -2)
+    self.assertEqual(representation.subseconds, 500000)
+    self.assertEqual(
+        logical_type.to_language_type(representation), Timestamp(-1.5))
+
+  def test_to_language_type_rejects_out_of_range_subseconds(self):
+    # Java's toInputType rejects these as likely data corruption.
+    logical_type = schemas.ParameterizedTimestamp(3)
+    with self.assertRaises(ValueError):
+      logical_type.to_language_type(
+          schemas.ParameterizedTimestampShortRepresentation(
+              np.int64(10), np.int16(5000)))
+    with self.assertRaises(ValueError):
+      logical_type.to_language_type(
+          schemas.ParameterizedTimestampShortRepresentation(
+              np.int64(10), np.int16(-1)))
+
+  def test_to_representation_type_guards_against_precision_loss(self):
+    # Mirrors the Java SDK's toBaseType check: a value that cannot be
+    # represented exactly at the logical type's precision is an error.
+    logical_type = schemas.ParameterizedTimestamp(6)
+    nanos_value = Timestamp(seconds=1, subseconds=123456789, precision=9)
+    with self.assertRaises(ValueError):
+      logical_type.to_representation_type(nanos_value)
+    # Lower-precision values are converted losslessly.
+    millis_value = Timestamp(seconds=1, subseconds=500, precision=3)
+    representation = logical_type.to_representation_type(millis_value)
+    self.assertEqual(representation.subseconds, 500000)
+
+
 class HypothesisTest(unittest.TestCase):
   # There is considerable variablility in runtime for this test, disable
   # deadline.
diff --git a/sdks/python/apache_beam/utils/timestamp.py 
b/sdks/python/apache_beam/utils/timestamp.py
index 3f585eecae0..2953541b42f 100644
--- a/sdks/python/apache_beam/utils/timestamp.py
+++ b/sdks/python/apache_beam/utils/timestamp.py
@@ -24,7 +24,9 @@ For internal use only; no backwards-compatibility guarantees.
 # mypy: disallow-untyped-defs
 
 import datetime
+import re
 import time
+from typing import Optional
 from typing import Union
 from typing import overload
 
@@ -41,28 +43,82 @@ TimestampTypes = Union[int, float, 'Timestamp']
 DurationTypes = Union[int, float, 'Duration']
 TimestampDurationTypes = Union[int, float, 'Duration', 'Timestamp']
 
+# Powers of ten indexed by exponent
+_POW_10 = {i: 10**i for i in range(10)}
+
 
 class Timestamp(object):
-  """Represents a Unix second timestamp with microsecond granularity.
+  """Represents a Unix second timestamp with configurable subsecond precision.
 
   Can be treated in common timestamp arithmetic operations as a numeric type.
 
-  Internally stores a time interval as an int of microseconds. This strategy
-  is necessary since floating point values lose precision when storing values,
-  especially after arithmetic operations (for example, 10000000 % 0.1 evaluates
-  to 0.0999999994448885).
+  Internally stores the timestamp as an int of floored seconds since the
+  epoch plus a non-negative int subsecond value interpreted with the
+  specified precision where (``0 <= subseconds < 10**precision``).
+
+  Integer storage is necessary since floating point
+  values lose precision when storing values, especially after arithmetic
+  operations (for example, 10000000 % 0.1 evaluates to 0.0999999994448885).
+
+  ``precision`` is the number of decimal digits used to represent the
+  fraction of a second (e.g. 3 for millis, 6 for micros, 9 for
+  nanos). Defaults to microseconds.
+  If ``seconds`` is a float, the fractional part will be captured up
+  to ``precision`` digits.
+
+  Lossy conversion operations will throw an error unless
+  ``allow_lossy_conversion=True`` is specified (e.g. see ``to_utc_datetime``).
   """
+  MICROS_PRECISION = 6
+  NANOS_PRECISION = 9
+
   def __init__(
       self,
       seconds: Union[int, float] = 0,
-      micros: Union[int, float] = 0) -> None:
+      subseconds: Union[int, float] = 0,
+      precision: int = MICROS_PRECISION,
+      *,
+      micros: Optional[Union[int, float]] = None) -> None:
     if not isinstance(seconds, (int, float)):
       raise TypeError(
           'Cannot interpret %s %s as seconds.' % (seconds, type(seconds)))
-    if not isinstance(micros, (int, float)):
+    if not isinstance(subseconds, (int, float)):
       raise TypeError(
-          'Cannot interpret %s %s as micros.' % (micros, type(micros)))
-    self.micros = int(seconds * 1000000) + int(micros)
+          'Cannot interpret %s %s as subseconds.' %
+          (subseconds, type(subseconds)))
+    if not isinstance(precision, int):
+      raise TypeError(
+          'Cannot interpret %s %s as precision.' % (precision, 
type(precision)))
+    if not 0 <= precision <= Timestamp.NANOS_PRECISION:
+      raise ValueError(
+          'Timestamp precision must be between 0 and %d (inclusive), '
+          'but was %d.' % (Timestamp.NANOS_PRECISION, precision))
+    if micros is not None:
+      if not isinstance(micros, (int, float)):
+        raise TypeError(
+            'Cannot interpret %s %s as micros.' % (micros, type(micros)))
+      if subseconds:
+        raise ValueError(
+            'micros and subseconds are mutually exclusive, got micros=%s, '
+            'subseconds=%s.' % (micros, subseconds))
+      if precision != Timestamp.MICROS_PRECISION:
+        raise ValueError(
+            'micros implies microsecond precision (6) but precision was %d; '
+            'use subseconds instead.' % precision)
+      subseconds = micros
+    self._precision = precision
+    total = int(seconds * _POW_10[precision]) + int(subseconds)
+    self._seconds, self._subseconds = divmod(total, _POW_10[precision])
+
+  def _total(self, precision: int) -> int:
+    """Returns the total time since the epoch in units of 10**-precision
+    seconds.
+
+    ``precision`` must be greater than or equal to this timestamp's
+    precision, so that scaling up is always lossless.
+    """
+    return self._seconds * _POW_10[precision] + (
+        self._subseconds * _POW_10[precision - self._precision])
 
   @staticmethod
   def of(seconds: TimestampTypes) -> 'Timestamp':
@@ -116,6 +172,10 @@ class Timestamp(object):
   def from_rfc3339(cls, rfc3339: str) -> 'Timestamp':
     """Create a ``Timestamp`` instance from an RFC 3339 compliant string.
 
+    Fractional seconds up to microseconds produce a microsecond-precision
+    Timestamp; a longer fraction (up to nanoseconds) produces a Timestamp
+    whose precision matches the number of fractional digits.
+
     .. note::
       All timezones are implicitly converted to UTC.
 
@@ -128,33 +188,108 @@ class Timestamp(object):
       raise ValueError(
           "Could not parse RFC 3339 string '{}' due to error: '{}'.".format(
               rfc3339, e))
-    return cls.from_utc_datetime(dt)
+    timestamp = cls.from_utc_datetime(dt)
+    # dateutil silently truncates fractional seconds to microseconds; parse
+    # any sub-microsecond digits ourselves to avoid losing precision.
+    fraction = re.search(r'[0-9]{2}[.,]([0-9]{7,})', rfc3339)
+    if fraction:
+      digits = fraction.group(1)
+      if len(digits) > cls.NANOS_PRECISION:
+        raise ValueError(
+            "Could not parse RFC 3339 string '%s': fractional seconds "
+            'beyond nanosecond precision are not supported.' % rfc3339)
+      precision = len(digits)
+      sub_micro = int(digits[cls.MICROS_PRECISION:])
+      return Timestamp(
+          timestamp.seconds(),
+          timestamp.subseconds() * _POW_10[precision - cls.MICROS_PRECISION] +
+          sub_micro,
+          precision)
+    return timestamp
 
   def seconds(self) -> int:
     """Returns the timestamp in seconds."""
-    return self.micros // 1000000
+    return self._seconds
+
+  def subseconds(self) -> int:
+    """Returns the fraction of a second, in units of 10**-precision seconds.
+
+    Always non-negative and less than 10**precision
+    """
+    return self._subseconds
+
+  def precision(self) -> int:
+    """Returns the precision of this Timestamp."""
+    return self._precision
+
+  @property
+  def micros(self) -> int:
+    """Returns the total number of microseconds since the epoch."""
+    if self._precision > Timestamp.MICROS_PRECISION:
+      raise ValueError(
+          '%r has greater than microsecond precision, converting it to '
+          'micros may lose precision. Use to_precision(6, '
+          'allow_lossy_conversion=True) to explicitly truncate it first, '
+          'or use nanos instead.' % self)
+    return self._total(Timestamp.MICROS_PRECISION)
+
+  @property
+  def nanos(self) -> int:
+    """Returns the total number of nanoseconds since the epoch."""
+    return self._total(Timestamp.NANOS_PRECISION)
+
+  def to_precision(
+      self,
+      precision: int,
+      allow_lossy_conversion: bool = False) -> 'Timestamp':
+    """Returns this Timestamp converted to the given precision.
+
+    Increasing precision is always lossless. Decreasing precision raises
+    ValueError if this timestamp has a non-zero component below the target
+    precision, unless allow_lossy_conversion is True, in which case the
+    timestamp is truncated (floored) to the target precision.
+    """
+    if precision == self._precision:
+      return self
+    if not 0 <= precision <= Timestamp.NANOS_PRECISION:
+      raise ValueError(
+          'Timestamp precision must be between 0 and %d (inclusive), '
+          'but was %d.' % (Timestamp.NANOS_PRECISION, precision))
+    if precision > self._precision:
+      scale = _POW_10[precision - self._precision]
+      return Timestamp(self._seconds, self._subseconds * scale, precision)
+    scale = _POW_10[self._precision - precision]
+    remainder = self._subseconds % scale
+    if remainder and not allow_lossy_conversion:
+      raise ValueError(
+          '%r cannot be represented exactly at precision %d. Set '
+          'allow_lossy_conversion=True to truncate it.' % (self, precision))
+    return Timestamp(self._seconds, self._subseconds // scale, precision)
 
   def predecessor(self) -> 'Timestamp':
-    """Returns the largest timestamp smaller than self."""
-    return Timestamp(micros=self.micros - 1)
+    """Returns the largest timestamp smaller than self, at this precision."""
+    return Timestamp(self._seconds, self._subseconds - 1, self._precision)
 
   def successor(self) -> 'Timestamp':
-    """Returns the smallest timestamp larger than self."""
-    return Timestamp(micros=self.micros + 1)
+    """Returns the smallest timestamp larger than self, at this precision."""
+    return Timestamp(self._seconds, self._subseconds + 1, self._precision)
 
   def __repr__(self) -> str:
-    micros = self.micros
+    total = self._total(self._precision)
     sign = ''
-    if micros < 0:
+    if total < 0:
       sign = '-'
-      micros = -micros
-    int_part = micros // 1000000
-    frac_part = micros % 1000000
+      total = -total
+    int_part, frac_part = divmod(total, _POW_10[self._precision])
     if frac_part:
-      return 'Timestamp(%s%d.%06d)' % (sign, int_part, frac_part)
+      return 'Timestamp(%s%d.%0*d)' % (
+          sign, int_part, self._precision, frac_part)
     return 'Timestamp(%s%d)' % (sign, int_part)
 
-  def to_utc_datetime(self, has_tz: bool = False) -> datetime.datetime:
+  def to_utc_datetime(
+      self,
+      has_tz: bool = False,
+      allow_lossy_conversion: bool = False) -> datetime.datetime:
     """Returns a ``datetime.datetime`` object of UTC for this Timestamp.
 
     Note that this method returns a ``datetime.datetime`` object without a
@@ -164,74 +299,112 @@ class Timestamp(object):
 
     Args:
       has_tz: whether the timezone info is attached, default to False.
+      allow_lossy_conversion: must be set to True to convert a timestamp
+        with precision above microseconds, since ``datetime.datetime`` only
+        supports microsecond resolution; the result is truncated (floored)
+        to whole microseconds.
 
     Returns:
       a ``datetime.datetime`` object of UTC for this Timestamp.
-    """
 
+    Raises:
+      ValueError: if this timestamp has precision above microseconds and
+        allow_lossy_conversion is not True.
+    """
+    if self._precision > Timestamp.MICROS_PRECISION:
+      if not allow_lossy_conversion:
+        raise ValueError(
+            'Converting %r to datetime truncates it to microseconds. Set '
+            'allow_lossy_conversion=True to allow this conversion.' % self)
+      micros_of_second = self._subseconds // _POW_10[self._precision -
+                                                     
Timestamp.MICROS_PRECISION]
+    else:
+      micros_of_second = self._subseconds * _POW_10[Timestamp.MICROS_PRECISION 
-
+                                                    self._precision]
     # We can't easily construct a datetime object from microseconds, so we
     # create one at the epoch and add an appropriate timedelta interval.
     epoch = self._epoch_datetime_utc()
     if not has_tz:
       epoch = epoch.replace(tzinfo=None)
-    return epoch + datetime.timedelta(microseconds=self.micros)
+    return epoch + datetime.timedelta(
+        seconds=self._seconds, microseconds=micros_of_second)
 
   def to_rfc3339(self) -> str:
-    # Append 'Z' for UTC timezone.
-    return self.to_utc_datetime().isoformat() + 'Z'
+    """Returns an RFC 3339 string for this Timestamp."""
+    if self._precision <= Timestamp.MICROS_PRECISION:
+      # Append 'Z' for UTC timezone.
+      return self.to_utc_datetime().isoformat() + 'Z'
+    # format the fractional second manually
+    whole_second_datetime = self._epoch_datetime_utc().replace(
+        tzinfo=None) + datetime.timedelta(seconds=self._seconds)
+    result = whole_second_datetime.isoformat()
+    if self._subseconds:
+      result = result + '.%0*d' % (self._precision, self._subseconds)
+
+    return result + 'Z'
 
   def to_proto(self) -> timestamp_pb2.Timestamp:
     """Returns the `google.protobuf.timestamp_pb2` representation."""
-    secs = self.micros // 1000000
-    nanos = (self.micros % 1000000) * 1000
-    return timestamp_pb2.Timestamp(seconds=secs, nanos=nanos)
+    return timestamp_pb2.Timestamp(
+        seconds=self._seconds,
+        nanos=self._subseconds *
+        _POW_10[Timestamp.NANOS_PRECISION - self._precision])
 
   @staticmethod
   def from_proto(timestamp_proto: timestamp_pb2.Timestamp) -> 'Timestamp':
     """Creates a Timestamp from a `google.protobuf.timestamp_pb2`.
 
-    Note that the google has a sub-second resolution of nanoseconds whereas 
this
-    class has a resolution of microsends. This class will truncate the
-    nanosecond resolution down to the microsecond.
+    The returned Timestamp has microsecond precision if the proto's nanos
+    are microsecond-aligned, and nanosecond precision otherwise.
     """
-
     if timestamp_proto.nanos % 1000 != 0:
-      # TODO(https://github.com/apache/beam/issues/19922): Better define
-      # timestamps.
-      raise ValueError(
-          "Cannot convert from nanoseconds to microseconds " +
-          "because this loses precision. Please make sure that " +
-          "this is the correct behavior you want and manually " +
-          "truncate the precision to the nearest microseconds. " +
-          "See [https://github.com/apache/beam/issues/19922] for " +
-          "more information.")
-
+      return Timestamp(
+          timestamp_proto.seconds,
+          timestamp_proto.nanos,
+          Timestamp.NANOS_PRECISION)
     return Timestamp(
-        seconds=timestamp_proto.seconds, micros=timestamp_proto.nanos // 1000)
+        timestamp_proto.seconds,
+        timestamp_proto.nanos // 1000,
+        Timestamp.MICROS_PRECISION)
 
   def __float__(self) -> float:
     # Note that the returned value may have lost precision.
-    return self.micros / 1000000
+    return self._total(Timestamp.NANOS_PRECISION) / 1000000000
 
   def __int__(self) -> int:
     # Note that the returned value may have lost precision.
-    return self.micros // 1000000
+    return self._seconds
 
   def __eq__(self, other: object) -> bool:
-    # Allow comparisons between Duration and Timestamp values.
-    if isinstance(other, (Duration, Timestamp)):
-      return self.micros == other.micros
+    if isinstance(other, Timestamp):
+      if self._precision == other._precision:
+        return (
+            self._seconds == other._seconds and
+            self._subseconds == other._subseconds)
+      precision = max(self._precision, other._precision)
+      return self._total(precision) == other._total(precision)
+    elif isinstance(other, Duration):
+      # Allow comparisons between Duration and Timestamp values.
+      return self._total(Timestamp.NANOS_PRECISION) == other.micros * 1000
     elif isinstance(other, (int, float)):
-      return self.micros == Timestamp.of(other).micros
+      return self == Timestamp.of(other)
     else:
       # Support equality with other types
       return NotImplemented
 
   def __lt__(self, other: TimestampDurationTypes) -> bool:
     # Allow comparisons between Duration and Timestamp values.
-    if not isinstance(other, Duration):
-      other = Timestamp.of(other)
-    return self.micros < other.micros
+    if isinstance(other, Duration):
+      return self._total(Timestamp.NANOS_PRECISION) < other.micros * 1000
+    other = Timestamp.of(other)
+    if self._seconds != other._seconds:
+      return self._seconds < other._seconds
+    if self._precision == other._precision:
+      return self._subseconds < other._subseconds
+    precision = max(self._precision, other._precision)
+    return (
+        self._subseconds * _POW_10[precision - self._precision]
+        < other._subseconds * _POW_10[precision - other._precision])
 
   def __gt__(self, other: TimestampDurationTypes) -> bool:
     return not (self < other or self == other)
@@ -243,11 +416,16 @@ class Timestamp(object):
     return not self < other
 
   def __hash__(self) -> int:
-    return hash(self.micros)
+    # Normalized to max precision
+    return hash(self._total(Timestamp.NANOS_PRECISION))
 
   def __add__(self, other: DurationTypes) -> 'Timestamp':
     other = Duration.of(other)
-    return Timestamp(micros=self.micros + other.micros)
+    precision = max(self._precision, Timestamp.MICROS_PRECISION)
+    return Timestamp(
+        subseconds=self._total(precision) +
+        other.micros * _POW_10[precision - Timestamp.MICROS_PRECISION],
+        precision=precision)
 
   def __radd__(self, other: DurationTypes) -> 'Timestamp':
     return self + other
@@ -264,13 +442,33 @@ class Timestamp(object):
       self, other: Union[DurationTypes,
                          'Timestamp']) -> Union['Timestamp', 'Duration']:
     if isinstance(other, Timestamp):
-      return Duration(micros=self.micros - other.micros)
+      diff_nanos = (
+          self._total(Timestamp.NANOS_PRECISION) -
+          other._total(Timestamp.NANOS_PRECISION))
+      if diff_nanos % 1000 != 0:
+        raise ValueError(
+            'The difference of %r and %r has sub-microsecond precision, '
+            'which Duration cannot represent. Truncate the operands with '
+            'to_precision(6, allow_lossy_conversion=True) first.' %
+            (self, other))
+      return Duration(micros=diff_nanos // 1000)
     other = Duration.of(other)
-    return Timestamp(micros=self.micros - other.micros)
+    precision = max(self._precision, Timestamp.MICROS_PRECISION)
+    return Timestamp(
+        subseconds=self._total(precision) -
+        other.micros * _POW_10[precision - Timestamp.MICROS_PRECISION],
+        precision=precision)
 
   def __mod__(self, other: DurationTypes) -> 'Duration':
     other = Duration.of(other)
-    return Duration(micros=self.micros % other.micros)
+    remainder_nanos = self._total(Timestamp.NANOS_PRECISION) % (
+        other.micros * 1000)
+    if remainder_nanos % 1000 != 0:
+      raise ValueError(
+          'The remainder of %r modulo %r has sub-microsecond precision, '
+          'which Duration cannot represent. Truncate this timestamp with '
+          'to_precision(6, allow_lossy_conversion=True) first.' % (self, 
other))
+    return Duration(micros=remainder_nanos // 1000)
 
 
 MIN_TIMESTAMP = Timestamp(
@@ -361,8 +559,10 @@ class Duration(object):
 
   def __eq__(self, other: object) -> bool:
     # Allow comparisons between Duration and Timestamp values.
-    if isinstance(other, (Duration, Timestamp)):
+    if isinstance(other, Duration):
       return self.micros == other.micros
+    elif isinstance(other, Timestamp):
+      return self.micros * 1000 == other.nanos
     elif isinstance(other, (int, float)):
       return self.micros == Duration.of(other).micros
     else:
@@ -371,8 +571,9 @@ class Duration(object):
 
   def __lt__(self, other: TimestampDurationTypes) -> bool:
     # Allow comparisons between Duration and Timestamp values.
-    if not isinstance(other, Timestamp):
-      other = Duration.of(other)
+    if isinstance(other, Timestamp):
+      return self.micros * 1000 < other.nanos
+    other = Duration.of(other)
     return self.micros < other.micros
 
   def __gt__(self, other: TimestampDurationTypes) -> bool:
@@ -385,7 +586,9 @@ class Duration(object):
     return not self < other
 
   def __hash__(self) -> int:
-    return hash(self.micros)
+    # Timestamps hash on their total nanoseconds. Hash equivalently so that
+    # a Duration and Timestamp that compare equal also hash equal.
+    return hash(self.micros * 1000)
 
   def __neg__(self) -> 'Duration':
     return Duration(micros=-self.micros)
diff --git a/sdks/python/apache_beam/utils/timestamp_test.py 
b/sdks/python/apache_beam/utils/timestamp_test.py
index f8d6cfdeafe..e1d120da471 100644
--- a/sdks/python/apache_beam/utils/timestamp_test.py
+++ b/sdks/python/apache_beam/utils/timestamp_test.py
@@ -172,11 +172,13 @@ class TimestampTest(unittest.TestCase):
     expected_ts = Timestamp(seconds=1234, micros=56)
     self.assertEqual(actual_ts, expected_ts)
 
-  def test_from_proto_fails_with_truncation(self):
-    # TODO(https://github.com/apache/beam/issues/19922): Better define
-    # timestamps.
-    with self.assertRaises(ValueError):
-      Timestamp.from_proto(timestamp_pb2.Timestamp(seconds=1234, nanos=56789))
+  def test_from_proto_with_sub_micro_nanos(self):
+    # Sub-microsecond protos produce a nanosecond-precision Timestamp
+    # instead of losing precision (or raising, as this method used to).
+    actual_ts = Timestamp.from_proto(
+        timestamp_pb2.Timestamp(seconds=1234, nanos=56789))
+    self.assertEqual(actual_ts.precision(), Timestamp.NANOS_PRECISION)
+    self.assertEqual(actual_ts.nanos, 1234 * 10**9 + 56789)
 
   def test_to_proto(self):
     ts = Timestamp(seconds=1234, micros=56)
@@ -193,6 +195,181 @@ class TimestampTest(unittest.TestCase):
         self.assertTrue(max_val >= min_val, "%s >= %s" % (max_val, min_val))
 
 
+class TimestampPrecisionTest(unittest.TestCase):
+  def test_constructor(self):
+    ts = Timestamp(seconds=1234, subseconds=123456789, precision=9)
+    self.assertEqual(ts.seconds(), 1234)
+    self.assertEqual(ts.subseconds(), 123456789)
+    self.assertEqual(ts.precision(), 9)
+    self.assertEqual(ts.nanos, 1234123456789)
+
+    # Subseconds overflowing a second carry into seconds.
+    ts = Timestamp(seconds=1, subseconds=1500, precision=3)
+    self.assertEqual(ts.seconds(), 2)
+    self.assertEqual(ts.subseconds(), 500)
+
+    # Negative timestamps floor seconds so subseconds stay non-negative.
+    ts = Timestamp(seconds=-2, subseconds=500000, precision=6)
+    self.assertEqual(ts.seconds(), -2)
+    self.assertEqual(ts.subseconds(), 500000)
+    self.assertEqual(ts, Timestamp(-1.5))
+
+  def test_constructor_validation(self):
+    with self.assertRaises(ValueError):
+      Timestamp(0, 0, precision=10)
+    with self.assertRaises(ValueError):
+      Timestamp(0, 0, precision=-1)
+    with self.assertRaises(TypeError):
+      Timestamp(0, 0, precision=6.0)  # type: ignore[arg-type]
+    # micros is an alias for subseconds at the default precision only.
+    with self.assertRaises(ValueError):
+      Timestamp(0, subseconds=1, micros=1)
+    with self.assertRaises(ValueError):
+      Timestamp(0, precision=9, micros=1)
+    self.assertEqual(Timestamp(1, micros=500000), Timestamp(1.5))
+
+  def test_default_precision_is_micros(self):
+    self.assertEqual(Timestamp(1.5).precision(), Timestamp.MICROS_PRECISION)
+    self.assertEqual(Timestamp.now().precision(), Timestamp.MICROS_PRECISION)
+
+  def test_equality_across_precisions(self):
+    self.assertEqual(
+        Timestamp(seconds=1, subseconds=500, precision=3),
+        Timestamp(seconds=1, subseconds=500000000, precision=9))
+    self.assertEqual(
+        hash(Timestamp(seconds=1, subseconds=500, precision=3)),
+        hash(Timestamp(seconds=1, subseconds=500000000, precision=9)))
+    # An equal Duration must hash equal too.
+    self.assertEqual(hash(Timestamp(micros=5)), hash(Duration(micros=5)))
+    self.assertNotEqual(
+        Timestamp(seconds=1, subseconds=500000001, precision=9),
+        Timestamp(seconds=1, subseconds=500000, precision=6))
+
+  def test_comparison(self):
+    self.assertLess(
+        Timestamp(seconds=1, subseconds=500000000, precision=9),
+        Timestamp(seconds=1, subseconds=500001, precision=6))
+    self.assertGreater(
+        Timestamp(seconds=1, subseconds=500000001, precision=9),
+        Timestamp(seconds=1, subseconds=500000, precision=6))
+    self.assertLess(Timestamp(2.1), Timestamp(3))
+    self.assertGreater(Timestamp(3), Timestamp(2.1))
+
+  def test_str(self):
+    self.assertEqual('Timestamp(1.500)', str(Timestamp(1, 500, precision=3)))
+    self.assertEqual(
+        'Timestamp(1.123456789)', str(Timestamp(1, 123456789, precision=9)))
+    self.assertEqual(
+        'Timestamp(-1.500000000)', str(Timestamp(-1.5, precision=9)))
+    self.assertEqual('Timestamp(1)', str(Timestamp(1, precision=9)))
+    self.assertEqual('Timestamp(1)', str(Timestamp(1, precision=0)))
+
+  def test_precision_conversion(self):
+    ts = Timestamp(seconds=1, subseconds=123456789, precision=9)
+    self.assertIs(ts.to_precision(9), ts)
+    # Lossless upward conversion.
+    up = Timestamp(1.5).to_precision(9)
+    self.assertEqual(up.precision(), 9)
+    self.assertEqual(up, Timestamp(1.5))
+    # Lossy downward conversion requires explicit permission and floors.
+    with self.assertRaises(ValueError):
+      ts.to_precision(6)
+    truncated = ts.to_precision(6, allow_lossy_conversion=True)
+    self.assertEqual(truncated.precision(), 6)
+    self.assertEqual(truncated.micros, 1123456)
+    # A lossless downward conversion doesn't require the flag.
+    self.assertEqual(Timestamp(1.5, precision=9).to_precision(3).precision(), 
3)
+    # Truncation of negative timestamps floors towards negative infinity.
+    negative = Timestamp(
+        seconds=-1, precision=9).predecessor().to_precision(
+            6, allow_lossy_conversion=True)
+    self.assertEqual(negative, Timestamp(-1) - Duration(micros=1))
+
+  def test_micros_guard_rail(self):
+    ts = Timestamp(seconds=1, subseconds=123456789, precision=9)
+    with self.assertRaises(ValueError):
+      _ = ts.micros
+    self.assertEqual(ts.nanos, 1123456789)
+    self.assertEqual(Timestamp(1, 500, precision=3).micros, 1500000)
+
+  def test_predecessor_successor(self):
+    ts = Timestamp(seconds=10, precision=9)
+    self.assertEqual(ts.predecessor().nanos, 10 * 10**9 - 1)
+    self.assertEqual(ts.successor().nanos, 10 * 10**9 + 1)
+    self.assertEqual(ts.predecessor().precision(), 9)
+    self.assertEqual(ts.successor().seconds(), 10)
+    self.assertEqual(ts.predecessor().seconds(), 9)
+    # Micros-precision timestamps keep their historical 1-micro step.
+    self.assertEqual(Timestamp(10).predecessor().micros, 10 * 10**6 - 1)
+
+  def test_to_utc_datetime_guard_rail(self):
+    ts = Timestamp(seconds=1234, subseconds=123456789, precision=9)
+    with self.assertRaises(ValueError):
+      ts.to_utc_datetime()
+    dt = ts.to_utc_datetime(allow_lossy_conversion=True)
+    self.assertEqual(dt.microsecond, 123456)
+    # Timestamps at or below microsecond precision convert freely.
+    Timestamp(1.5).to_utc_datetime()
+
+  def test_to_rfc3339_is_lossless(self):
+    ts = Timestamp(seconds=1458343379, subseconds=123456789, precision=9)
+    self.assertEqual(ts.to_rfc3339(), '2016-03-18T23:22:59.123456789Z')
+    self.assertEqual(Timestamp.from_rfc3339(ts.to_rfc3339()), ts)
+    # Whole seconds have no fractional digits, as before.
+    self.assertEqual(
+        Timestamp(1458343379, precision=9).to_rfc3339(), 
'2016-03-18T23:22:59Z')
+
+  def test_from_rfc3339_with_nanos(self):
+    ts = Timestamp.from_rfc3339('2016-03-18T23:22:59.123456789Z')
+    self.assertEqual(ts.precision(), 9)
+    self.assertEqual(ts.subseconds(), 123456789)
+    # Precision matches the number of fractional digits (above 6).
+    ts = Timestamp.from_rfc3339('2016-03-18T23:22:59.1234567Z')
+    self.assertEqual(ts.precision(), 7)
+    self.assertEqual(ts.subseconds(), 1234567)
+    self.assertEqual(Timestamp.from_rfc3339(ts.to_rfc3339()), ts)
+    self.assertEqual(
+        Timestamp.from_rfc3339('2016-03-18T23:22:59.123Z').precision(), 6)
+    with self.assertRaises(ValueError):
+      Timestamp.from_rfc3339('2016-03-18T23:22:59.1234567891Z')
+
+  def test_from_rfc3339_with_comma_separator(self):
+    # ISO 8601 also allows ',' as the decimal separator.
+    ts = Timestamp.from_rfc3339('2016-03-18T23:22:59,123456789Z')
+    self.assertEqual(
+        ts, Timestamp.from_rfc3339('2016-03-18T23:22:59.123456789Z'))
+    self.assertEqual(ts.subseconds(), 123456789)
+
+  def test_proto_round_trip_with_nanos(self):
+    ts = Timestamp(seconds=1234, subseconds=123456789, precision=9)
+    self.assertEqual(Timestamp.from_proto(ts.to_proto()), ts)
+    # Negative timestamps use non-negative proto nanos.
+    ts = Timestamp(-1.5, precision=9).predecessor()
+    proto = ts.to_proto()
+    self.assertEqual(proto.seconds, -2)
+    self.assertEqual(proto.nanos, 499999999)
+    self.assertEqual(Timestamp.from_proto(proto), ts)
+
+  def test_arithmetic_preserves_precision(self):
+    ts = Timestamp(seconds=1, subseconds=123456789, precision=9)
+    self.assertEqual((ts + 1).nanos, ts.nanos + 10**9)
+    self.assertEqual((ts + 1).precision(), 9)
+    self.assertEqual((ts - Duration(micros=1)).nanos, ts.nanos - 1000)
+    # Sub-micros precision arithmetic results are widened to micros.
+    self.assertEqual((Timestamp(1, 500, precision=3) + 0.5).precision(), 6)
+
+  def test_duration_arithmetic_guard_rail(self):
+    ts = Timestamp(seconds=1, subseconds=123456789, precision=9)
+    # Exact differences are representable as a (micros) Duration.
+    self.assertEqual(ts - (ts - Duration(micros=3)), Duration(micros=3))
+    self.assertEqual(ts + Duration(micros=5) - ts, Duration(micros=5))
+    # Sub-microsecond differences are not.
+    with self.assertRaises(ValueError):
+      _ = ts - ts.predecessor()
+    with self.assertRaises(ValueError):
+      _ = ts % Duration(seconds=1)
+
+
 class DurationTest(unittest.TestCase):
   def test_of(self):
     interval = Duration(123)

Reply via email to