Skip to content

Commit b91e8aa

Browse files
committed
fix test failure
1 parent 4ac6c18 commit b91e8aa

2 files changed

Lines changed: 15 additions & 0 deletions

File tree

sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryIOTranslation.java

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -438,6 +438,7 @@ static class BigQueryIOWriteTranslator implements TransformPayloadTranslator<Wri
438438
.addNullableByteArrayField("row_mutation_information_fn")
439439
.addNullableByteArrayField("bad_record_error_handler")
440440
.addNullableByteArrayField("bad_record_router")
441+
.addNullableInt64Field("auto_schema_update_strict_timeout_ms")
441442
.build();
442443

443444
public static final String BIGQUERY_WRITE_TRANSFORM_URN =
@@ -582,6 +583,9 @@ public Row toConfigRow(Write<?> transform) {
582583
fieldValues.put("use_beam_schema", transform.getUseBeamSchema());
583584
fieldValues.put("auto_sharding", transform.getAutoSharding());
584585
fieldValues.put("auto_schema_update", transform.getAutoSchemaUpdate());
586+
fieldValues.put(
587+
"auto_schema_update_strict_timeout_ms",
588+
transform.getAutoSchemaUpdateStrictTimeout().getMillis());
585589
if (transform.getWriteProtosClass() != null) {
586590
fieldValues.put("write_protos_class", toByteArray(transform.getWriteProtosClass()));
587591
}
@@ -870,6 +874,15 @@ public Write<?> fromConfigRow(Row configRow, PipelineOptions options) {
870874
if (autoSchemaUpdate != null) {
871875
builder = builder.setAutoSchemaUpdate(autoSchemaUpdate);
872876
}
877+
Long autoSchemaUpdateStrictTimeoutMillis =
878+
configRow.getInt64("auto_schema_update_strict_timeout_ms");
879+
if (autoSchemaUpdateStrictTimeoutMillis != null) {
880+
builder =
881+
builder
882+
.setAutoSchemaUpdate(true)
883+
.setAutoSchemaUpdateStrictTimeout(
884+
org.joda.time.Duration.millis(autoSchemaUpdateStrictTimeoutMillis));
885+
}
873886
byte[] writeProtosClasses = configRow.getBytes("write_protos_class");
874887
if (writeProtosClasses != null) {
875888
builder =

sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery/BigQueryIOTranslationTest.java

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -132,6 +132,8 @@ public class BigQueryIOTranslationTest {
132132
WRITE_TRANSFORM_SCHEMA_MAPPING.put("getAutoSharding", "auto_sharding");
133133
WRITE_TRANSFORM_SCHEMA_MAPPING.put("getPropagateSuccessful", "propagate_successful");
134134
WRITE_TRANSFORM_SCHEMA_MAPPING.put("getAutoSchemaUpdate", "auto_schema_update");
135+
WRITE_TRANSFORM_SCHEMA_MAPPING.put(
136+
"getAutoSchemaUpdateStrictTimeout", "auto_schema_update_strict_timeout_ms");
135137
WRITE_TRANSFORM_SCHEMA_MAPPING.put("getWriteProtosClass", "write_protos_class");
136138
WRITE_TRANSFORM_SCHEMA_MAPPING.put("getDirectWriteProtos", "direct_write_protos");
137139
WRITE_TRANSFORM_SCHEMA_MAPPING.put("getDeterministicRecordIdFn", "deterministic_record_id_fn");

0 commit comments

Comments
 (0)