nankeChen75 commented on issue #9788:
URL: https://github.com/apache/seatunnel/issues/9788#issuecomment-5970492737

   Thanks @DanielLeens. I completed the requested boundary controls and was 
able to isolate the failure to the Doris-source Arrow -> `SeaTunnelRow` 
conversion.
   
   ### Environment
   
   I reproduced this on current `dev` at:
   
   ```text
   SeaTunnel commit: bd09f6d9ae97de7230a6659962bb06f0b8eef10e
   Doris FE/BE:      doris-2.1.8-rc01-52396fbc88
   Java:             OpenJDK 17.0.20
   Doris system TZ:  Etc/UTC
   Doris session TZ: Etc/UTC
   Worker JVM TZ:    tested with both Etc/UTC and Asia/Shanghai
   ```
   
   The failure reproduces with both worker time zones. Changing the worker 
timezone changes the displayed value by 8 hours, but the large-year corruption 
remains.
   
   ---
   
   ### 1. Source DDL and exact input
   
   The source column is Doris `DATETIMEV2(6)`:
   
   ```sql
   CREATE TABLE `dt_source` (
     `id` bigint NULL,
     `event_time` datetime(6) NULL
   ) ENGINE=OLAP
   DUPLICATE KEY(`id`)
   DISTRIBUTED BY HASH(`id`) BUCKETS 1;
   ```
   
   Exact source rows:
   
   ```text
   1  2025-08-28 12:34:56.123456
   2  2022-01-01 00:00:00.000001
   ```
   
   I also created a source-side DATE/DATETIMEV2 control:
   
   ```sql
   CREATE TABLE `dt_source_control` (
     `id` bigint NOT NULL,
     `event_date` date NULL,
     `event_time` datetime(6) NULL
   ) ENGINE=OLAP
   DUPLICATE KEY(`id`)
   DISTRIBUTED BY HASH(`id`) BUCKETS 1;
   ```
   
   Baseline in Doris:
   
   ```text
   1  2025-08-28  2025-08-28 12:34:56.123456
   2  2022-01-01  2022-01-01 00:00:00.000001
   ```
   
   ---
   
   ### 2. Control A: Doris Source -> Console
   
   Configuration:
   
   ```hocon
   env {
     parallelism = 1
     job.mode = "BATCH"
   }
   
   source {
     Doris {
       fenodes = "172.30.0.10:8030"
       username = "root"
       password = ""
       database = "repro9788"
       table = "dt_source"
       doris.read.field = "id,event_time"
     }
   }
   
   sink {
     Console {
       log.print.data = true
     }
   }
   ```
   
   With the worker JVM explicitly set to `Etc/UTC`, the Console output is:
   
   ```text
   output rowType: id<BIGINT>, event_time<TIMESTAMP>
   
   1, +57626-09-11T22:15:23.456
   2, +53970-02-26T16:00:00.001
   ```
   
   The job itself finishes successfully:
   
   ```text
   Total Read Count:  2
   Total Write Count: 2
   Total Failed Count: 0
   ```
   
   With the worker JVM set to `Asia/Shanghai`, the first value becomes:
   
   ```text
   +57626-09-12T06:15:23.456
   ```
   
   so the timezone only changes the expected 8-hour offset; it does not explain 
the invalid year.
   
   The DATE/DATETIMEV2 source control gives:
   
   ```text
   1, 2025-08-28, +57626-09-11T22:15:23.456
   2, 2022-01-01, +53970-02-26T16:00:00.001
   ```
   
   Therefore the row and the `DATE` field are decoded correctly; only the 
`DATETIMEV2(6)` field is corrupted.
   
   This also reproduces without any Doris sink or target table involved.
   
   ---
   
   ### 3. Control B: fixed SeaTunnel row -> Doris Sink
   
   For the reverse boundary control I used `FakeSource` to construct known 
SeaTunnel `DATE` and `TIMESTAMP` values:
   
   ```hocon
   source {
     FakeSource {
       schema = {
         fields {
           id = bigint
           event_date = date
           event_time = timestamp
         }
       }
   
       rows = [
         {
           kind = INSERT
           fields = [
             1,
             "2025-08-28",
             "2025-08-28 12:34:56.123456"
           ]
         },
         {
           kind = INSERT
           fields = [
             2,
             "2022-01-01",
             "2022-01-01 00:00:00.000001"
           ]
         }
       ]
     }
   }
   ```
   
   The Doris target has the corresponding `DATE` and `DATETIMEV2(6)` columns:
   
   ```sql
   CREATE TABLE `dt_sink_control` (
     `id` bigint NOT NULL,
     `event_date` date NULL,
     `event_time` datetime(6) NULL
   ) ENGINE=OLAP
   DUPLICATE KEY(`id`)
   DISTRIBUTED BY HASH(`id`) BUCKETS 1;
   ```
   
   Sink configuration:
   
   ```hocon
   sink {
     Doris {
       fenodes = "172.30.0.10:8030"
       username = "root"
       password = ""
       database = "repro9788"
       table = "dt_sink_control"
   
       sink.enable-2pc = false
       sink.label-prefix = "repro9788_sink_control"
       doris.batch.size = 1
   
       doris.config = {
         format = "json"
         read_json_by_line = "true"
       }
     }
   }
   ```
   
   Both Stream Loads report:
   
   ```text
   Status: Success
   NumberTotalRows: 1
   NumberLoadedRows: 1
   NumberFilteredRows: 0
   ```
   
   The exact stored values are:
   
   ```text
   1  2025-08-28  2025-08-28 12:34:56.123456
   2  2022-01-01  2022-01-01 00:00:00.000001
   ```
   
   So a known-correct SeaTunnel `LocalDateTime` passes through the Doris sink 
