JazzFan opened a new issue, #40110:
URL: https://github.com/apache/beam/issues/40110

   ### What happened?
   
   ## What happened?
   
   **Version:** Observed on Beam 2.75.0 (Python SDK pipeline using 
`WriteToBigQuery(method=STORAGE_WRITE_API)`, Java expansion service 2.75.0, 
Dataflow runner).
   
   **Affected releases:** 2.70.0 through 2.76.0 (current latest). 2.69.0 is the 
last working release. The sweep below runs 2.70.0, 2.71.0, 2.74.0, 2.75.0 and 
2.76.0; 2.72.0 and 2.73.0 were not run but carry both offending call sites in 
source. The two offending call sites are unchanged on `master` at the time of 
filing (`TableRowToStorageApiProto.java` lines 1993 and 2003).
   
   When BigQuery rejects individual rows in an AppendRows request 
(`AppendSerializationError`), 
`StorageApiWritesShardedRecords.handleAppendFailure` converts the rejected 
protos back to `TableRow`s and emits them on the failed-rows output. In the 
cross-language `BigQueryStorageWriteApiSchemaTransformProvider`, that output is 
converted to a Beam `Row` via `BigQueryUtils.toBeamRow`. For any table with a 
TIMESTAMP column this conversion throws:
   
   ```
   org.apache.beam.sdk.util.UserCodeException: 
java.time.format.DateTimeParseException: Text '2026-09-02 T18:51:43.417' could 
not be parsed at index 10
       at java.time.Instant.parse(Instant.java:397)
       at 
org.apache.beam.sdk.io.gcp.bigquery.BigQueryUtils.toBeamValue(BigQueryUtils.java:932)
       at 
org.apache.beam.sdk.io.gcp.bigquery.BigQueryUtils.toBeamRow(BigQueryUtils.java:859)
       at 
org.apache.beam.sdk.io.gcp.bigquery.providers.BigQueryStorageWriteApiSchemaTransformProvider$BigQueryStorageWriteApiSchemaTransform.lambda$expand$328833dc$1(BigQueryStorageWriteApiSchemaTransformProvider.java:256)
       at 
org.apache.beam.sdk.transforms.MapElements$2.processElement(MapElements.java:151)
       ...
       at 
org.apache.beam.sdk.io.gcp.bigquery.StorageApiWritesShardedRecords$WriteRecordsDoFn.handleAppendFailure(StorageApiWritesShardedRecords.java:593)
       at 
org.apache.beam.sdk.io.gcp.bigquery.RetryManager.run(RetryManager.java:295)
       at 
org.apache.beam.sdk.io.gcp.bigquery.StorageApiWritesShardedRecords$WriteRecordsDoFn.process(StorageApiWritesShardedRecords.java:1069)
   ```
   
   Note the value: `2026-09-02 T18:51:43.417` — a space **and** a `T` between 
