From 9b4dcbf15bc5ae003c6dee66118ab479439c372b Mon Sep 17 00:00:00 2001 From: Zihan Dai Date: Sat, 15 Aug 2026 19:57:29 +1000 Subject: [PATCH] Tag the Kafka write error output with the schema it actually emits ErrorCounterFn emits ErrorHandling.errorRecord(errorSchema, ...), where errorSchema is already ErrorHandling.errorSchema(inputSchema). The error PCollection was tagged with the wrapper applied a second time, declaring a shape no emitted element can match. Same defect as apache/beam#39759 in TFRecordWriteSchemaTransformProvider. The existing tests apply ErrorCounterFn directly and tag the output themselves, so none of them reach the transform's expand(). Signed-off-by: Zihan Dai --- .../KafkaWriteSchemaTransformProvider.java | 3 +-- ...KafkaWriteSchemaTransformProviderTest.java | 24 +++++++++++++++++++ 2 files changed, 25 insertions(+), 2 deletions(-) diff --git a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaWriteSchemaTransformProvider.java b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaWriteSchemaTransformProvider.java index b9c41746240a..ea03ea5ff147 100644 --- a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaWriteSchemaTransformProvider.java +++ b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaWriteSchemaTransformProvider.java @@ -297,8 +297,7 @@ public PCollectionRowTuple expand(PCollectionRowTuple input) { } // TODO: include output from KafkaIO Write once updated from PDone - PCollection errorOutput = - outputTuple.get(ERROR_TAG).setRowSchema(ErrorHandling.errorSchema(errorSchema)); + PCollection errorOutput = outputTuple.get(ERROR_TAG).setRowSchema(errorSchema); return PCollectionRowTuple.of( handleErrors ? configuration.getErrorHandling().getOutput() : "errors", errorOutput); } diff --git a/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaWriteSchemaTransformProviderTest.java b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaWriteSchemaTransformProviderTest.java index fed783f70a00..dfba9ab59da6 100644 --- a/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaWriteSchemaTransformProviderTest.java +++ b/sdks/java/io/kafka/src/test/java/org/apache/beam/sdk/io/kafka/KafkaWriteSchemaTransformProviderTest.java @@ -49,6 +49,7 @@ import org.apache.beam.sdk.transforms.SerializableFunction; import org.apache.beam.sdk.values.KV; import org.apache.beam.sdk.values.PCollection; +import org.apache.beam.sdk.values.PCollectionRowTuple; import org.apache.beam.sdk.values.PCollectionTuple; import org.apache.beam.sdk.values.Row; import org.apache.beam.sdk.values.TupleTag; @@ -271,6 +272,29 @@ public void testBuildTransformWithManaged() { } } + @Test + public void testErrorOutputCarriesTheSchemaErrorCounterFnEmits() { + // The output schema is fixed while the graph is built, so this needs no runner. + p.enableAbandonedNodeEnforcement(false); + + Schema inputSchema = Schema.builder().addByteArrayField("bytes").build(); + KafkaWriteSchemaTransformProvider.KafkaWriteSchemaTransformConfiguration configuration = + KafkaWriteSchemaTransformProvider.KafkaWriteSchemaTransformConfiguration.builder() + .setFormat("RAW") + .setTopic("test-topic") + .setBootstrapServers("host:9092") + .setErrorHandling(ErrorHandling.builder().setOutput("errors").build()) + .build(); + + PCollectionRowTuple output = + PCollectionRowTuple.of("input", p.apply(Create.empty(inputSchema))) + .apply(new KafkaWriteSchemaTransformProvider().from(configuration)); + + // ErrorCounterFn emits ErrorHandling.errorRecord(errorSchema, ...), where errorSchema is + // already ErrorHandling.errorSchema(inputSchema). + assertEquals(ErrorHandling.errorSchema(inputSchema), output.get("errors").getSchema()); + } + @Test public void testKafkaWriteSchemaTransformConfigurationSchema() throws NoSuchSchemaException { Schema schema =