and is stored as `DATETIMEV2(6)` without losing or changing the six-digit 
fractional value.
   
   This excludes the normal Doris sink serialization / target-side 
interpretation path for valid `LocalDateTime` values.
   
   ---
   
   ### 4. Arrow -> SeaTunnelRow boundary
   
   I then traced the failing source reader directly.
   
   At runtime the Doris source provides:
   
   ```text
   Arrow type:   Timestamp(MICROSECOND, +08:00)
   MinorType:    TIMESTAMPMICROTZ
   Vector class: TimeStampMicroTZVector
   ```
   
   For the first row:
   
   ```text
   fieldVector.getObject() = 1756355696123456
   Java type               = Long
   ```
   
   The resulting SeaTunnel field is:
   
   ```text
   +57626-09-11T22:15:23.456
   ```
   
   The relevant path is:
   
   ```text
   DorisSourceReader.pollNext
     -> DorisValueReader.hasNext
     -> ArrowToSeatunnelRowReader.readArrow
     -> ArrowToSeatunnelRowReader.convertSeatunnelRow
     -> ArrowToSeatunnelRowReader.convertArrowData
     -> ArrowToSeatunnelRowReader.convertSeatunnelRowValue
   ```
   
   `TimeStampMicroConverter.support()` currently handles:
   
   ```text
   TIMESTAMPMICRO
   ```
   
   but not:
   
   ```text
   TIMESTAMPMICROTZ
   ```
   
   Therefore `TIMESTAMPMICROTZ` falls back to `DefaultConverter`, whose 
`fieldVector.getObject()` returns the epoch value as a microsecond `Long`.
   
   Later, the timestamp conversion currently distinguishes the second-based 
timestamp types and otherwise uses milliseconds:
   
   ```java
   if (Types.MinorType.TIMESTAMPSEC == minorType
           || Types.MinorType.TIMESTAMPSECTZ == minorType) {
       return Instant.ofEpochSecond((Long) fieldValue)
               .atZone(ZoneId.systemDefault())
               .toLocalDateTime();
   } else {
       return Instant.ofEpochMilli((Long) fieldValue)
               .atZone(ZoneId.systemDefault())
               .toLocalDateTime();
   }
   ```
   
   For the exact runtime value:
   
   ```text
   1756355696123456 interpreted as microseconds
     -> 2025-08-28T04:34:56.123456Z
   
   1756355696123456 interpreted as milliseconds
     -> +57626-09-11T22:15:23.456Z
   ```
   
   The latter exactly matches the value observed in the UTC Console run.
   
   The second input provides the same unit evidence:
   
   ```text
   source: 2022-01-01 00:00:00.000001
   result: +53970-02-26T16:00:00.001
   ```
   
   i.e. `.000001` becomes `.001` when the microsecond epoch value is 
interpreted as milliseconds.
   
   I also added a focused `TimeStampMicroTZVector` reader regression test. 
Before a fix, it fails with the same large-year symptom.
   
   ---
   
   ### 5. Full Doris -> SeaTunnel -> Doris result
   
   For completeness I also ran the original round-trip shape with matching 
source and target `DATETIMEV2(6)` columns.
   
   The SeaTunnel job finishes successfully and both one-row Stream Loads report:
   
   ```text
   Status: Success
   NumberLoadedRows: 1
   NumberFilteredRows: 0
   ```
   
   However, the exact target rows are:
   
   ```text
   1  NULL
   2  NULL
   ```
   
   The Source -> Console control proves the values are already corrupted before 
they reach the sink, while the fixed-row -> Doris control proves that valid 
`LocalDateTime` values are preserved by the sink.
   
   So the `NULL`s in the full round trip are a downstream consequence of the 
already-invalid large-year values, rather than the original corruption point.
   
   ---
   
   Based on these controls, the failing boundary appears to be specifically:
   
   ```text
   TimeStampMicroTZVector / TIMESTAMPMICROTZ
     -> no microsecond-aware converter match
     -> DefaultConverter returns microsecond Long
     -> timestamp value is interpreted with Instant.ofEpochMilli()
     -> corrupted LocalDateTime
   ```
   
   This keeps the existing `DATETIMEV2` local-date-time behavior intact outside 
the demonstrated failure path.
   
   For a minimal fix, I am considering adding explicit `TIMESTAMPMICROTZ` 
support to the microsecond timestamp conversion path, together with the focused 
regression test, rather than changing the shared fallback conversion behavior.
   
   Would that be the preferred scope for the fix?


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