apache / apache/beam

Apache Beam BigQuery Storage Write Api error checkTimestamp(SimpleDoFnRunner.java:259)

Open
#22,314 1 comment 0 reactions 0 assignees View on GitHub
Dominant language
Java
Stars
8.7k
Forks
4.7k
Avg merge
1d 20h
Merged PRs (30d)
196

Description

Hi,
I used Apache Beam in version 2.40.0
I Write pipeline:
` Pipeline pipeline = Pipeline.create(options);

PCollection messages = pipeline .apply("Read PubSub Messages", PubsubIO.readMessagesWithAttributesAndMessageId().fromSubscription("subscription"));

WriteResult out = messages.apply("WriteSuccessfulRecourds",
BigQueryIO.write()
.withoutValidation()
.withFormatFunction((PubsubMessage elem) -> {
System.out.println("Start transform: " + elem);
TableRow tableRow = new TableRow().set("ID", "1").set("NAME", "TEST").set("DATE", new String(elem.getPayload()));//"2022-01-01T23:15");
tableRow.set("messageCustom", elem);
return tableRow;
})
.withCreateDisposition(CreateDisposition.CREATE_NEVER)
.withWriteDisposition(WriteDisposition.WRITE_APPEND)
.withExtendedErrorInfo()
.withMethod(BigQueryIO.Write.Method.STORAGE_WRITE_API)
//.withTriggeringFrequency(Duration.millis(500))
.withNumStorageWriteApiStreams(20)
.withFailedInsertRetryPolicy(InsertRetryPolicy.retryTransientErrors())
.ignoreUnknownValues()
//.withExtendedErrorInfo()
.to(String.format("%s:%s.%s", project", "TEST_BIG_QUERY", "TEST_BIG_QUERY_TABLE")));

out.getFailedStorageApiInserts().apply("Error", new PubsubMessageToTableRow(options));

return pipeline.run();`

And after 60 minutes i have this error:

`Exception in thread "main" org.apache.beam.sdk.Pipeline$PipelineExecutionException: java.lang.IllegalArgumentException: Cannot output with timestamp 294247-01-10T04:00:54.776Z. Output timestamps must be no earlier than the timestamp of the current input or timer (294247-01-10T04:00:54.776Z) minus the allowed skew (0 milliseconds) and no later than 294247-01-10T04:00:54.775Z. See the DoFn#getAllowedTimestampSkew() Javadoc for details on changing the allowed skew.
at org.apache.beam.runners.direct.DirectRunner$DirectPipelineResult.waitUntilFinish(DirectRunner.java:373)
at org.apache.beam.runners.direct.DirectRunner$DirectPipelineResult.waitUntilFinish(DirectRunner.java:341)
at org.apache.beam.runners.direct.DirectRunner.run(DirectRunner.java:218)
at org.apache.beam.runners.direct.DirectRunner.run(DirectRunner.java:67)
at org.apache.beam.sdk.Pipeline.run(Pipeline.java:323)
at org.apache.beam.sdk.Pipeline.run(Pipeline.java:309)
at pl.woxtech.dataflow.pubsub.to.bigquery.GCPPubSubToBigQuery.run(GCPPubSubToBigQuery.java:132)
at pl.woxtech.dataflow.pubsub.to.bigquery.GCPPubSubToBigQuery.main(GCPPubSubToBigQuery.java:73)
Caused by: java.lang.IllegalArgumentException: Cannot output with timestamp 294247-01-10T04:00:54.776Z. Output timestamps must be no earlier than the timestamp of the current input or timer (294247-01-10T04:00:54.776Z) minus the allowed skew (0 milliseconds) and no later than 294247-01-10T04:00:54.775Z. See the DoFn#getAllowedTimestampSkew() Javadoc for details on changing the allowed skew.
at org.apache.beam.repackaged.direct_java.runners.core.SimpleDoFnRunner.checkTimestamp(SimpleDoFnRunner.java:259)
at org.apache.beam.repackaged.direct_java.runners.core.SimpleDoFnRunner.access$1300(SimpleDoFnRunner.java:85)
at org.apache.beam.repackaged.direct_java.runners.core.SimpleDoFnRunner$OnTimerArgumentProvider.output(SimpleDoFnRunner.java:843)
at org.apache.beam.sdk.transforms.DoFnOutputReceivers$WindowedContextOutputReceiver.output(DoFnOutputReceivers.java:76)
at org.apache.beam.sdk.io.gcp.bigquery.StorageApiWritesShardedRecords$WriteRecordsDoFn.finalizeStream(StorageApiWritesShardedRecords.java:536)
at org.apache.beam.sdk.io.gcp.bigquery.StorageApiWritesShardedRecords$WriteRecordsDoFn.onTimer(StorageApiWritesShardedRecords.java:550)`

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.