apache / apache/beam

BigQueryIO withFailedInsertRetryPolicy is endlessly retrying "invalid" rows

Open
#18,531 0 comments 0 reactions 0 assignees View on GitHub
bug gcp io java P3
Dominant language
Java
Stars
8.7k
Forks
4.7k
Avg merge
2d 2h
Merged PRs (30d)
205

Description

Using the InsertRetryPolicy.retryTransientErrors() on streaming data into a BigQuery table is endlessly retrying "invalid" rows.

To quote Eugene Kirpichov [~kirpichov]
bq. Upon talking to the BigQuery team, it became clear that this is indeed a bug in BigQueryIO. This error is not reported via InsertErrors because the InsertAll request specifies the table once rather than per row, and the table is invalid, so all rows in the batch are invalid. Beam should handle this.

```

p.apply(BigQueryIO.writeTableRows()
.to(new DatePartitionedTableSpecifier(tableReference,
"tracking data"))
.withSchema(schema)
.withCreateDisposition(BigQueryIO.Write.CreateDisposition.CREATE_IF_NEEDED)

.withWriteDisposition(BigQueryIO.Write.WriteDisposition.WRITE_APPEND)

.withFailedInsertRetryPolicy(InsertRetryPolicy.retryTransientErrors())
// write all failed inserts
to a DMQ
.getFailedInserts().apply(MapElements.via(new SimpleFunction()
{
public PubsubMessage apply(final TableRow _row) {
try {
return new PubsubMessage(JacksonFactory.getDefaultInstance().toByteArray(_row),

Collections.emptyMap());
} catch (IOException e)
{
throw new RuntimeException("failed to write to DMQ", e);
}
}
})).apply(PubsubIO.writeMessages().to("projects/gameduell-bits-bigquery-poc/topics/dmq"));

```

