Skip to content

Commit 6eade1a

Browse files
committed
fix another gemini review
1 parent 06cee19 commit 6eade1a

2 files changed

Lines changed: 11 additions & 3 deletions

File tree

sdks/java/io/datadog/src/main/java/org/apache/beam/sdk/io/datadog/DatadogWriteSchemaTransformProvider.java

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -22,11 +22,13 @@
2222
import com.google.auto.service.AutoService;
2323
import java.util.Collections;
2424
import java.util.List;
25+
import org.apache.beam.sdk.coders.RowCoder;
2526
import org.apache.beam.sdk.schemas.Schema;
2627
import org.apache.beam.sdk.schemas.transforms.SchemaTransform;
2728
import org.apache.beam.sdk.schemas.transforms.SchemaTransformProvider;
2829
import org.apache.beam.sdk.schemas.transforms.TypedSchemaTransformProvider;
2930
import org.apache.beam.sdk.schemas.transforms.providers.ErrorHandling;
31+
import org.apache.beam.sdk.transforms.Create;
3032
import org.apache.beam.sdk.transforms.DoFn;
3133
import org.apache.beam.sdk.transforms.ParDo;
3234
import org.apache.beam.sdk.values.PCollection;
@@ -172,7 +174,11 @@ public PCollectionRowTuple expand(PCollectionRowTuple input) {
172174
return PCollectionRowTuple.of(errorHandling.getOutput(), allErrors);
173175
} else {
174176
writeErrors.apply("Fail on Write Error", ParDo.of(new FailOnWriteErrorFn()));
175-
return PCollectionRowTuple.empty(input.getPipeline());
177+
PCollection<Row> emptyErrors =
178+
input
179+
.getPipeline()
180+
.apply("Empty Errors Placeholder", Create.empty(RowCoder.of(dynamicErrorSchema)));
181+
return PCollectionRowTuple.of(ERROR, emptyErrors);
176182
}
177183
}
178184
}

sdks/java/io/datadog/src/test/java/org/apache/beam/sdk/io/datadog/DatadogWriteSchemaTransformProviderTest.java

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -158,7 +158,8 @@ public void testWriteBuildTransformAndRun() {
158158
PCollection<Row> input = p.apply("Create", Create.of(ROWS).withRowSchema(SCHEMA));
159159
PCollectionRowTuple inputTuple = PCollectionRowTuple.of(INPUT, input);
160160
PCollectionRowTuple output = transform.expand(inputTuple);
161-
assertTrue(output.getAll().isEmpty());
161+
assertEquals(1, output.getAll().size());
162+
assertTrue(output.has(ERROR));
162163

163164
assertThrows(PipelineExecutionException.class, () -> p.run().waitUntilFinish());
164165
}
@@ -298,7 +299,8 @@ public void testBuildTransformWithAllParameters() {
298299
PCollection<Row> input = p.apply("Create", Create.of(ROWS).withRowSchema(SCHEMA));
299300
PCollectionRowTuple inputTuple = PCollectionRowTuple.of(INPUT, input);
300301
PCollectionRowTuple output = transform.expand(inputTuple);
301-
assertTrue(output.getAll().isEmpty());
302+
assertEquals(1, output.getAll().size());
303+
assertTrue(output.has(ERROR));
302304

303305
assertThrows(PipelineExecutionException.class, () -> p.run().waitUntilFinish());
304306
}

0 commit comments

Comments
 (0)