date and time, and no zone. (The trace is from the original production failure; 
the reproductions below use `2026-09-03`, which is immaterial to the behaviour.)
   
   ### Root cause
   
   **Producer.** `TableRowToStorageApiProto.tableRowFromMessage`, `case 
TIMESTAMP`, at 
[#L1993](https://github.com/apache/beam/blob/v2.76.0/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/TableRowToStorageApiProto.java#L1993)
 and 
[#L2003](https://github.com/apache/beam/blob/v2.76.0/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/TableRowToStorageApiProto.java#L2003):
   
   ```java
   return LocalDateTime.ofInstant(instant, 
ZoneOffset.UTC).format(TIMESTAMP_FORMATTER);
   ```
   
   `BigQueryUtils.TIMESTAMP_FORMATTER` 
([#L217](https://github.com/apache/beam/blob/v2.76.0/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryUtils.java#L217))
 is a *parsing* formatter. It accepts either date/time separator by declaring 
two optional sections that each contain only a literal:
   
   ```java
   .optionalStart().appendLiteral(' ').optionalEnd()
   .optionalStart().appendLiteral('T').optionalEnd()
   ```
   
   That works for parsing. But `java.time` *prints* an optional section 
whenever all of its fields are available, and a literal-only section has no 
fields — so both separators are always printed. The zone/offset sections are 
skipped because a `LocalDateTime` has none. Result: `2026-09-02 T18:51:43.417`.
   
   **Consumer.** `BigQueryUtils.toBeamValue`, `SqlTypes.TIMESTAMP` branch, 
[#L932](https://github.com/apache/beam/blob/v2.76.0/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryUtils.java#L932),
 calls `java.time.Instant.parse`, which requires strict ISO-8601. It fails at 
index 10.
   
   The string is not usable by any consumer in the codebase. Beam's own round 
trip (`Instant.from(TIMESTAMP_FORMATTER.parse(s))`, as used at 
`TableRowToStorageApiProto.java#L392`) parses it but cannot resolve it to an 
instant, since it carries neither offset nor zone:
   
   ```
   java.time.DateTimeException: Unable to obtain Instant from TemporalAccessor:
       {},ISO resolved to 2026-09-03T18:51:43.417 of type 
java.time.format.Parsed
   ```
   
   **A producer-only fix is not sufficient.** `toBeamValue` handles TIMESTAMP 
in two places — `SqlTypes.TIMESTAMP` via `Instant.parse` and 
`TypeName.DATETIME` via `JSON_VALUE_PARSERS` — and no string satisfies both. 
Measured against 2.76.0:
   
   | `TableRow` string | `SqlTypes.TIMESTAMP` | `TypeName.DATETIME` |
   | --- | --- | --- |
   | `2026-09-03 18:51:43.417 UTC` (BigQuery canonical) | 
`DateTimeParseException` | OK |
   | `2026-09-03T18:51:43.417Z` (ISO-8601) | OK | `NumberFormatException` |
   | `2026-09-03 T18:51:43.417` (current output) | `DateTimeParseException` | 
`DateTimeParseException` |
   
   Whichever format the producer emits, one consumer branch rejects it, so a 
consumer-side change is needed too.
   
   ### Regression
   
   Introduced by #36425 ("Fix issues in tableRowFromMessage", merged 
2025-10-31, commit d46a013), first shipped in **2.70.0**. It is the only commit 
touching `TableRowToStorageApiProto.java` between `v2.69.0` and `v2.70.0`, and 
it added both `format(TIMESTAMP_FORMATTER)` call sites.
   
   Before it, `jsonValueFromMessageValue` switched on the *proto* field type 
rather than the BigQuery logical type, so an int64 TIMESTAMP fell through to 
the `INT64`/`default` branch and the raw epoch-micros integer was copied into 
the `TableRow`. The `Long.parseLong` branch of `toBeamValue` handles that, so 
the round trip worked.
   
   Verified by round-tripping `TableRow` → `messageFromTableRow` → 
`tableRowFromMessage` → `toBeamRow` against released artifacts (JDK 17), schema 
`id STRING, ts TIMESTAMP`:
   
   | Beam | `tableRowFromMessage` produced | `toBeamRow` |
   | --- | --- | --- |
   | 2.68.0, 2.69.0 | `1788461503417000` | OK |
   | 2.70.0, 2.71.0, 2.74.0, 2.75.0, 2.76.0 | `2026-09-03 T18:51:43.417` | 
`DateTimeParseException` at index 10 |
   
   2.69.0 is the last working release for `SqlTypes.TIMESTAMP` consumers, which 
is the path the schema transform uses. (The pre-2.70 epoch-micros output was 
separately mishandled by the `TypeName.DATETIME` branch, whose contract is 
epoch *seconds*. That is moot from 2.70 onward and needs no fix.)
   
   ### The malformed format was absorbed into the test suite
   
   Three separate places, all from #36425. This matters practically: correcting 
the producer turns two dozen existing tests red until each is addressed.
   
   **1. The expected value is computed with the production expression.** 
`TableRowToStorageApiProtoTest`'s `normalizeSingularField` helper:
   
   ```java
   case TIMESTAMP:
     ...
     return LocalDateTime.ofInstant(instant, 
ZoneOffset.UTC).format(TIMESTAMP_FORMATTER);
   ```
   
   The assertion compares the output against itself, so it passes for any 
output the formatter produces. On an unmodified `v2.76.0` checkout all 25 tests 
in that class pass, including `testTableRowFromMessageNoF`, 
`testTableRowFromMessageWithF` and `testMessageFromTableRow`, which each carry 
TIMESTAMP fields through this path. `TableRowToStorageApiProtoIT` does not 
close the gap either — it compares via `SELECT FORMAT_TIMESTAMP(...)` in 
BigQuery rather than round-tripping through `toBeamRow`.
   
   **2. The malformed string appears as a literal expected value.** 
`"1970-01-01 T00:00:00.000043"` at 
[`BigQueryIOWriteTest.java#L3879`](https://github.com/apache/beam/blob/v2.76.0/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryIOWriteTest.java#L3879),
 
[#L4430](https://github.com/apache/beam/blob/v2.76.0/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryIOWriteTest.java#L4430),
 
[#L4588](https://github.com/apache/beam/blob/v2.76.0/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryIOWriteTest.java#L4588)
 and 
[#L4718](https://github.com/apache/beam/blob/v2.76.0/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryIOWriteTest.java#L4718),
 used as an input row and then compared against round-tripped output. 
`TIMESTAMP_FORMATTER` parses it, since both separators are optional on the 
parse side, so the assertions pass. #3
 6425 uses the correct `1970-01-01T00:00:00.000043` in two other tests.
   
   **3. The discrepancy is encoded as a per-path difference.** 
[`BigQueryIOWriteTest.java#L1338-L1351`](https://github.com/apache/beam/blob/v2.76.0/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryIOWriteTest.java#L1338-L1351):
   
   ```java
   .set(
       "instantval",
       useStorageApi || useStorageApiApproximate
           ? "2019-01-01 T00:00:00"
           : "2019-01-01 00:00:00 UTC"),
   ```
   
   The non-Storage-API write paths produce BigQuery's canonical form; the 
Storage Write API path produces the malformed one, and the assertion branches 
on that. Once the producer is fixed both branches agree and the conditional 
collapses.
   
   Correcting the producer locally on a `v2.76.0` checkout and running 
`org.apache.beam.sdk.io.gcp.bigquery.*` (973 tests) shows how much has to move 
with it:
   
   | State | Failures |
   | --- | --- |
   | Unmodified baseline | 1 (`testReadTransformProtoTranslation`, an unrelated 
live-BigQuery test) |
   | Production fix only | 25 (24 caused by the change, plus the baseline 
failure) |
   | \+ four literals corrected | 7 |
   | \+ `runTestWriteAvro` conditional collapsed | 0 |
   
   ### Reproduction
   
   The formatter behaviour reproduces on a plain JDK with no Beam dependency:
   
   ```java
   import java.time.*;
   import java.time.format.*;
   
   public class Repro {
     public static void main(String[] a) {
       // Copied from BigQueryUtils.DATETIME_SPACE_FORMATTER / 
TIMESTAMP_FORMATTER (v2.76.0)
       DateTimeFormatter dateTimeSpace = new DateTimeFormatterBuilder()
           .append(DateTimeFormatter.ISO_LOCAL_DATE)
           .optionalStart().appendLiteral(' ').optionalEnd()
           .optionalStart().appendLiteral('T').optionalEnd()
           .append(DateTimeFormatter.ISO_LOCAL_TIME)
           .toFormatter().withZone(ZoneOffset.UTC);
       DateTimeFormatter timestamp = new DateTimeFormatterBuilder()
           .append(dateTimeSpace)
           .optionalStart().appendOffsetId().optionalEnd()
           .optionalStart().appendOffset("+HH:mm", "+00:00").optionalEnd()
           .optionalStart().appendLiteral(' 
').parseCaseSensitive().appendZoneRegionId().optionalEnd()
           .toFormatter();
   
       Instant in = Instant.ofEpochSecond(1788461503L, 417_000_000L);
       String s = LocalDateTime.ofInstant(in, ZoneOffset.UTC).format(timestamp);
       System.out.println("formatted: '" + s + "'");      // '2026-09-03 
T18:51:43.417'
       Instant.parse(s);                                  // 
DateTimeParseException at index 10
     }
   }
   ```
   
   The version sweep in the table above was produced by round-tripping through 
the released `beam-sdks-java-io-google-cloud-platform` artifacts for each 
version; that harness can be supplied if useful.
   
   ### Impact
   
   The single producer bug surfaces two ways, depending on how the failed-rows 
output is consumed.
   
   **A. Cross-language / schema-transform consumers — hard failure.** Every 
Python and YAML `WriteToBigQuery` with `STORAGE_WRITE_API`, where 
`BigQueryStorageWriteApiSchemaTransformProvider` maps failed rows through 
`toBeamRow`. For any table with a TIMESTAMP column that conversion throws, so 
no row reaches the dead-letter output.
   
   BigQuery's row-level diagnostic is lost with it. `handleAppendFailure` does 
propagate the message, and the provider does put it on the output row — but 
`error_message` and `failed_row` are populated in the same `Row.withFieldValue` 
chain 
([#L250-L257](https://github.com/apache/beam/blob/v2.76.0/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/providers/BigQueryStorageWriteApiSchemaTransformProvider.java#L250-L257)),
 so the throw destroys the element that would have explained the rejection.
   
   Because the conversion is fused into `WriteRecordsDoFn.process`, the 
exception fails the whole bundle rather than the individual row. On streaming 
runners the bundle is retried, the same rows are rejected again, and the 
conversion throws again. In the Dataflow pipeline where this was observed the 
affected shard stopped making progress — one deterministically rejected row (a 
null in a REQUIRED column) blocked the pipeline instead of being dead-lettered. 
*That last part is a production observation; the reproductions cover the 
conversion failure only.*
   
   **B. Java consumers of `getFailedStorageApiInserts()` — silent corruption.** 
No exception. The user receives a `TableRow` whose TIMESTAMP field holds 
`2026-09-03 T18:51:43.417`, which no Beam parser, no BigQuery parser and no 
standard ISO-8601 parser accepts. Re-inserting dead-lettered rows, persisting 
them, or parsing them downstream yields corrupt data or a failure far from the 
cause.
   
   For tables with a TIMESTAMP column the dead-letter output is non-functional: 
mode A delivers no rows at all, mode B delivers rows whose timestamp no parser 
accepts. Both require the failed-rows path to be exercised, which is an error 
condition rather than the normal write path, so the blast radius is bounded — 
but for a pipeline that relies on dead-lettering to make progress, the feature 
simply does not work.
   
   ### Constraints on a fix
   
   Recording what was measured, not proposing a design — the choice below is 
the maintainers' to make.
   
   **The producer needs a printing formatter.** `TIMESTAMP_FORMATTER` is built 
for parsing and cannot
   be used to print, for the reason given under *Root cause*.
   
   **A producer-only change will not fix it.** The two consumer branches accept 
disjoint formats (see
   the table under *Root cause*), so whichever format the producer emits, one 
branch has to be adapted
   to it. There are two candidate targets, and they push the required change to 
opposite sides:
   
   - **BigQuery canonical** (`yyyy-MM-dd HH:mm:ss[.ffffff] UTC`) is already 
accepted by
     `TypeName.DATETIME`, is what the BigQuery REST API returns, and is what 
Beam's non-Storage-API
     write paths already emit — the conditional at 
`BigQueryIOWriteTest.java#L1338-L1351` above shows
     that difference explicitly. It would require widening 
`SqlTypes.TIMESTAMP`, which currently uses
     strict `Instant.parse`.
   - **ISO-8601** is already accepted by `SqlTypes.TIMESTAMP`. It would require 
changing
     `TypeName.DATETIME`, whose non-`UTC` fallback is `Double.parseDouble` over 
epoch seconds, and
     would leave the Storage Write API path emitting a different format from 
Beam's other write paths.
   
   **If canonical form is chosen, fractional precision has to cap at six 
digits.**
   `BIGQUERY_TIMESTAMP_PARSER` is built with `appendFractionOfSecond(1, 6)`, 
and it is what the
   `TypeName.DATETIME` branch uses for strings ending in `UTC`. The second 
producer call site (line
   2003) reads `nanos` straight off a `google.protobuf.Timestamp`, which can 
carry sub-microsecond
   values, so a nine-digit fraction arrives as an `IllegalArgumentException`:
   
   | `TableRow` string | `TypeName.DATETIME` |
   | --- | --- |
   | `2026-09-03 18:51:43.417123 UTC` | OK |
   | `2026-09-03 18:51:43.417123456 UTC` | `IllegalArgumentException` |
   
   The int64-micros call site (line 1993) is inherently micro-aligned and 
unaffected. This particular
   cap does not apply to the ISO-8601 option, since `Instant.parse` accepts 
nine digits and
   `BIGQUERY_TIMESTAMP_PARSER` would not see the value — though BigQuery 
TIMESTAMP is microsecond
   precision either way, so truncating at the producer is worth considering 
regardless.
   
   **The test corrections above are part of the work.** All three places need 
addressing, and the
   regression test that replaces them needs its expected value written as a 
literal rather than
   computed with the production formatter — otherwise it reproduces the same 
blind spot.
   
   cc @reuvenlax (author of #36425)
   
   
   ### Issue Priority
   
   Priority: 2 (default / most bugs should be filed as P2)
   
   ### Issue Components
   
   - [x] Component: Python SDK
   - [x] Component: Java SDK
   - [ ] Component: Go SDK
   - [ ] Component: Typescript SDK
   - [x] Component: IO connector
   - [ ] Component: Beam YAML
   - [ ] Component: Beam examples
   - [ ] Component: Beam playground
   - [ ] Component: Beam katas
   - [ ] Component: Website
   - [ ] Component: Infrastructure
   - [ ] Component: Spark Runner
   - [ ] Component: Flink Runner
   - [ ] Component: Prism Runner
   - [ ] Component: Twister2 Runner
   - [ ] Component: Hazelcast Jet Runner
   - [x] Component: Google Cloud Dataflow Runner


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