```

(1a04bdb0d43aca9c): java.lang.RuntimeException: com.google.api.client.googleapis.json.GoogleJsonResponseException:
400 Bad Request
{
"code" : 400,
"errors" : [ {
"domain" : "global",
"message" : "The
destination table's partition rum$20170925 is outside the allowed bounds. You can only stream to partitions
within 31 days in the past and 16 days in the future relative to the current date.",
"reason" :
"invalid"
} ],
"message" : "The destination table's partition rum$20170925 is outside the allowed
bounds. You can only stream to partitions within 31 days in the past and 16 days in the future relative
to the current date.",
"status" : "INVALID_ARGUMENT"
}
org.apache.beam.sdk.io.gcp.bigquery.BigQueryServicesImpl$DatasetServiceImpl.insertAll(BigQueryServicesImpl.java:774)

org.apache.beam.sdk.io.gcp.bigquery.BigQueryServicesImpl$DatasetServiceImpl.insertAll(BigQueryServicesImpl.java:809)

org.apache.beam.sdk.io.gcp.bigquery.StreamingWriteFn.flushRows(StreamingWriteFn.java:126)

org.apache.beam.sdk.io.gcp.bigquery.StreamingWriteFn.finishBundle(StreamingWriteFn.java:96)
Caused
by: com.google.api.client.googleapis.json.GoogleJsonResponseException: 400 Bad Request
{
"code"
: 400,
"errors" : [ {
"domain" : "global",
"message" : "The destination table's partition
rum$20170925 is outside the allowed bounds. You can only stream to partitions within 31 days in the
past and 16 days in the future relative to the current date.",
"reason" : "invalid"
} ],

"message" : "The destination table's partition rum$20170925 is outside the allowed bounds. You can only
stream to partitions within 31 days in the past and 16 days in the future relative to the current date.",

"status" : "INVALID_ARGUMENT"
}
com.google.api.client.googleapis.json.GoogleJsonResponseException.from(GoogleJsonResponseException.java:146)

com.google.api.client.googleapis.services.json.AbstractGoogleJsonClientRequest.newExceptionOnError(AbstractGoogleJsonClientRequest.java:113)

com.google.api.client.googleapis.services.json.AbstractGoogleJsonClientRequest.newExceptionOnError(AbstractGoogleJsonClientRequest.java:40)

com.google.api.client.googleapis.services.AbstractGoogleClientRequest$1.interceptResponse(AbstractGoogleClientRequest.java:321)

com.google.api.client.http.HttpRequest.execute(HttpRequest.java:1065)
com.google.api.client.googleapis.services.AbstractGoogleClientRequest.executeUnparsed(AbstractGoogleClientRequest.java:419)

com.google.api.client.googleapis.services.AbstractGoogleClientRequest.executeUnparsed(AbstractGoogleClientRequest.java:352)

com.google.api.client.googleapis.services.AbstractGoogleClientRequest.execute(AbstractGoogleClientRequest.java:469)

org.apache.beam.sdk.io.gcp.bigquery.BigQueryServicesImpl$DatasetServiceImpl$1.call(BigQueryServicesImpl.java:720)

org.apache.beam.sdk.io.gcp.bigquery.BigQueryServicesImpl$DatasetServiceImpl$1.call(BigQueryServicesImpl.java:712)

java.util.concurrent.FutureTask.run(FutureTask.java:266)
java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142)

java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617)
java.lang.Thread.run(Thread.java:745)
java.lang.RuntimeException:
com.google.api.client.googleapis.json.GoogleJsonResponseException: 400 Bad Request
{
"code" : 400,

"errors" : [ {
"domain" : "global",
"message" : "The destination table's partition rum$20170925
is outside the allowed bounds. You can only stream to partitions within 31 days in the past and 16 days
in the future relative to the current date.",
"reason" : "invalid"
} ],
"message" : "The
destination table's partition rum$20170925 is outside the allowed bounds. You can only stream to partitions
within 31 days in the past and 16 days in the future relative to the current date.",
"status" : "INVALID_ARGUMENT"
}

org.apache.beam.sdk.io.gcp.bigquery.BigQueryServicesImpl$DatasetServiceImpl.insertAll(BigQueryServicesImpl.java:774)

org.apache.beam.sdk.io.gcp.bigquery.BigQueryServicesImpl$DatasetServiceImpl.insertAll(BigQueryServicesImpl.java:809)

org.apache.beam.sdk.io.gcp.bigquery.StreamingWriteFn.flushRows(StreamingWriteFn.java:126)

org.apache.beam.sdk.io.gcp.bigquery.StreamingWriteFn.finishBundle(StreamingWriteFn.java:96)
Caused
by: com.google.api.client.googleapis.json.GoogleJsonResponseException: 400 Bad Request
{
"code"
: 400,
"errors" : [ {
"domain" : "global",
"message" : "The destination table's partition
rum$20170925 is outside the allowed bounds. You can only stream to partitions within 31 days in the
past and 16 days in the future relative to the current date.",
"reason" : "invalid"
} ],

"message" : "The destination table's partition rum$20170925 is outside the allowed bounds. You can only
stream to partitions within 31 days in the past and 16 days in the future relative to the current date.",

"status" : "INVALID_ARGUMENT"
}
com.google.api.client.googleapis.json.GoogleJsonResponseException.from(GoogleJsonResponseException.java:146)

com.google.api.client.googleapis.services.json.AbstractGoogleJsonClientRequest.newExceptionOnError(AbstractGoogleJsonClientRequest.java:113)

com.google.api.client.googleapis.services.json.AbstractGoogleJsonClientRequest.newExceptionOnError(AbstractGoogleJsonClientRequest.java:40)

com.google.api.client.googleapis.services.AbstractGoogleClientRequest$1.interceptResponse(AbstractGoogleClientRequest.java:321)

com.google.api.client.http.HttpRequest.execute(HttpRequest.java:1065)
com.google.api.client.googleapis.services.AbstractGoogleClientRequest.executeUnparsed(AbstractGoogleClientRequest.java:419)

com.google.api.client.googleapis.services.AbstractGoogleClientRequest.executeUnparsed(AbstractGoogleClientRequest.java:352)

com.google.api.client.googleapis.services.AbstractGoogleClientRequest.execute(AbstractGoogleClientRequest.java:469)

org.apache.beam.sdk.io.gcp.bigquery.BigQueryServicesImpl$DatasetServiceImpl$1.call(BigQueryServicesImpl.java:720)

org.apache.beam.sdk.io.gcp.bigquery.BigQueryServicesImpl$DatasetServiceImpl$1.call(BigQueryServicesImpl.java:712)

java.util.concurrent.FutureTask.run(FutureTask.java:266)
java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142)

java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617)
java.lang.Thread.run(Thread.java:745)
java.lang.RuntimeException:
com.google.api.client.googleapis.json.GoogleJsonResponseException: 400 Bad Request
{
"code" : 400,

"errors" : [ {
"domain" : "global",
"message" : "The destination table's partition rum$20170925
is outside the allowed bounds. You can only stream to partitions within 31 days in the past and 16 days
in the future relative to the current date.",
"reason" : "invalid"
} ],
"message" : "The
destination table's partition rum$20170925 is outside the allowed bounds. You can only stream to partitions
within 31 days in the past and 16 days in the future relative to the current date.",
"status" : "INVALID_ARGUMENT"
}

org.apache.beam.sdk.io.gcp.bigquery.BigQueryServicesImpl$DatasetServiceImpl.insertAll(BigQueryServicesImpl.java:774)

org.apache.beam.sdk.io.gcp.bigquery.BigQueryServicesImpl$DatasetServiceImpl.insertAll(BigQueryServicesImpl.java:809)

org.apache.beam.sdk.io.gcp.bigquery.StreamingWriteFn.flushRows(StreamingWriteFn.java:126)

org.apache.beam.sdk.io.gcp.bigquery.StreamingWriteFn.finishBundle(StreamingWriteFn.java:96)
Caused
by: com.google.api.client.googleapis.json.GoogleJsonResponseException: 400 Bad Request
{
"code"
: 400,
"errors" : [ {
"domain" : "global",
"message" : "The destination table's partition
rum$20170925 is outside the allowed bounds. You can only stream to partitions within 31 days in the
past and 16 days in the future relative to the current date.",
"reason" : "invalid"
} ],

"message" : "The destination table's partition rum$20170925 is outside the allowed bounds. You can only
stream to partitions within 31 days in the past and 16 days in the future relative to the current date.",

"status" : "INVALID_ARGUMENT"
}
com.google.api.client.googleapis.json.GoogleJsonResponseException.from(GoogleJsonResponseException.java:146)

com.google.api.client.googleapis.services.json.AbstractGoogleJsonClientRequest.newExceptionOnError(AbstractGoogleJsonClientRequest.java:113)

com.google.api.client.googleapis.services.json.AbstractGoogleJsonClientRequest.newExceptionOnError(AbstractGoogleJsonClientRequest.java:40)

com.google.api.client.googleapis.services.AbstractGoogleClientRequest$1.interceptResponse(AbstractGoogleClientRequest.java:321)

com.google.api.client.http.HttpRequest.execute(HttpRequest.java:1065)
com.google.api.client.googleapis.services.AbstractGoogleClientRequest.executeUnparsed(AbstractGoogleClientRequest.java:419)

com.google.api.client.googleapis.services.AbstractGoogleClientRequest.executeUnparsed(AbstractGoogleClientRequest.java:352)

com.google.api.client.googleapis.services.AbstractGoogleClientRequest.execute(AbstractGoogleClientRequest.java:469)

org.apache.beam.sdk.io.gcp.bigquery.BigQueryServicesImpl$DatasetServiceImpl$1.call(BigQueryServicesImpl.java:720)

org.apache.beam.sdk.io.gcp.bigquery.BigQueryServicesImpl$DatasetServiceImpl$1.call(BigQueryServicesImpl.java:712)

java.util.concurrent.FutureTask.run(FutureTask.java:266)
java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142)

java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617)
java.lang.Thread.run(Thread.java:745)
java.lang.RuntimeException:
com.google.api.client.googleapis.json.GoogleJsonResponseException: 400 Bad Request
{
"code" : 400,

"errors" : [ {
"domain" : "global",
"message" : "The destination table's partition rum$20170925
is outside the allowed bounds. You can only stream to partitions within 31 days in the past and 16 days
in the future relative to the current date.",
"reason" : "invalid"
} ],
"message" : "The
destination table's partition rum$20170925 is outside the allowed bounds. You can only stream to partitions
within 31 days in the past and 16 days in the future relative to the current date.",
"status" : "INVALID_ARGUMENT"
}

org.apache.beam.sdk.io.gcp.bigquery.BigQueryServicesImpl$DatasetServiceImpl.insertAll(BigQueryServicesImpl.java:774)

org.apache.beam.sdk.io.gcp.bigquery.BigQueryServicesImpl$DatasetServiceImpl.insertAll(BigQueryServicesImpl.java:809)

org.apache.beam.sdk.io.gcp.bigquery.StreamingWriteFn.flushRows(StreamingWriteFn.java:126)

org.apache.beam.sdk.io.gcp.bigquery.StreamingWriteFn.finishBundle(StreamingWriteFn.java:96)
Caused
by: com.google.api.client.googleapis.json.GoogleJsonResponseException: 400 Bad Request
{
"code"
: 400,
"errors" : [ {
"domain" : "global",
"message" : "The destination table's partition
rum$20170925 is outside the allowed bounds. You can only stream to partitions within 31 days in the
past and 16 days in the future relative to the current date.",
"reason" : "invalid"
} ],

"message" : "The destination table's partition rum$20170925 is outside the allowed bounds. You can only
stream to partitions within 31 days in the past and 16 days in the future relative to the current date.",

"status" : "INVALID_ARGUMENT"
}
com.google.api.client.googleapis.json.GoogleJsonResponseException.from(GoogleJsonResponseException.java:146)

com.google.api.client.googleapis.services.json.AbstractGoogleJsonClientRequest.newExceptionOnError(AbstractGoogleJsonClientRequest.java:113)

com.google.api.client.googleapis.services.json.AbstractGoogleJsonClientRequest.newExceptionOnError(AbstractGoogleJsonClientRequest.java:40)

com.google.api.client.googleapis.services.AbstractGoogleClientRequest$1.interceptResponse(AbstractGoogleClientRequest.java:321)

com.google.api.client.http.HttpRequest.execute(HttpRequest.java:1065)
com.google.api.client.googleapis.services.AbstractGoogleClientRequest.executeUnparsed(AbstractGoogleClientRequest.java:419)

com.google.api.client.googleapis.services.AbstractGoogleClientRequest.executeUnparsed(AbstractGoogleClientRequest.java:352)

com.google.api.client.googleapis.services.AbstractGoogleClientRequest.execute(AbstractGoogleClientRequest.java:469)

org.apache.beam.sdk.io.gcp.bigquery.BigQueryServicesImpl$DatasetServiceImpl$1.call(BigQueryServicesImpl.java:720)

org.apache.beam.sdk.io.gcp.bigquery.BigQueryServicesImpl$DatasetServiceImpl$1.call(BigQueryServicesImpl.java:712)

java.util.concurrent.FutureTask.run(FutureTask.java:266)
java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142)

java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617)
java.lang.Thread.run(Thread.java:745)
java.lang.RuntimeException:
com.google.api.client.googleapis.json.GoogleJsonResponseException: 400 Bad Request
{
"code" : 400,

"errors" : [ {
"domain" : "global",
"message" : "The destination table's partition rum$20170925
is outside the allowed bounds. You can only stream to partitions within 31 days in the past and 16 days
in the future relative to the current date.",
"reason" : "invalid"
} ],
"message" : "The
destination table's partition rum$20170925 is outside the allowed bounds. You can only stream to partitions
within 31 days in the past and 16 days in the future relative to the current date.",
"status" : "INVALID_ARGUMENT"
}

org.apache.beam.sdk.io.gcp.bigquery.BigQueryServicesImpl$DatasetServiceImpl.insertAll(BigQueryServicesImpl.java:774)

org.apache.beam.sdk.io.gcp.bigquery.BigQueryServicesImpl$DatasetServiceImpl.insertAll(BigQueryServicesImpl.java:809)

org.apache.beam.sdk.io.gcp.bigquery.StreamingWriteFn.flushRows(StreamingWriteFn.java:126)

org.apache.beam.sdk.io.gcp.bigquery.StreamingWriteFn.finishBundle(StreamingWriteFn.java:96)
Caused
by: com.google.api.client.googleapis.json.GoogleJsonResponseException: 400 Bad Request
{
"code"
: 400,
"errors" : [ {
"domain" : "global",
"message" : "The destination table's partition
rum$20170925 is outside the allowed bounds. You can only stream to partitions within 31 days in the
past and 16 days in the future relative to the current date.",
"reason" : "invalid"
} ],

"message" : "The destination table's partition rum$20170925 is outside the allowed bounds. You can only
stream to partitions within 31 days in the past and 16 days in the future relative to the current date.",

"status" : "INVALID_ARGUMENT"
}
com.google.api.client.googleapis.json.GoogleJsonResponseException.from(GoogleJsonResponseException.java:146)

com.google.api.client.googleapis.services.json.AbstractGoogleJsonClientRequest.newExceptionOnError(AbstractGoogleJsonClientRequest.java:113)

com.google.api.client.googleapis.services.json.AbstractGoogleJsonClientRequest.newExceptionOnError(AbstractGoogleJsonClientRequest.java:40)

com.google.api.client.googleapis.services.AbstractGoogleClientRequest$1.interceptResponse(AbstractGoogleClientRequest.java:321)

com.google.api.client.http.HttpRequest.execute(HttpRequest.java:1065)
com.google.api.client.googleapis.services.AbstractGoogleClientRequest.executeUnparsed(AbstractGoogleClientRequest.java:419)

com.google.api.client.googleapis.services.AbstractGoogleClientRequest.executeUnparsed(AbstractGoogleClientRequest.java:352)

com.google.api.client.googleapis.services.AbstractGoogleClientRequest.execute(AbstractGoogleClientRequest.java:469)

org.apache.beam.sdk.io.gcp.bigquery.BigQueryServicesImpl$DatasetServiceImpl$1.call(BigQueryServicesImpl.java:720)

org.apache.beam.sdk.io.gcp.bigquery.BigQueryServicesImpl$DatasetServiceImpl$1.call(BigQueryServicesImpl.java:712)

java.util.concurrent.FutureTask.run(FutureTask.java:266)
java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142)

java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617)
java.lang.Thread.run(Thread.java:745)
java.lang.RuntimeException:
com.google.api.client.googleapis.json.GoogleJsonResponseException: 400 Bad Request
{
"code" : 400,

"errors" : [ {
"domain" : "global",
"message" : "The destination table's partition rum$20170925
is outside the allowed bounds. You can only stream to partitions within 31 days in the past and 16 days
in the future relative to the current date.",
"reason" : "invalid"
} ],
"message" : "The
destination table's partition rum$20170925 is outside the allowed bounds. You can only stream to partitions
within 31 days in the past and 16 days in the future relative to the current date.",
"status" : "INVALID_ARGUMENT"
}

org.apache.beam.sdk.io.gcp.bigquery.BigQueryServicesImpl$DatasetServiceImpl.insertAll(BigQueryServicesImpl.java:774)

org.apache.beam.sdk.io.gcp.bigquery.BigQueryServicesImpl$DatasetServiceImpl.insertAll(BigQueryServicesImpl.java:809)

org.apache.beam.sdk.io.gcp.bigquery.StreamingWriteFn.flushRows(StreamingWriteFn.java:126)

org.apache.beam.sdk.io.gcp.bigquery.StreamingWriteFn.finishBundle(StreamingWriteFn.java:96)
Caused
by: com.google.api.client.googleapis.json.GoogleJsonResponseException: 400 Bad Request
{
"code"
: 400,
"errors" : [ {
"domain" : "global",
"message" : "The destination table's partition
rum$20170925 is outside the allowed bounds. You can only stream to partitions within 31 days in the
past and 16 days in the future relative to the current date.",
"reason" : "invalid"
} ],

"message" : "The destination table's partition rum$20170925 is outside the allowed bounds. You can only
stream to partitions within 31 days in the past and 16 days in the future relative to the current date.",

"status" : "INVALID_ARGUMENT"
}
com.google.api.client.googleapis.json.GoogleJsonResponseException.from(GoogleJsonResponseException.java:146)

com.google.api.client.googleapis.services.json.AbstractGoogleJsonClientRequest.newExceptionOnError(AbstractGoogleJsonClientRequest.java:113)

com.google.api.client.googleapis.services.json.AbstractGoogleJsonClientRequest.newExceptionOnError(AbstractGoogleJsonClientRequest.java:40)

com.google.api.client.googleapis.services.AbstractGoogleClientRequest$1.interceptResponse(AbstractGoogleClientRequest.java:321)

com.google.api.client.http.HttpRequest.execute(HttpRequest.java:1065)
com.google.api.client.googleapis.services.AbstractGoogleClientRequest.executeUnparsed(AbstractGoogleClientRequest.java:419)

com.google.api.client.googleapis.services.AbstractGoogleClientRequest.executeUnparsed(AbstractGoogleClientRequest.java:352)

com.google.api.client.googleapis.services.AbstractGoogleClientRequest.execute(AbstractGoogleClientRequest.java:469)

org.apache.beam.sdk.io.gcp.bigquery.BigQueryServicesImpl$DatasetServiceImpl$1.call(BigQueryServicesImpl.java:720)

org.apache.beam.sdk.io.gcp.bigquery.BigQueryServicesImpl$DatasetServiceImpl$1.call(BigQueryServicesImpl.java:712)

java.util.concurrent.FutureTask.run(FutureTask.java:266)
java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142)

java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617)
java.lang.Thread.run(Thread.java:745)

```

Imported from Jira [BEAM-3271](https://issues.apache.org/jira/browse/BEAM-3271). Original Jira may contain additional context.
Reported by: cw_krebs.

Contributor guide

Open the contributing guide

Research direction

Start with StreamingWriteFn.flushRows and finishBundle, then inspect BigQueryServicesImpl.DatasetServiceImpl.insertAll from the stack trace. Reproduce the invalid partition-table response with retryTransientErrors(); done means table-level invalid errors stop retrying endlessly while failed row inserts still follow the configured policy.

Written by the indexing model from the issue text.

Assessment

Tech stack
google-cloud, java
Domain
cloud, data
Issue type
Bug
Difficulty
3/5
Estimated time
1-2 days
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
35/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.