tchivs opened a new pull request, #4530:
URL: https://github.com/apache/flink-cdc/pull/4530

   ## What is the purpose of this pull request?
   
   Fixes [FLINK-40594](https://issues.apache.org/jira/browse/FLINK-40594): the 
pipeline runtime stores TIME as millisecond-of-day, so every `TIME(p)` with `p 
> 3` is silently truncated even though `TimeType` accepts a precision of up to 
9 and `DebeziumEventDeserializationSchema` already builds values through 
`fromMicroOfDay` / `fromNanoOfDay`.
   
   Beyond the lost digits, `TimeData.equals`, `hashCode` and `compareTo` 
compared `millisOfDay`, so two `TIME(6)` values differing only in microseconds 
were indistinguishable and collapsed during deduplication, ordering and key 
comparison.
   
   ## Brief change log
   
   - `TimeData` stores nanosecond-of-day and exposes `toMicroOfDay()` / 
`toNanoOfDay()`; equality, hash and comparison now use that value. 
`toMillisOfDay()` keeps its signature and truncating behaviour, so existing 
callers stay source and binary compatible.
   - `TimeDataSerializer` becomes precision-aware, following the existing shape 
of `TimestampDataSerializer` / `LocalZonedTimestampDataSerializer` instead of 
introducing a new pattern:
     - `precision <= 3` keeps the historical four-byte millisecond encoding and 
resolves as `isCompatibleAsIs`, so existing state keeps working;
     - `precision > 3` uses an eight-byte nanosecond-of-day encoding and 
resolves as `isCompatibleAfterMigration`, through a versioned 
`TypeSerializerSnapshot` that still reads the legacy snapshot envelope and 
four-byte payload.
   - The binary readers/writers (`BinaryRecordData`, `BinaryArrayData`, 
`AbstractBinaryWriter`, `BinaryArrayWriter`), the converters and 
`InternalSerializers` thread the declared precision through.
   - Two existing expectations that pinned the truncation are corrected rather 
than relaxed:
     - the `TIME(6)` / `TIME(9)` transform expectation now asserts the retained 
fraction instead of `21:48:25.123`;
     - the temporal-function check compares `LOCALTIME` at whole-second 
granularity, because `LOCALTIME` and `CURRENT_TIME` are `TIME(0)` and 
previously carried a fraction their declared type does not have.
   
   ---
   
   ## Verifying this change
   
   This change added tests and can be verified as follows:
   
   - *Added unit tests in 
`flink-cdc-common/src/test/java/org/apache/flink/cdc/common/data/TimeDataTest.java`*
 covering nanosecond retention and the equality/ordering behaviour.
   - *Updated unit tests in 
`flink-cdc-runtime/.../serializer/data/TimeDataSerializerTest.java`*: it 
becomes an abstract base with concrete `TimeDataSerializer3Test` / `6Test` / 
`9Test`, plus a `TimeDataSerializerCompatibilityTest` asserting 
`isCompatibleAsIs` for the millisecond encoding, `isCompatibleAfterMigration` 
for the widened one, and successful reading of the legacy snapshot envelope and 
four-byte payload.
   - *Updated unit tests in 
`flink-cdc-common/.../converter/InternalObjectConverterTest.java` and 
`JavaObjectConverterTest.java`.*
   - *Updated integration tests in 
`flink-cdc-composer/.../FlinkPipelineTransformITCase.java` and 
`flink-cdc-pipeline-connector-postgres/.../PostgresFullTypesITCase.java`.* The 
Postgres case gains a direct microsecond-of-day assertion, because the existing 
`TIME(6)` expectation compared two equally truncated values and therefore 
passed vacuously.
   
   Local results on top of `9f23c0356`:
   
   | Module | Result |
   |---|---|
   | `flink-cdc-common` | 46 tests, 0 failures |
   | `flink-cdc-runtime` | 992 tests, 0 failures |
   | `flink-cdc-composer` (`FlinkPipelineTransformITCase`) | 71 tests, 0 
failures |
   | Spotless | passes |
   
   Two notes for reviewers:
   
   - Existing `TimeDataSerializerTest` data included 
`TimeData.fromNanoOfDay(102400)`, `204800` and `409600`. Under millisecond 
storage all three collapsed to zero at construction, so the previous round-trip 
assertion passed while the intended distinct values never existed; the 
precision-parameterised tests replace it.
   - These runs used JDK 17. I do not have the JDK 11 baseline installed 
locally, so the baseline is covered by CI rather than by a local run.
   
   ## Documentation
   
   - Does this pull request introduce a new feature? (no)
   - If yes, how is the feature documented? (not applicable)
   
   ---
   
   ##### Was generative AI tooling used to co-author this PR?
   
   - [X] Yes (please specify the tool below)
   
   Generated-by: Claude Code (Anthropic)
   
   The change was prepared, verified and described with AI assistance; I have 
reviewed every hunk, ran the suites listed above, and take responsibility for 
the correctness of the change.


-- 
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]

Reply via email to